diff --git a/gather_program/.idea/misc.xml b/gather_program/.idea/misc.xml new file mode 100644 index 000000000..2e1a5699d --- /dev/null +++ b/gather_program/.idea/misc.xml @@ -0,0 +1,67 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + 1.8 + + + + + + + + \ No newline at end of file diff --git a/gather_program/.idea/workspace.xml b/gather_program/.idea/workspace.xml new file mode 100644 index 000000000..464ce9784 --- /dev/null +++ b/gather_program/.idea/workspace.xml @@ -0,0 +1,215 @@ + + + + + + + + + + + + + true + DEFINITION_ORDER + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + 1492270889131 + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/gather_program/.project b/gather_program/.project index 22f0f1b7c..3be1413dc 100644 --- a/gather_program/.project +++ b/gather_program/.project @@ -10,6 +10,11 @@ + + org.springframework.ide.eclipse.core.springbuilder + + + org.eclipse.m2e.core.maven2Builder @@ -17,6 +22,7 @@ + org.springframework.ide.eclipse.core.springnature org.eclipse.jdt.core.javanature org.eclipse.m2e.core.maven2Nature diff --git a/gather_program/bin/resources/applicationContext-myBatis.xml b/gather_program/bin/resources/applicationContext-myBatis.xml index 910d3e8ae..a1aa5cab9 100644 --- a/gather_program/bin/resources/applicationContext-myBatis.xml +++ b/gather_program/bin/resources/applicationContext-myBatis.xml @@ -39,9 +39,9 @@ destroy-method="close"> + value="jdbc:mysql://172.16.128.36:3306/ossean_production?characterEncoding=UTF-8&zeroDateTimeBehavior=convertToNull&autoReconnect=true" /> - + diff --git a/gather_program/bin/resources/gather_projects.xml b/gather_program/bin/resources/gather_projects.xml index bbd0a819b..6a675f038 100644 --- a/gather_program/bin/resources/gather_projects.xml +++ b/gather_program/bin/resources/gather_projects.xml @@ -3,7 +3,7 @@ TableFlow pointers - oschina_project,openhub_project,sourceforge_project,apache,freecode_project + github,oschina_project,openhub_project,sourceforge_project,apache,freecode_project gather_projects id,name,tags,url,url_md5,description,language,source,license,homepage,now(),extracted_time,created_time diff --git a/gather_program/src/main/java/org/ossean/gather/process/GatherProcess.java b/gather_program/src/main/java/org/ossean/gather/process/GatherProcess.java index 799f3c137..d4af98e07 100644 --- a/gather_program/src/main/java/org/ossean/gather/process/GatherProcess.java +++ b/gather_program/src/main/java/org/ossean/gather/process/GatherProcess.java @@ -67,8 +67,8 @@ public class GatherProcess { configureName = args[0].toString(); } else { logger.error("404 configure!"); - //configureName = "relative_memos"; - configureName = "gather_projects"; + configureName = "relative_memos"; + //configureName = "gather_projects"; //return; } ApplicationContext applicationContext = new ClassPathXmlApplicationContext("classpath:/applicationContext*.xml"); diff --git a/gather_program/src/main/java/org/ossean/gather/process/GatherThreadNew.java b/gather_program/src/main/java/org/ossean/gather/process/GatherThreadNew.java index 8e0028860..c1c4253a6 100644 --- a/gather_program/src/main/java/org/ossean/gather/process/GatherThreadNew.java +++ b/gather_program/src/main/java/org/ossean/gather/process/GatherThreadNew.java @@ -135,11 +135,15 @@ public class GatherThreadNew implements Runnable { // 插入Id段数据,忽略重复值 try { + logger.info("准备查询批量数据"); String selectItems = getSelectItems(conf); //这里的目标表是relative_memos if (conf.getTargetTableName().equals("relative_memos")) { + logger.info("对帖子进行汇总"); + List dataGet = gatherDao.getPostGatherData(sourceTableName, selectItems, beginId, endId, - conf.getAndWhere()); + conf.getAndWhere()); + logger.info("从数据源表中提取到数据"); gatherPosts(dataGet); } else @@ -163,6 +167,7 @@ public class GatherThreadNew implements Runnable { } catch (Exception ex) { // 数据迁移过程可能发生异常情况 + ex.printStackTrace(); logger.error(ex); System.exit(0); } @@ -190,6 +195,9 @@ public class GatherThreadNew implements Runnable { public void handleUpdateGatherPosts(int id, RelativeMemo model_new) { targetDao.updateRelativeMemo(gatherPostsTableName, model_new, id);// 更新数据relative_memos表 } + public void handleUpdateGatherPostsByUrlMD5(int id, RelativeMemo model_new) { + targetDao.updateRelativeMemoByUrlMD5(gatherPostsTableName, model_new, id);// 更新数据relative_memos表 + } // 处理URL不存在的项目 插入gather_projects表 public void handleInsertGatherProjects(GatherProject model, Configure conf) { @@ -255,6 +263,16 @@ public class GatherThreadNew implements Runnable { } } + /** + * 最新修改的程序更好地考虑到了程序的有效性,主要分以下情况: + * 1.pk_control_posts中已存在(查询到pk_model),则先根据url_md5查询targetTable中是否存在改对象(设为model)。 + * 如果model不为空,则判断pk_model.getId == model.getId是否成立,成立则直接更新,不成立则将pk_control中的 + * 更新为model.getId. + * 2.pk_control_posts中不存在,则直接插入生成对象并再次查询到pk_model,先根据url_md5查询targetTable中是否存在改对象(设为model)。 + * 如果model不为空,则判断pk_model.getId == model.getId是否成立,成立则直接更新,不成立则将pk_control中的 + * 更新为model.getId. + * @param dataGet + */ @Transactional(propagation=Propagation.REQUIRED) public void gatherPosts(List dataGet){ for (int i = 0; i < dataGet.size(); i++) { @@ -262,35 +280,122 @@ public class GatherThreadNew implements Runnable { String urlMD5 = model.getUrl_md5();// 通过urlMD5判断是不是已经存在该帖子、是否更新 int postId = 0; + logger.info("开始汇总:" + model.getId()); //查找pk_control_posts表以判定此帖子是否已经存在 PKControlPosts pkControlModel = pkControlPostsDao.selectItemByUrlMD5( - pkControlPostsTableName, urlMD5); - + pkControlPostsTableName, urlMD5); //帖子已经存在则进行更新操作 if(pkControlModel != null){ - model.setId(pkControlModel.getId()); + try{ - if(targetDao.selectRelativeMemosItem(conf.getTargetTableName(), model.getId())!=null) - handleUpdateGatherPosts(pkControlModel.getId(), model); + logger.info("帖子已经在pk_control_posts存在"); + RelativeMemo memosModel = targetDao.selectRelativeMemosItemByUrlMD5(conf.getTargetTableName(), model.getUrl_md5()); + if(memosModel!=null) + { + logger.info("帖子确实已经插入到了relative_memos表中"); + try + { + if(memosModel.getId() == pkControlModel.getId()) + { + logger.error("memosModel.getId() == pkControlModel.getId()"); + model.setId(pkControlModel.getId()); + handleUpdateGatherPosts(pkControlModel.getId(), model); + } + else + { + logger.error("memosModel.getId() != pkControlModel.getId()"); + gatherDao.updatePkControlTable(pkControlPostsTableName, + memosModel.getId(), memosModel.getUrl_md5()); + pkControlModel.setId(model.getId()); + model.setId(memosModel.getId()); + targetDao.deleteWhileExist(gatherPostsTableName,model); + handleUpdateGatherPostsByUrlMD5(pkControlModel.getId(), model); + } + + logger.info("更新完成"); + } + catch(Exception e) + { + logger.error(e); + logger.error("更新失败:295"); + System.exit(0); + } + } else - handleInsertGatherPosts(model, conf); + { + logger.info("帖子实际上并未插入到relative_memos表中"); + try{ + model.setId(pkControlModel.getId()); + targetDao.deleteWhileExist(gatherPostsTableName,model); + handleInsertGatherPosts(model, conf); + logger.info("更新转为插入"); + } + catch (Exception e) + { + logger.error("更新转为插入时出错:298"); + System.exit(0); + + } + } }catch(Exception e){ logger.error("更新帖子时出错:" + e); + System.exit(0); } } - else{//帖子不存在就插入pk_control_posts表并为帖子生成唯一固定ID + else + {//帖子不存在就插入pk_control_posts表并为帖子生成唯一固定ID + logger.info("准备将帖子插入pk_control_posts"); + try{ pkControlPostsDao.insertOneItem(pkControlPostsTableName, urlMD5); + + }catch(Exception e) + { + logger.info("插入pkcontorl时失败"); + //System.exit(0); + continue; + } PKControlPosts controlItem = pkControlPostsDao.selectItemByUrlMD5( pkControlPostsTableName, urlMD5); - model.setId(controlItem.getId()); - try{ - handleInsertGatherPosts(model, conf); + RelativeMemo memosModel2 = targetDao.selectRelativeMemosItemByUrlMD5(conf.getTargetTableName(), model.getUrl_md5()); + + if(memosModel2 == null) + { + model.setId(controlItem.getId()); + targetDao.deleteWhileExist(gatherPostsTableName,model); + handleInsertGatherPosts(model, conf); + } + else + { + if(controlItem.getId() == memosModel2.getId()) + { + model.setId(controlItem.getId()); + int sucess2 = targetDao.updateRelativeMemoByUrlMD5(gatherPostsTableName, model, model.getId()); + if(sucess2 != 1) + { + logger.error("365"); + System.exit(0); + } + } + else + { + gatherDao.updatePkControlTable(pkControlPostsTableName, memosModel2.getId(), memosModel2.getUrl_md5()); + model.setId(memosModel2.getId()); + controlItem.setId(memosModel2.getId()); + targetDao.deleteWhileExist(gatherPostsTableName,model); + handleUpdateGatherPosts(controlItem.getId(), model); + + } + + } + logger.info("插入relative_memos成功"); }catch(Exception e){ logger.error("插入帖子时出错:" + e); + int sucess = targetDao.updateRelativeMemoByUrlMD5(gatherPostsTableName, model, model.getId()); + if(sucess != 1)System.exit(0); } } diff --git a/gather_program/src/main/java/org/ossean/gather/sourceDao/GatherDao.java b/gather_program/src/main/java/org/ossean/gather/sourceDao/GatherDao.java index b1b0268b0..d8a5ee523 100644 --- a/gather_program/src/main/java/org/ossean/gather/sourceDao/GatherDao.java +++ b/gather_program/src/main/java/org/ossean/gather/sourceDao/GatherDao.java @@ -24,7 +24,16 @@ public interface GatherDao { @Param("taskTableName") String taskTableName, @Param("sourceTableName") String sourceTableName, @Param("targetTableName") String targetTableName); - + + /** + * 更新pk_control_posts + * @param targettable + * @param id + * @param url_md5 + */ + @Update("update ${targettable} set id = #{id} where url_md5 = #{url_md5}") + public void updatePkControlTable(@Param("targettable")String targettable,@Param("id")int id, + @Param("url_md5")String url_md5); // 查询未处理的task @Select("SELECT BeginId,EndId,BeginTime,EndTime " + "FROM ${migrationTask} " diff --git a/gather_program/src/main/java/org/ossean/gather/targetDao/TargetDao.java b/gather_program/src/main/java/org/ossean/gather/targetDao/TargetDao.java index 284e2bce2..7f9452eb4 100644 --- a/gather_program/src/main/java/org/ossean/gather/targetDao/TargetDao.java +++ b/gather_program/src/main/java/org/ossean/gather/targetDao/TargetDao.java @@ -32,6 +32,22 @@ public interface TargetDao { public void updateRelativeMemo( @Param("targetTable") String targetTableName, @Param("model") RelativeMemo model, @Param("id") int id); + + // 对urlMD5相同的数据进行update操作 + @Update("update ${targetTable} set id=#{model.id},title=#{model.title},content=#{model.content},created_time=#{model.created_time},updated_time=#{model.updated_time}," + + "memo_type=#{model.memo_type},tags=#{model.tags},source=#{model.source},url=#{model.url},url_md5=#{model.url_md5},author=#{model.author},author_url=#{model.author_url}," + + "view_num=#{model.view_num},review_num=#{model.review_num},extracted_time=#{model.extracted_time} where url_md5=#{model.url_md5}") + public Integer updateRelativeMemoByUrlMD5( + @Param("targetTable") String targetTableName, + @Param("model") RelativeMemo model, @Param("id") int id); + @Update("update ${targettable} set id = #{id} where url_md5 = #{url_md5}") + public void updatePkControlTable(@Param("targettable")String targettable,@Param("id")int id, + @Param("url_md5")String url_md5); + +// 查找relative_memos表对应url_md5的记录 + @Select("select * from ${table} where url_md5=#{url_md5}") + public RelativeMemo selectRelativeMemosItemByUrlMD5( + @Param("table") String table, @Param("url_md5") String url_md5); // 向tag表存储数据 @Insert("insert ignore into ${table} (name) values (#{name})") @@ -122,5 +138,7 @@ public interface TargetDao { @Select("select * from ${table} where url_md5=#{urlMD5}") public JobRequirement findJobByUrlMD5(@Param("table") String table, @Param("urlMD5") String urlMD5); + @Delete("delete from ${gatherPostsTableName} where id = #{model.id} and url_md5 != #{model.url_md5}") + public void deleteWhileExist(@Param("gatherPostsTableName")String gatherPostsTableName, @Param("model")RelativeMemo model); } diff --git a/gather_program/src/main/resource/applicationContext-myBatis.xml b/gather_program/src/main/resource/applicationContext-myBatis.xml index dc36f7a0a..abb8e8115 100644 --- a/gather_program/src/main/resource/applicationContext-myBatis.xml +++ b/gather_program/src/main/resource/applicationContext-myBatis.xml @@ -19,7 +19,7 @@ destroy-method="close"> + value="jdbc:mysql://localhost:3306/ossean_gather?characterEncoding=UTF-8&zeroDateTimeBehavior=convertToNull&autoReconnect=true" /> @@ -39,7 +39,7 @@ destroy-method="close"> + value="jdbc:mysql://localhost:3306/ossean_gather?characterEncoding=UTF-8&zeroDateTimeBehavior=convertToNull&autoReconnect=true" /> diff --git a/gather_program/src/main/resource/relative_memos.xml b/gather_program/src/main/resource/relative_memos.xml index 9c3e13cf4..6522dd261 100644 --- a/gather_program/src/main/resource/relative_memos.xml +++ b/gather_program/src/main/resource/relative_memos.xml @@ -3,10 +3,10 @@ TableFlow pointers - 51cto_blog + stack_over_flow relative_memos id,title,content,created_time,now(),type,tags,source,url,url_md5,author,author_url,view_num,review_num,extracted_time - id,title,content,created_time,updated_time,type,tags,source,url,url_md5,author,author_url,view_num,review_num,extracted_time + id,title,content,created_time,updated_time,memo_type,tags,source,url,url_md5,author,author_url,view_num,review_num,extracted_time 10000 diff --git a/match_program/lukeall-5.3.1.jar b/match_program/lukeall-5.3.1.jar new file mode 100644 index 000000000..9ca45eb75 Binary files /dev/null and b/match_program/lukeall-5.3.1.jar differ diff --git a/match_program/src/main/java/com/ossean/match/lucene/LuceneSearch.java b/match_program/src/main/java/com/ossean/match/lucene/LuceneSearch.java index 66794e196..b6463768c 100644 --- a/match_program/src/main/java/com/ossean/match/lucene/LuceneSearch.java +++ b/match_program/src/main/java/com/ossean/match/lucene/LuceneSearch.java @@ -91,7 +91,7 @@ public class LuceneSearch { String prjId = d.get(LuceneIndex.prjIdFieldName); String[] prjNames = d.getValues(searchField); for(String prjName : prjNames){ - if (keyWords.contains(prjName)) { + if (keyWords.contains(prjName.toLowerCase())) { int pId = Integer.parseInt(prjId); if (matchMap.containsKey(pId)) { matchMap.put(pId, matchMap.get(pId) + weight + sd.score/1000); diff --git a/match_program/src/main/java/com/ossean/match/matchprocess/Match.java b/match_program/src/main/java/com/ossean/match/matchprocess/Match.java index a7c25df14..ddb3dd62d 100644 --- a/match_program/src/main/java/com/ossean/match/matchprocess/Match.java +++ b/match_program/src/main/java/com/ossean/match/matchprocess/Match.java @@ -172,6 +172,8 @@ public class Match { logger.error("insertPrjToMemoMatchResult error: " + e); } } + if(list.size()!=0) + batchInsertJDBC(list,getTargetTable(prjId)); //System.out.println(prjId+" current insert time cost:"+(System.currentTimeMillis()-start)/1000+" seconds"); } diff --git a/project_manager/bin/Multithreadprojectsfilter.sh b/project_manager/bin/Multithreadprojectsfilter.sh new file mode 100644 index 000000000..122e75949 --- /dev/null +++ b/project_manager/bin/Multithreadprojectsfilter.sh @@ -0,0 +1,20 @@ +#!/bin/bash + +find ./target/classes -name "*.properties"|xargs rm -f +find ./target/classes -name "*.xml"|xargs rm -f +find ./target/classes -name "*.dic"|xargs rm -f + +#export CLASSPATH=$CURR_DIR/lib:$CURR_DIR:$JAVA_HOME/lib:$JAVA_HOME/jre/lib + +tmp='./target/classes':$tmp +tmp='./target/project_manager-0.0.1-SNAPSHOT-jar-with-dependencies-without-resources/*':$tmp +tmp='./bin/resources':$tmp +CLASSPATH=$tmp:$CLASSPATH + + +echo $CLASSPATH +JVM_ARGS="-Xmn98m -Xmx1024m -Xms512m -XX:NewRatio=4 -XX:SurvivorRatio=4 -XX:MaxTenuringThreshold=2" +#echo JVM_ARGS=$JVM_ARGS +#ulimit -n 400000 +#echo "" > nohup.out +java $JVM_ARGS -classpath $CLASSPATH com.ossean.projectmanager.ProjectsFilterProcessMain > log/projectsfilter.log 2>&1 & \ No newline at end of file diff --git a/project_manager/bin/resources/applicationContext_mybatis.xml b/project_manager/bin/resources/applicationContext_mybatis.xml index 6c5596ea6..793511843 100644 --- a/project_manager/bin/resources/applicationContext_mybatis.xml +++ b/project_manager/bin/resources/applicationContext_mybatis.xml @@ -39,9 +39,9 @@ destroy-method="close"> + value="jdbc:mysql://172.16.128.30:3306/ossean_production?characterEncoding=UTF-8&zeroDateTimeBehavior=convertToNull&autoReconnect=true" /> - + diff --git a/project_manager/src/main/java/com/ossean/projectmanager/AppContext.java b/project_manager/src/main/java/com/ossean/projectmanager/AppContext.java new file mode 100644 index 000000000..680a61afd --- /dev/null +++ b/project_manager/src/main/java/com/ossean/projectmanager/AppContext.java @@ -0,0 +1,11 @@ +package com.ossean.projectmanager; + +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.ClassPathXmlApplicationContext; + +public class AppContext { + + static ApplicationContext appContext = new ClassPathXmlApplicationContext( + "classpath:/applicationContext*.xml"); +} + diff --git a/project_manager/src/main/java/com/ossean/projectmanager/ProjectsFilterProcessMain.java b/project_manager/src/main/java/com/ossean/projectmanager/ProjectsFilterProcessMain.java new file mode 100644 index 000000000..7e4aba256 --- /dev/null +++ b/project_manager/src/main/java/com/ossean/projectmanager/ProjectsFilterProcessMain.java @@ -0,0 +1,133 @@ +package com.ossean.projectmanager; +import java.text.SimpleDateFormat; +import java.util.Date; +/** + * 将筛选过程写成多线程,这是入口程序 + * 在这里调用筛选线程 + */ +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import javax.annotation.Resource; +import org.apache.log4j.Logger; +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.stereotype.Component; + +import com.ossean.projectmanager.lasttabledao.OpenSourceProjectDao; +import com.ossean.projectmanager.lasttabledao.RelativeMemoToOpenSourceProjectDao; +import com.ossean.projectmanager.parttabledao.PartProjectDao; +import com.ossean.projectmanager.projectsfilter.ProjectsFilterThread; +@Component +public class ProjectsFilterProcessMain { + + @Resource + private OpenSourceProjectDao lastProjectDao; + + @Resource + private PartProjectDao partProjectDao; + + @Resource + private RelativeMemoToOpenSourceProjectDao matchResultDao; + + //线程池,开20个线程 + private ExecutorService pool = Executors.newFixedThreadPool(30); + + int minId=0 ; + int maxId=0 ; + + Logger logger = Logger.getLogger(this.getClass()); + + public static void main(String[] args){ + @SuppressWarnings("resource") + ApplicationContext applicationContext = new ClassPathXmlApplicationContext("classpath:/applicationContext*.xml"); + ProjectsFilterProcessMain mainClass = applicationContext.getBean(ProjectsFilterProcessMain.class); + mainClass.start(); + } + + + private void start() + { + /** + * 实现了对数据的监测,实时监测是否有新的数据到来 + * 循环查询数据库看是否有可以筛选的数据(即filteration=0) + * 当没有需要筛选的数据时,程序睡眠 + */ + long startTime = System.currentTimeMillis(); +// while(true) +// { + // TODO Auto-generated method stub + minId = lastProjectDao.getMinId(0); + maxId = lastProjectDao.getMaxId(0); + logger.info(minId + ":" + maxId); + + //获取当前时间,用于记录线程活动 + SimpleDateFormat df = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss"); + String date = df.format(new Date()); + startFilter(); + + /** + * 等待进程池中所有进程执行完毕才开始再一次扫描检查是否有新的数据到来 + * 即程序在某一时刻检查到有需要筛选的数据时就一次性将这些数据筛选完毕 + * 再进行下一次扫描检查是否有新的数据到来 + */ + pool.shutdown(); + while(!pool.isTerminated()) + { + + } + + logger.info(date + "时刻开始的筛选已经完成"); + try { + //休眠20秒等待所有操作完成 + Thread.sleep(20000); + } catch (InterruptedException e) { + logger.error(e); + } + System.gc(); //手动垃圾回收 + +// try { +// //本次扫描的结果处理完毕就休眠两个小时 +// Thread.sleep(2*60*60*1000); +// } catch (InterruptedException e) { +// logger.error(e); +// } + //} + + } + + + //进行具体的筛选操作 + private void startFilter() { + int startId = minId; + int endId = minId; + int sumPerThread = (maxId - minId)/30 + 1;//每个线程处理的数据量 + while(true) + { + /** + * 每个线程负责startId到endId部分的数据的筛选 + */ + endId = startId + sumPerThread; + + if(endId < maxId) + { + //若minId = maxId,说明筛选已完成,暂时没有需要筛选的数据 + ProjectsFilterThread filterThread = (ProjectsFilterThread)AppContext.appContext.getBean("ProjectsFilterThread"); + filterThread.setBorder(startId,endId); + pool.execute(filterThread); + } + else + { + endId = maxId; + //若minId = maxId,说明给各线程分配数据已完成。 + ProjectsFilterThread filterThread = (ProjectsFilterThread)AppContext.appContext.getBean("ProjectsFilterThread"); + filterThread.setBorder(startId,endId); + pool.execute(filterThread); + break; + } + + startId = endId + 1; + } + logger.info("线程分配完毕"); + } + +} diff --git a/project_manager/src/main/java/com/ossean/projectmanager/lasttabledao/OpenSourceProjectDao.java b/project_manager/src/main/java/com/ossean/projectmanager/lasttabledao/OpenSourceProjectDao.java index bb60a54c9..353a323c5 100644 --- a/project_manager/src/main/java/com/ossean/projectmanager/lasttabledao/OpenSourceProjectDao.java +++ b/project_manager/src/main/java/com/ossean/projectmanager/lasttabledao/OpenSourceProjectDao.java @@ -51,9 +51,14 @@ public interface OpenSourceProjectDao { public List getBatchPrjs(@Param("startId") int startId, @Param("batchSize") int batchSize); + // 根据开始和结束位置批量获取项目 + @Select("select id,source,url,filtration from open_source_projects where filtration=0 and id >= #{startId} and id <= #{endId}") + public List getPrjsByBorder(@Param("startId") int startId, + @Param("endId") int endId); @Select("select min(id) from open_source_projects where filtration=#{filtration}") public int getMinId(@Param("filtration") int filtration); - + @Select("select max(id) from open_source_projects where filtration=#{filtration}") + public int getMaxId(@Param("filtration") int filtration); // 删除项目 @Update("delete from open_source_projects where id=#{id}") public void deleteProject(@Param("id") int id); diff --git a/project_manager/src/main/java/com/ossean/projectmanager/model/GithubProject.java b/project_manager/src/main/java/com/ossean/projectmanager/model/GithubProject.java new file mode 100644 index 000000000..141246bbb --- /dev/null +++ b/project_manager/src/main/java/com/ossean/projectmanager/model/GithubProject.java @@ -0,0 +1,40 @@ +package com.ossean.projectmanager.model; + +public class GithubProject { + private int id; + private int forks; + private String url; + private String url_mad5; + private String name; + public int getId() { + return id; + } + public void setId(int id) { + this.id = id; + } + public int getForks() { + return forks; + } + public void setForks(int forks) { + this.forks = forks; + } + public String getUrl() { + return url; + } + public void setUrl(String url) { + this.url = url; + } + public String getUrl_mad5() { + return url_mad5; + } + public void setUrl_mad5(String url_mad5) { + this.url_mad5 = url_mad5; + } + public String getName() { + return name; + } + public void setName(String name) { + this.name = name; + } + +} diff --git a/project_manager/src/main/java/com/ossean/projectmanager/parttabledao/PartProjectDao.java b/project_manager/src/main/java/com/ossean/projectmanager/parttabledao/PartProjectDao.java index cfa9bb21b..18c40e046 100644 --- a/project_manager/src/main/java/com/ossean/projectmanager/parttabledao/PartProjectDao.java +++ b/project_manager/src/main/java/com/ossean/projectmanager/parttabledao/PartProjectDao.java @@ -3,6 +3,7 @@ package com.ossean.projectmanager.parttabledao; import org.apache.ibatis.annotations.Param; import org.apache.ibatis.annotations.Select; +import com.ossean.projectmanager.model.GithubProject; import com.ossean.projectmanager.model.OpenhubProject; import com.ossean.projectmanager.model.SourceForgeProject; @@ -14,5 +15,7 @@ public interface PartProjectDao { @Select("select name,description,download_num,favor_num from sourceforge_project where url = #{url} group by url_md5 order by extracted_time desc") public SourceForgeProject getSourceForgePrjByUrl( @Param("url") String url); + @Select("select name,forks from github where url = #{url} group by url_md5 order by extracted_time desc") + public GithubProject getGithubPrjByUrl( @Param("url") String url); } diff --git a/project_manager/src/main/java/com/ossean/projectmanager/projectsfilter/ProjectsFilter.java b/project_manager/src/main/java/com/ossean/projectmanager/projectsfilter/ProjectsFilter.java index 8a0ec811c..dd44c74f0 100644 --- a/project_manager/src/main/java/com/ossean/projectmanager/projectsfilter/ProjectsFilter.java +++ b/project_manager/src/main/java/com/ossean/projectmanager/projectsfilter/ProjectsFilter.java @@ -12,6 +12,7 @@ import com.ossean.projectmanager.lasttabledao.OpenSourceProjectDao; import com.ossean.projectmanager.lasttabledao.PointersDao; import com.ossean.projectmanager.lasttabledao.RelativeMemoToOpenSourceProjectDao; import com.ossean.projectmanager.model.OpenhubProject; +import com.ossean.projectmanager.model.GithubProject; import com.ossean.projectmanager.model.OpenSourceProject; import com.ossean.projectmanager.model.SourceForgeProject; import com.ossean.projectmanager.parttabledao.PartProjectDao; @@ -37,11 +38,13 @@ public class ProjectsFilter { */ public void filtratePrjs() { logger.info("Reading projects......"); + long startTime = System.currentTimeMillis(); int startId = lastProjectDao.getMinId(0); while (true) { List prjsList = lastProjectDao .getBatchPrjs(startId,batchsize); if(prjsList.size()==0){ + logger.info((System.currentTimeMillis() - startTime)/1000); logger.info("Filter done......sleeping......"); try { Thread.sleep(1000*60*15);// 筛选完成,休息15分钟 @@ -126,7 +129,30 @@ public class ProjectsFilter { getTargetTable(project.getId()), project.getId()); // 删除该项目的匹配结果,确保无之前的匹配结果 } - } else { + } else + if(source.equals("github")){ + GithubProject githubProject = partProjectDao + .getGithubPrjByUrl(url); //根据URL从github分表中获取github + //github项目的筛选规则是,forks>0,且名字不为空 + if(githubProject != null + && !"".equals(githubProject.getName()) + &&githubProject.getName() != null + &&githubProject.getForks() > 0) + { + lastProjectDao.updateFiltratedPrj(project.getId(),1); // 筛选标识从0变为1,表示该项目经过筛选新增的 + matchResultDao.deleteMatchResult( + getTargetTable(project.getId()), + project.getId()); // 删除该项目的匹配结果,确保无之前的匹配结果 + } + else + { + matchResultDao.deleteMatchResult( + getTargetTable(project.getId()), + project.getId()); // 删除该项目的匹配结果,确保无之前的匹配结果 + } + + }else + { logger.info("Unknown source... source = " + source); } } @@ -144,13 +170,13 @@ public class ProjectsFilter { * @return */ public static String getTargetTable(int ospId) { - String targetTableName = ""; - if (ospId >= 770000) { - targetTableName = "relative_memo_to_open_source_projects_70"; - } else { - int a = 1 + ospId / 11000; - targetTableName = "relative_memo_to_open_source_projects_" + a; - } + String targetTableName = "relative_memo_to_open_source_projects_70_lqk_test"; +// if (ospId >= 770000) { +// targetTableName = "relative_memo_to_open_source_projects_70"; +// } else { +// int a = 1 + ospId / 11000; +// targetTableName = "relative_memo_to_open_source_projects_" + a; +// } // if (osp_id < 500) { // targetTableName = "relative_memo_to_open_source_projects_1"; // } diff --git a/project_manager/src/main/java/com/ossean/projectmanager/projectsfilter/ProjectsFilterThread.java b/project_manager/src/main/java/com/ossean/projectmanager/projectsfilter/ProjectsFilterThread.java new file mode 100644 index 000000000..dcd80b0e8 --- /dev/null +++ b/project_manager/src/main/java/com/ossean/projectmanager/projectsfilter/ProjectsFilterThread.java @@ -0,0 +1,207 @@ +package com.ossean.projectmanager.projectsfilter; +/** + * 将筛选程序写成多线程 + * 基本思路是:程序初始化时,读取未完成筛选的数据的总量SUM,然后设置batchSize = SUM/20; + * 每20个线程中,每个线程负责处理batchSize条数据。 + */ +import java.util.List; + +import javax.annotation.Resource; + +import org.apache.commons.lang3.StringUtils; +import org.apache.log4j.Logger; +import org.springframework.context.annotation.Scope; +import org.springframework.stereotype.Component; + +import com.ossean.projectmanager.lasttabledao.OpenSourceProjectDao; +import com.ossean.projectmanager.lasttabledao.PointersDao; +import com.ossean.projectmanager.lasttabledao.RelativeMemoToOpenSourceProjectDao; +import com.ossean.projectmanager.model.OpenhubProject; +import com.ossean.projectmanager.model.GithubProject; +import com.ossean.projectmanager.model.OpenSourceProject; +import com.ossean.projectmanager.model.SourceForgeProject; +import com.ossean.projectmanager.parttabledao.PartProjectDao; + +@Component("ProjectsFilterThread") +@Scope("prototype") +public class ProjectsFilterThread implements Runnable { + + //用于标志是否终止当前进程 + boolean keepRunning = true; + + @Resource + private OpenSourceProjectDao lastProjectDao; + + @Resource + private PartProjectDao partProjectDao; + + @Resource + private RelativeMemoToOpenSourceProjectDao matchResultDao; + + Logger logger = Logger.getLogger(this.getClass()); + private final int batchsize = 5000; + private int MaxId;// + private int MinId; + + //设置当前线程所处理数据的开始记录的id和结束记录的Id + public void setBorder(int MinId,int MaxId) + { + this.MinId = MinId; + this.MaxId = MaxId; + } + + @Override + public void run() { + filterPrjs(); + } + + /** + * 对open_source_projects表根据各个社区的特定字段做筛选 + */ + public void filterPrjs() { + Thread.currentThread().setName("Thread:" + this.MinId + " to " + this.MaxId + ":"); + int startId = this.MinId;//上次批量筛选操作开始的位置 + int endId = this.MinId;//上次批量筛选操作结束的位置 + logger.info(Thread.currentThread().getName() + " begin Reading projects......"); + while (true) { + endId = startId + batchsize; + + if(endId > MaxId) + endId = MaxId; + + List prjsList = lastProjectDao + .getPrjsByBorder(startId,endId); + if(prjsList.size()==0){ + startId = endId + 1; + continue; + } + else{ + //logger.info("Filtering projects......"); + for (OpenSourceProject project : prjsList) { + //logger.info("project id: " + project.getId()); + String prjUrl = project.getUrl(); + String source = ""; + String url = ""; + if (prjUrl == null || "".equals(prjUrl)) { + lastProjectDao.updateFiltratedPrj(project.getId(), 0); + continue; + } + if (prjUrl.contains("|,|")) { // 即url中包含多个项目来源 + String firstUrl = StringUtils.splitByWholeSeparator(prjUrl, + "|,|")[0];// 只对第一个,即去重时保留的最热的项目来源做筛选。 + source = StringUtils.splitByWholeSeparator(firstUrl, "|:|")[0]; // 从url字段中取得第一个来源社区。 + url = StringUtils.splitByWholeSeparator(firstUrl, "|:|")[1]; // 获得第一个url + } else { // url只有一个项目来源 + source = StringUtils.splitByWholeSeparator(prjUrl, "|:|")[0]; + url = StringUtils.splitByWholeSeparator(prjUrl, "|:|")[1]; + } + if (source != null && source.length() > 0) { + source = source.toLowerCase(); + } else { + continue; + } + if (source.equals("openhub")) { + OpenhubProject openhubProject = partProjectDao + .getOpenHubPrjByUrl(url); // 根据url从openhub的项目分表获得项目信息 + if (openhubProject != null + && openhubProject.getName() != null + && !"".equals(openhubProject.getName()) + && openhubProject.getDescription() != null + && !"".equals(openhubProject.getDescription()) + && openhubProject.getCode_repository() != null + && !openhubProject.getCode_repository().contains( + "Add a code location")) { // openhub的筛选条件为name、description不为空,且该项目有版本库 + lastProjectDao.updateFiltratedPrj(project.getId(), + 1); // 筛选标识从0变为1,表示该项目经过筛选新增的 + matchResultDao.deleteMatchResult( + getTargetTable(project.getId()), + project.getId()); // 删除该项目的匹配结果,确保无之前的匹配结果 + } else { + //lastProjectDao.deleteProject(project.getId()); + matchResultDao.deleteMatchResult( + getTargetTable(project.getId()), + project.getId()); // 删除该项目的匹配结果 + } + } else if (source.equals("sourceforge")) { + SourceForgeProject sourceforgeProject = partProjectDao + .getSourceForgePrjByUrl(url); // 根据url从SourceForge的项目分表获得项目信息 + if (sourceforgeProject != null + && sourceforgeProject.getName() != null + && !"".equals(sourceforgeProject.getName()) + && sourceforgeProject.getDescription() != null + && !"".equals(sourceforgeProject.getDescription()) + && ((sourceforgeProject.getDownload_num() > 0) || (sourceforgeProject + .getStars() > 0))) { + lastProjectDao.updateFiltratedPrj(project.getId(), + 1); // 筛选标识从0变为1,表示该项目经过筛选新增的 + matchResultDao.deleteMatchResult( + getTargetTable(project.getId()), + project.getId()); // 删除该项目的匹配结果,确保无之前的匹配结果 + } else { + //lastProjectDao.deleteProject(project.getId()); + matchResultDao.deleteMatchResult( + getTargetTable(project.getId()), + project.getId()); // 删除该项目的匹配结果 + } + } else if (source.equals("oschina") || source.equals("apache")) { + if (project.getFilration() == 0) { + lastProjectDao.updateFiltratedPrj(project.getId(), 1); // 筛选标识从0变为1,表示该项目经过筛选新增的 + matchResultDao.deleteMatchResult( + getTargetTable(project.getId()), + project.getId()); // 删除该项目的匹配结果,确保无之前的匹配结果 + } + } else + if(source.equals("github")){ + GithubProject githubProject = partProjectDao + .getGithubPrjByUrl(url); //根据URL从github分表中获取github + if(githubProject != null + && !"".equals(githubProject.getName()) + &&githubProject.getName() != null + &&githubProject.getForks() > 0) + { + lastProjectDao.updateFiltratedPrj(project.getId(),1); // 筛选标识从0变为1,表示该项目经过筛选新增的 + matchResultDao.deleteMatchResult( + getTargetTable(project.getId()), + project.getId()); // 删除该项目的匹配结果,确保无之前的匹配结果 + } + else + { + matchResultDao.deleteMatchResult( + getTargetTable(project.getId()), + project.getId()); // 删除该项目的匹配结果,确保无之前的匹配结果 + } + + }else + { + logger.info("Unknown source... source = " + source); + } + } + } + + if(endId >= MaxId) + break; + + startId = endId + 1; + } + + logger.info( Thread.currentThread().getName() + "所负责的部分筛选完毕!"); + } + + /** + * get the match result table's name + * + * @param osp_id + * @return + */ + public static String getTargetTable(int ospId) { + String targetTableName = ""; + if (ospId >= 770000) { + targetTableName = "relative_memo_to_open_source_projects_70"; + } else { + int a = 1 + ospId / 11000; + targetTableName = "relative_memo_to_open_source_projects_" + a; + } + return targetTableName; + } + +} diff --git a/project_manager/src/main/resource/applicationContext-myBatis.xml b/project_manager/src/main/resource/applicationContext-myBatis.xml index aea32198d..58f9e7f78 100644 --- a/project_manager/src/main/resource/applicationContext-myBatis.xml +++ b/project_manager/src/main/resource/applicationContext-myBatis.xml @@ -19,7 +19,7 @@ destroy-method="close"> + value="jdbc:mysql://localhost:3306/test?characterEncoding=UTF-8&zeroDateTimeBehavior=convertToNull&autoReconnect=true" /> @@ -39,7 +39,7 @@ destroy-method="close"> + value="jdbc:mysql://localhost:3306/test?characterEncoding=UTF-8&zeroDateTimeBehavior=convertToNull&autoReconnect=true" /> diff --git a/project_match/bin/start_merge_projects_github.sh b/project_match/bin/start_merge_projects_github.sh new file mode 100644 index 000000000..eabd27011 --- /dev/null +++ b/project_match/bin/start_merge_projects_github.sh @@ -0,0 +1,16 @@ +#!/bin/bash + +find ./target/classes -name "*.properties"|xargs rm -f +find ./target/classes -name "*.xml"|xargs rm -f +find ./target/classes -name "*.dic"|xargs rm -f + +#export CLASSPATH=$CURR_DIR/lib:$CURR_DIR:$JAVA_HOME/lib:$JAVA_HOME/jre/lib + +tmp='./target/classes':$tmp +tmp='./target/Project_Match-0.0.1-SNAPSHOT-jar-with-dependencies-without-resources/*':$tmp +tmp='./bin/resources':$tmp +CLASSPATH=$tmp:$CLASSPATH + + +echo $CLASSPATH +java -classpath $CLASSPATH com.ossean.MergeProjectsForGithub >>log/merge_projects_2017.log 2>&1 & diff --git a/project_match/bin/start_multi_transfer_projects.sh b/project_match/bin/start_multi_transfer_projects.sh new file mode 100644 index 000000000..4a702b179 --- /dev/null +++ b/project_match/bin/start_multi_transfer_projects.sh @@ -0,0 +1,16 @@ +#!/bin/bash + +find ./target/classes -name "*.properties"|xargs rm -f +find ./target/classes -name "*.xml"|xargs rm -f +find ./target/classes -name "*.dic"|xargs rm -f + +#export CLASSPATH=$CURR_DIR/lib:$CURR_DIR:$JAVA_HOME/lib:$JAVA_HOME/jre/lib + +tmp='./target/classes':$tmp +tmp='./target/Project_Match-0.0.1-SNAPSHOT-jar-with-dependencies-without-resources/*':$tmp +tmp='./bin/resources':$tmp +CLASSPATH=$tmp:$CLASSPATH + + +echo $CLASSPATH +java -classpath $CLASSPATH com.ossean.TransferManageProcess >>log/multi_transfer_projects.log 2>&1 & diff --git a/project_match/src/main/java/com/ossean/MergeProjectsForGithub.java b/project_match/src/main/java/com/ossean/MergeProjectsForGithub.java new file mode 100644 index 000000000..d927444e5 --- /dev/null +++ b/project_match/src/main/java/com/ossean/MergeProjectsForGithub.java @@ -0,0 +1,101 @@ +package com.ossean; + +import java.text.DateFormat; +import java.text.ParseException; +import java.text.SimpleDateFormat; +import java.util.Date; +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +import javax.annotation.Resource; + +import org.apache.log4j.Logger; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.stereotype.Component; + +import com.ossean.dao.DBSource; +import com.ossean.dao.GatherDao; +import com.ossean.dao.UpdateControlProjectsDao; +import com.ossean.model.GatherProjectsModel; +import com.ossean.util.MergeProjectNew2; + +@Component +public class MergeProjectsForGithub { + Logger logger = Logger.getLogger(this.getClass()); + @Resource + private DBSource dbSource; + @Resource + private GatherDao gatherDao; + @Resource + private UpdateControlProjectsDao updateControlDao; + @Qualifier("mergeProjectNew2") + @Autowired + private MergeProjectNew2 mergeProjectNew2; + private static Date extractedTime; + + private static String pointerTableName = TableName.pointerTableName;//记录去重进行到哪个地方了 + private static String sourceTableName = TableName.gatherProjectsTableName;//去重项目的来源,来自汇总程序生成的表 + private static String targetTableName = TableName.eddRelationTableName;//去重结果存储位置 + private static String eddRelationTableName = TableName.eddRelationTableName; + private static int batchSize = 500; + //读指针 + public int readPointer(String table, String source, String target, int minId){ + int pointer = minId; + try { + pointer = dbSource.getPointer(table, source, target); + } catch(Exception e) { + logger.info("No such pointer! Create one"); + dbSource.insertPointer(table, source, target, pointer); + } + return pointer; + } + + public void start(){ + logger.info("start remove projects!"); + int count=0; + int tmpCount = 0; + count = readPointer(pointerTableName,sourceTableName,targetTableName, count);//指针表计量处理的项目数 + long start = System.currentTimeMillis(); + while(true){ + //取完别名 + List gpmList = gatherDao.selectGPMListNew(sourceTableName,1,batchSize); + if(gpmList.size()==0){ + logger.info("Remove Duplicate Projects:all projects have been done. Sleep 10 mins"); + logger.warn("deal with "+tmpCount+" projects cost: "+(float)(System.currentTimeMillis() - start)/60000+" minutes"); + try { + Thread.sleep(600*1000); + } catch (InterruptedException e) { + e.printStackTrace(); + } + } + else{ + for(GatherProjectsModel model:gpmList){ + int handleCount = 0; + logger.info("Duplicate remove:handling project " + model.getId()); + //handleCount = mergeProjectNew2.handleNewProject(model,false); + String relationStr = "," + model.getId() + ","; + dbSource.insertEddRelations(eddRelationTableName,relationStr,0); + gatherDao.updateMark(sourceTableName, 2,model.getId()); + count++; + tmpCount++; + + } + dbSource.updatePointer(pointerTableName, sourceTableName, targetTableName, count); + if(tmpCount%10000==0) + logger.warn("deal with:"+tmpCount+" projects cost: "+(float)(System.currentTimeMillis() - start)/60000+" minutes"); + } + } + + } + + + public static void main(String[] args){ + ApplicationContext applicationContext = new ClassPathXmlApplicationContext("classpath:/applicationContext*.xml"); + MergeProjectsForGithub Main = applicationContext.getBean(MergeProjectsForGithub.class); + Main.start(); + } +} diff --git a/project_match/src/main/java/com/ossean/MultiThreadTransferProjects.java b/project_match/src/main/java/com/ossean/MultiThreadTransferProjects.java new file mode 100644 index 000000000..e2fd342d6 --- /dev/null +++ b/project_match/src/main/java/com/ossean/MultiThreadTransferProjects.java @@ -0,0 +1,107 @@ +package com.ossean; + +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +import javax.annotation.Resource; + +import org.apache.log4j.Logger; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.ApplicationContext; +import org.springframework.context.annotation.Scope; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.stereotype.Component; + +import com.ossean.dao.DBDest; +import com.ossean.dao.DBSource; +import com.ossean.dao.TransferPrjDao; +import com.ossean.model.EddRelations; +import com.ossean.util.TransferProjectsUtil; +import com.ossean.util.TransferProjectsUtil2; + + +@Component("multiThreadTransferProjects") +@Scope("prototype") +public class MultiThreadTransferProjects implements Runnable{ + Logger logger = Logger.getLogger(this.getClass()); + + @Resource + private DBSource dbSource; + @Resource + private DBDest dbDest; + @Resource + private TransferPrjDao transferPrjDao; + @Qualifier("transferProjectsUtil2") + @Autowired + private TransferProjectsUtil2 transferProjectsUtil2; + + + private String pointerTableName = TableName.pointerTableName; + private String sourceTableName = TableName.eddRelationTableName; + private String targetTableName = TableName.openSourceProjectsTableName; + private String taggingTableName = TableName.taggingsTableName; + + private int batchSize = 500; + private int startId; + private int endId; + private String threadName; + public void setParameters(int startId,int endId,String threadName) { + this.startId = startId; + this.endId = endId; + this.threadName = threadName; + } + + + @Override + public void run(){ + + logger.info(threadName+" is begin transfer projects~"); + Thread.currentThread().setName(threadName); + int beginId = startId; + while(beginId < endId){ + List eddRelationList = transferPrjDao.getEddRelationListInterval(sourceTableName, beginId,endId, batchSize); + for(EddRelations relation:eddRelationList){ + boolean isUpdate = false; + String gather_projects_ids = relation.getGather_projects_ids(); + try { + gather_projects_ids = gather_projects_ids.substring(1, gather_projects_ids.length() - 1); + } catch (Exception e1) { + System.out.println(relation.getGather_projects_ids()); + e1.printStackTrace(); + } + String[] idsArray = gather_projects_ids.split(","); + for(int i = 0; i < idsArray.length; i++){ + int id = Integer.parseInt(idsArray[i]); + if(null != dbDest.selectOpenSourceProjectsItem(targetTableName,id)){ + isUpdate = true; + break; + } + } + int prjId = transferProjectsUtil2.handleOneRelation(relation,isUpdate); + //System.out.println(prjId); + } + beginId = eddRelationList.get(eddRelationList.size()-1).getId(); + } + TransferManageProcess.gatherState.put(threadName, false); + logger.info(threadName+" is stop"); + } + + public static String getTargetTable(int ospId){ + String targetTableName = ""; + if(ospId >= 770000){ + targetTableName = "relative_memo_to_open_source_projects_70"; + } + else{ + int a = 1 + ospId/11000; + targetTableName = "relative_memo_to_open_source_projects_" + a; + } + + return targetTableName; + } + + + + +} diff --git a/project_match/src/main/java/com/ossean/TransferManageProcess.java b/project_match/src/main/java/com/ossean/TransferManageProcess.java new file mode 100644 index 000000000..d264598cb --- /dev/null +++ b/project_match/src/main/java/com/ossean/TransferManageProcess.java @@ -0,0 +1,82 @@ +package com.ossean; + +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.ThreadPoolExecutor; + +import javax.annotation.Resource; +import org.apache.log4j.Logger; +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.stereotype.Component; + +import com.ossean.dao.TransferPrjDao; + +@Component +public class TransferManageProcess { + + @Resource + private TransferPrjDao transferPrjDao; + private String sourceTableName = TableName.eddRelationTableName; + private static Logger logger = Logger.getLogger(TransferManageProcess.class); + public static Map gatherState = new HashMap(); // + private ExecutorService pool = Executors.newFixedThreadPool(20); + int batchSize = 500000; + public void start(){ + ApplicationContext applicationContext = new ClassPathXmlApplicationContext("classpath:/applicationContext*.xml"); + long startTime = System.currentTimeMillis(); + int maxId = transferPrjDao.selectMaxRelationId(sourceTableName); + int loop; + if (maxId % batchSize == 0) + loop = maxId/batchSize; + else + loop = maxId/batchSize+1; + //每个线程处理150000个项目 + for(int i = 0;i < loop;i++){ + Boolean state = gatherState.get("PrjThread"+(i+1)); + if(state == null || state == false){ + gatherState.put("PrjThread"+(i+1), true); + } + else if(state == true){ + continue; + } + MultiThreadTransferProjects thread = (MultiThreadTransferProjects)applicationContext.getBean("multiThreadTransferProjects"); + + if(i==(loop-1)) + thread.setParameters(i*batchSize,maxId,"PrjThread"+(i+1)); + else + thread.setParameters(i*batchSize,(i+1)*batchSize,"PrjThread"+(i+1)); + pool.execute(thread); + } + logger.info("current env thread nums: "+((ThreadPoolExecutor)pool).getActiveCount()); + System.out.println("已经开启所有的子线程"); + pool.shutdown(); + System.out.println("shutdown():启动一次顺序关闭,执行以前提交的任务,但不接受新任务。"); + while(true){ + if(pool.isTerminated()){ + System.out.println("所有的子线程都结束了!"); + System.out.println("cost time: "+(System.currentTimeMillis()-startTime)/1000+"seconds"); + break; + } + try { + logger.info(".......sleeping......"); + Thread.sleep(1*60*1000); + } catch (InterruptedException e) { + e.printStackTrace(); + logger.error(e); + } + } + System.gc(); //手动垃圾回收 + + } + + public static void main(String[] args) throws InterruptedException { + + ApplicationContext applicationContext = new ClassPathXmlApplicationContext("classpath:/applicationContext*.xml"); + TransferManageProcess main = applicationContext.getBean(TransferManageProcess.class); + main.start(); + } + +} diff --git a/project_match/src/main/java/com/ossean/dao/TransferPrjDao.java b/project_match/src/main/java/com/ossean/dao/TransferPrjDao.java index 35be10d6d..3c8d7ea15 100644 --- a/project_match/src/main/java/com/ossean/dao/TransferPrjDao.java +++ b/project_match/src/main/java/com/ossean/dao/TransferPrjDao.java @@ -15,6 +15,12 @@ public interface TransferPrjDao { @Select("select * from ${table} where id>=#{start} and flag =0 limit #{size}") public List getEddRelationList(@Param("table") String table, @Param("start") int start, @Param("size") int size); + @Select("select * from ${table} where id>#{start} and id <= #{end} and flag =0 limit #{size}") + public List getEddRelationListInterval(@Param("table") String table, @Param("start") int start,@Param("end") int end, @Param("size") int size); + + @Select("select max(id) from ${table}") + public int selectMaxRelationId(@Param("table") String table); + //更新edd_relations表的osp_id字段 @Update("update ${table} set osp_id=#{osp_id} where id=#{id}") public void updateEddRelationOspId(@Param("table") String table, @Param("id") int id, @Param("osp_id") int osp_id);