Jelajahi Sumber

feat:添加v2同步逻辑

zhaohaipeng 1 bulan lalu
induk
melakukan
8eb3b265b9

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

@@ -20,7 +20,6 @@ import org.apache.spark.api.java.JavaSparkContext;
 
 
 import java.util.HashMap;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Map;
-import java.util.Objects;
 
 
 /**
 /**
  * @author dyp
  * @author dyp
@@ -126,10 +125,10 @@ public class ODPSToRedis {
         int finalPartationNum = Math.min(maxPartitionNum, calcPartationNum);
         int finalPartationNum = Math.min(maxPartitionNum, calcPartationNum);
 
 
         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);
         fieldValues.repartition(finalPartationNum).foreachPartition(iterator -> {
         fieldValues.repartition(finalPartationNum).foreachPartition(iterator -> {
-            redisService.mSetV2(iterator, config);
-            redisService.mSetV2(iterator, dtsConfigV2);
+            redisService.mSetV2(iterator, config, false);
+            redisService.mSetV2(iterator, dtsConfigV2, true);
         });
         });
 
 
         redisService.setMonitor(config, argMap);
         redisService.setMonitor(config, argMap);

+ 13 - 7
recommend-feature-produce/src/main/java/com/tzld/piaoquan/recommend/feature/produce/service/RedisService.java

@@ -4,6 +4,7 @@ import com.google.common.reflect.TypeToken;
 import com.tzld.piaoquan.recommend.feature.produce.model.DTSConfig;
 import com.tzld.piaoquan.recommend.feature.produce.model.DTSConfig;
 import com.tzld.piaoquan.recommend.feature.produce.util.CompressionUtil;
 import com.tzld.piaoquan.recommend.feature.produce.util.CompressionUtil;
 import com.tzld.piaoquan.recommend.feature.produce.util.JSONUtils;
 import com.tzld.piaoquan.recommend.feature.produce.util.JSONUtils;
+import lombok.extern.slf4j.Slf4j;
 import org.apache.commons.collections.CollectionUtils;
 import org.apache.commons.collections.CollectionUtils;
 import org.apache.commons.collections.MapUtils;
 import org.apache.commons.collections.MapUtils;
 import org.apache.commons.lang3.StringUtils;
 import org.apache.commons.lang3.StringUtils;
@@ -20,6 +21,7 @@ import java.util.concurrent.TimeUnit;
  *
  *
  * @author dyp
  * @author dyp
  */
  */
+@Slf4j
 public class RedisService implements Serializable {
 public class RedisService implements Serializable {
     private int port = 6379;
     private int port = 6379;
     private String password = "";
     private String password = "";
@@ -35,7 +37,8 @@ public class RedisService implements Serializable {
         }
         }
     }
     }
 
 
-    public void mSetV2(Iterator<Map<String, String>> dataIte, DTSConfig config) {
+    public void mSetV2(Iterator<Map<String, String>> dataIte, DTSConfig config, boolean isV3Optimize) {
+        log.info("mSetV2 param: {}, isV3Optimize: {}", config, isV3Optimize);
         Map<String, String> batch = new HashMap<>();
         Map<String, String> batch = new HashMap<>();
         Jedis jedis = new Jedis(hostName, port);
         Jedis jedis = new Jedis(hostName, port);
         jedis.auth(password);
         jedis.auth(password);
@@ -57,8 +60,8 @@ public class RedisService implements Serializable {
 
 
             String value = JSONUtils.toJson(record);
             String value = JSONUtils.toJson(record);
             batch.put(redisKey, value);
             batch.put(redisKey, value);
-            if (batch.size() >= 500) {
-                mSet(jedis, batch, expire, TimeUnit.SECONDS);
+            if (batch.size() >= 1000) {
+                mSet(jedis, batch, expire, TimeUnit.SECONDS, isV3Optimize);
                 batch.clear();
                 batch.clear();
                 try {
                 try {
                     TimeUnit.MILLISECONDS.sleep(10);
                     TimeUnit.MILLISECONDS.sleep(10);
@@ -68,7 +71,7 @@ public class RedisService implements Serializable {
             }
             }
         }
         }
         if (MapUtils.isNotEmpty(batch)) {
         if (MapUtils.isNotEmpty(batch)) {
-            mSet(jedis, batch, expire, TimeUnit.SECONDS);
+            mSet(jedis, batch, expire, TimeUnit.SECONDS, isV3Optimize);
         }
         }
         jedis.close();
         jedis.close();
     }
     }
@@ -81,14 +84,17 @@ public class RedisService implements Serializable {
         jedis.setex(redisKey, expire, "");
         jedis.setex(redisKey, expire, "");
     }
     }
 
 
-    private void mSet(Jedis jedis, Map<String, String> batch, long expire, TimeUnit timeUnit) {
+    private void mSet(Jedis jedis, Map<String, String> batch, long expire, TimeUnit timeUnit, boolean isV3Optimize) {
         long expireSeconds = timeUnit.toSeconds(expire);
         long expireSeconds = timeUnit.toSeconds(expire);
 
 
         Pipeline pipeline = jedis.pipelined();
         Pipeline pipeline = jedis.pipelined();
         for (Map.Entry<String, String> e : batch.entrySet()) {
         for (Map.Entry<String, String> e : batch.entrySet()) {
             try {
             try {
-                String newKey = String.format("%s:v2", e.getKey());
-                pipeline.setex(newKey.getBytes(StandardCharsets.UTF_8), expireSeconds, CompressionUtil.snappyCompressV2(e.getValue()));
+                String key = e.getKey();
+                if (!isV3Optimize) {
+                    key = String.format("%s:v2", e.getKey());
+                }
+                pipeline.setex(key.getBytes(StandardCharsets.UTF_8), expireSeconds, CompressionUtil.snappyCompressV2(e.getValue()));
             } catch (Exception ignore) {
             } catch (Exception ignore) {
             }
             }
         }
         }