Ver código fonte

feat:添加v2同步逻辑

zhaohaipeng 1 mês atrás
pai
commit
d559e2b74a

+ 6 - 4
recommend-feature-produce/src/main/java/com/tzld/piaoquan/recommend/feature/produce/ODPSToRedis.java

@@ -20,6 +20,7 @@ import org.apache.spark.api.java.JavaSparkContext;
 
 import java.util.HashMap;
 import java.util.Map;
+import java.util.Objects;
 
 /**
  * @author dyp
@@ -126,10 +127,11 @@ public class ODPSToRedis {
 
         log.info("argMap.maxPartitionNum: {} calcPartationNum: {}, finalPartationNum: {}", maxPartitionNum, calcPartationNum, finalPartationNum);
         log.info("config: {}, dtsConfigV2: {}", config, dtsConfigV2);
-        fieldValues.repartition(finalPartationNum).foreachPartition(iterator -> {
-            redisService.mSetV2(iterator, config, false);
-            redisService.mSetV2(iterator, dtsConfigV2, true);
-        });
+        JavaRDD<Map<String, String>> repartitionedFieldValues = fieldValues.repartition(finalPartationNum);
+        repartitionedFieldValues.foreachPartition(iterator -> redisService.mSetV2(iterator, config, false));
+        if (Objects.nonNull(dtsConfigV2) && dtsConfigV2.selfCheck()) {
+            repartitionedFieldValues.foreachPartition(iterator -> redisService.mSetV2(iterator, dtsConfigV2, true));
+        }
 
         redisService.setMonitor(config, argMap);