|
@@ -128,7 +128,7 @@ public class ODPSToRedis {
|
|
|
log.info("argMap.maxPartitionNum: {} calcPartationNum: {}, finalPartationNum: {}", maxPartitionNum, calcPartationNum, finalPartationNum);
|
|
log.info("argMap.maxPartitionNum: {} calcPartationNum: {}, finalPartationNum: {}", maxPartitionNum, calcPartationNum, finalPartationNum);
|
|
|
log.info("config: {}, dtsConfigV2: {}", config, dtsConfigV2);
|
|
log.info("config: {}, dtsConfigV2: {}", config, dtsConfigV2);
|
|
|
JavaRDD<Map<String, String>> repartitionedFieldValues = fieldValues.repartition(finalPartationNum);
|
|
JavaRDD<Map<String, String>> repartitionedFieldValues = fieldValues.repartition(finalPartationNum);
|
|
|
- repartitionedFieldValues.foreachPartition(iterator -> redisService.mSetV2(iterator, config, false));
|
|
|
|
|
|
|
+ // repartitionedFieldValues.foreachPartition(iterator -> redisService.mSetV2(iterator, config, false));
|
|
|
if (Objects.nonNull(dtsConfigV2) && dtsConfigV2.selfCheck()) {
|
|
if (Objects.nonNull(dtsConfigV2) && dtsConfigV2.selfCheck()) {
|
|
|
repartitionedFieldValues.foreachPartition(iterator -> redisService.mSetV2(iterator, dtsConfigV2, true));
|
|
repartitionedFieldValues.foreachPartition(iterator -> redisService.mSetV2(iterator, dtsConfigV2, true));
|
|
|
}
|
|
}
|