Pārlūkot izejas kodu

feat:redis store optimize

zhaohaipeng 4 nedēļas atpakaļ
vecāks
revīzija
fcf8bb6b6c

+ 2 - 8
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.Map;
-import java.util.Objects;
 
 /**
  * @author dyp
@@ -59,8 +58,6 @@ public class ODPSToRedis {
             return;
         }
 
-        DTSConfig dtsConfigV2 = dtsConfigService.getV2DTSConfig(argMap);
-
         long count = 0;
         int retry = NumberUtils.toInt(argMap.getOrDefault("retry", "50"), 50);
         while (count <= 0 && retry-- >= 0) {
@@ -126,12 +123,9 @@ public class ODPSToRedis {
         int finalPartationNum = Math.min(maxPartitionNum, calcPartationNum);
 
         log.info("argMap.maxPartitionNum: {} calcPartationNum: {}, finalPartationNum: {}", maxPartitionNum, calcPartationNum, finalPartationNum);
-        log.info("config: {}, dtsConfigV2: {}", config, dtsConfigV2);
+        log.info("config: {}", config);
         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));
-        }
+        repartitionedFieldValues.foreachPartition(iterator -> redisService.mSetV2(iterator, config, false));
 
         redisService.setMonitor(config, argMap);
 

+ 9 - 29
recommend-feature-service/src/main/java/com/tzld/piaoquan/recommend/feature/service/FeatureV2Service.java

@@ -32,26 +32,17 @@ import java.util.stream.Collectors;
 @Service
 @Slf4j
 public class FeatureV2Service {
-    @Qualifier("redisTemplate")
-    @Autowired
-    private RedisTemplate<String, String> redisTemplate;
 
     @Qualifier("byteRedisTemplate")
     @Autowired
     private RedisTemplate<String, byte[]> byteRedisTemplate;
 
     @ApolloJsonValue("${dts.config:}")
-    private List<DTSConfig> newDtsConfigs;
-
-    @ApolloJsonValue("${dts.config.v2:}")
-    private List<DTSConfig> dtsConfigsV2;
+    private List<DTSConfig> dtsConfigs;
 
     @Value("${feature.mget.batch.size:300}")
     private Integer batchSize;
 
-    @Value("${feature.service.optimize.switch:false}")
-    private Boolean optimizeV3Switch;
-
     public MultiGetFeatureResponse multiGetFeature(MultiGetFeatureRequest request) {
         int keyCount = request.getFeatureKeyCount();
         if (keyCount == 0) {
@@ -62,17 +53,10 @@ public class FeatureV2Service {
 
         long startTime = System.currentTimeMillis();
 
-        List<String> redisKeys;
-        if (Boolean.TRUE.equals(optimizeV3Switch)) {
-            redisKeys = request.getFeatureKeyList().stream()
-                    .map(this::buildRedisKeyV2)
-                    .collect(Collectors.toList());
-        } else {
-            redisKeys = request.getFeatureKeyList().stream()
-                    .map(this::redisKey)
-                    .map(s -> String.format("%s:v2", s))
-                    .collect(Collectors.toList());
-        }
+        List<String> redisKeys = request.getFeatureKeyList()
+                .stream()
+                .map(this::redisKey)
+                .collect(Collectors.toList());
 
         int safeBatchSize = (batchSize == null || batchSize <= 0) ? 300 : batchSize;
 
@@ -106,9 +90,9 @@ public class FeatureV2Service {
                 }
             } catch (TimeoutException e) {
                 bf.future.cancel(true);
-                log.warn("Batch mGetV2 timeout, startIndex={}, size={}", bf.startIndex, bf.size);
+                log.warn("Batch mGet timeout, startIndex={}, size={}", bf.startIndex, bf.size);
             } catch (Exception e) {
-                log.error("Batch mGetV2 failed, startIndex={}, size={}", bf.startIndex, bf.size, e);
+                log.error("Batch mGet failed, startIndex={}, size={}", bf.startIndex, bf.size, e);
             }
         }
 
@@ -140,15 +124,11 @@ public class FeatureV2Service {
     }
 
     private String redisKey(FeatureKeyProto fk) {
-        return this.buildRedisKey(fk, newDtsConfigs);
-    }
-
-    private String buildRedisKeyV2(FeatureKeyProto fk) {
-        return this.buildRedisKey(fk, dtsConfigsV2);
+        return this.buildRedisKey(fk, dtsConfigs);
     }
 
     // Note:写入和读取的key生成规则应保持一致
-    private String buildRedisKey(FeatureKeyProto fk, List<DTSConfig> dtsConfigs){
+    private String buildRedisKey(FeatureKeyProto fk, List<DTSConfig> dtsConfigs) {
         Optional<DTSConfig> optional = dtsConfigs.stream()
                 .filter(c -> c.getOdps() != null && StringUtils.equals(c.getOdps().getTable(), fk.getTableName()))
                 .findFirst();