Browse Source

删除更新时间次序更新

xufei 4 years ago
parent
commit
74af06b990
1 changed files with 3 additions and 1 deletions
  1. 3 1
      src/main/java/com/winhc/task/PersonMergeIncremnetTask.java

+ 3 - 1
src/main/java/com/winhc/task/PersonMergeIncremnetTask.java

@@ -86,8 +86,10 @@ public class PersonMergeIncremnetTask {
 
 
     private void sendKafka(MergeParam param) throws InterruptedException, ExecutionException, TimeoutException {
     private void sendKafka(MergeParam param) throws InterruptedException, ExecutionException, TimeoutException {
         //加载文件发送kafka
         //加载文件发送kafka
-        loadCSVSendKafka(param.getPathPre(), param.getMergePath(), param.getTopic(), "0");
         loadCSVSendKafka(param.getPathPre(), param.getDeletedPath(), param.getTopic(), "1");
         loadCSVSendKafka(param.getPathPre(), param.getDeletedPath(), param.getTopic(), "1");
+        //保证先移除,再更新
+        Thread.sleep(10*1000);
+        loadCSVSendKafka(param.getPathPre(), param.getMergePath(), param.getTopic(), "0");
     }
     }
 
 
     private void mergeAndExport(MergeParam param) {
     private void mergeAndExport(MergeParam param) {