yumeng 1 ヶ月 前
コミット
57bab1f731

+ 2 - 0
src/main/java/com/adx/tencent/AppConfiguration.java

@@ -39,6 +39,7 @@ import com.zaxxer.hikari.HikariConfig;
 import com.zaxxer.hikari.HikariDataSource;
 import io.lettuce.core.ClientOptions;
 import io.lettuce.core.SocketOptions;
+import io.lettuce.core.protocol.ProtocolVersion;
 import org.apache.ibatis.session.SqlSessionFactory;
 import org.mybatis.spring.SqlSessionFactoryBean;
 import org.slf4j.Logger;
@@ -98,6 +99,7 @@ public class AppConfiguration {
 
         LettuceClientConfiguration clientCfg = LettuceClientConfiguration.builder()
                 .clientOptions(ClientOptions.builder()
+                        .protocolVersion(ProtocolVersion.RESP2)
                         .socketOptions(SocketOptions.builder()
                                 .connectTimeout(props.getRedisConnectTimeout())
                                 .build())

+ 17 - 0
src/main/java/com/adx/tencent/honor/store/HonorHotStore.java

@@ -159,6 +159,10 @@ public class HonorHotStore {
     public void ensureGroup(String group) {
         if (group == null || group.isBlank()) return;
         if (initializedGroups.contains(group)) return;
+        if (groupExists(group)) {
+            initializedGroups.add(group);
+            return;
+        }
         try {
             redis.opsForStream().createGroup(stream, ReadOffset.from("0"), group);
             initializedGroups.add(group);
@@ -171,6 +175,10 @@ public class HonorHotStore {
                 var id = redis.opsForStream().add(org.springframework.data.redis.connection.stream.StreamRecords
                         .newRecord().ofMap(Collections.singletonMap("_init", "1")).withStreamKey(stream));
                 if (id != null) redis.opsForStream().delete(stream, id);
+                if (groupExists(group)) {
+                    initializedGroups.add(group);
+                    return;
+                }
                 redis.opsForStream().createGroup(stream, ReadOffset.from("0"), group);
                 initializedGroups.add(group);
             } catch (Exception retryEx) {
@@ -255,6 +263,15 @@ public class HonorHotStore {
         return result;
     }
 
+    private boolean groupExists(String group) {
+        try {
+            var groups = redis.opsForStream().groups(stream);
+            return groups != null && groups.stream().anyMatch(info -> group.equals(info.groupName()));
+        } catch (Exception e) {
+            return false;
+        }
+    }
+
     private boolean isBusyGroupError(Throwable e) {
         while (e != null) {
             if (e.getMessage() != null && e.getMessage().contains("BUSYGROUP")) return true;

+ 17 - 0
src/main/java/com/adx/tencent/kuaishou/store/KuaishouHotStore.java

@@ -202,6 +202,10 @@ public class KuaishouHotStore {
     public void ensureGroup(String group) {
         if (group == null || group.isBlank()) return;
         if (initializedGroups.contains(group)) return;
+        if (groupExists(group)) {
+            initializedGroups.add(group);
+            return;
+        }
         try {
             redis.opsForStream().createGroup(stream, ReadOffset.from("0"), group);
             initializedGroups.add(group);
@@ -214,6 +218,10 @@ public class KuaishouHotStore {
                 var id = redis.opsForStream().add(org.springframework.data.redis.connection.stream.StreamRecords
                         .newRecord().ofMap(Collections.singletonMap("_init", "1")).withStreamKey(stream));
                 if (id != null) redis.opsForStream().delete(stream, id);
+                if (groupExists(group)) {
+                    initializedGroups.add(group);
+                    return;
+                }
                 redis.opsForStream().createGroup(stream, ReadOffset.from("0"), group);
                 initializedGroups.add(group);
             } catch (Exception retryEx) {
@@ -339,6 +347,15 @@ public class KuaishouHotStore {
         return result;
     }
 
+    private boolean groupExists(String group) {
+        try {
+            var groups = redis.opsForStream().groups(stream);
+            return groups != null && groups.stream().anyMatch(info -> group.equals(info.groupName()));
+        } catch (Exception e) {
+            return false;
+        }
+    }
+
     private boolean isBusyGroupError(Throwable e) {
         while (e != null) {
             if (e.getMessage() != null && e.getMessage().contains("BUSYGROUP")) return true;

+ 17 - 0
src/main/java/com/adx/tencent/storage/RedisHotStore.java

@@ -252,6 +252,10 @@ public class RedisHotStore {
     public void ensureGroup(String group) {
         if (group == null || group.isBlank()) return;
         if (initializedGroups.contains(group)) return;
+        if (groupExists(group)) {
+            initializedGroups.add(group);
+            return;
+        }
         try {
             redis.opsForStream().createGroup(stream, ReadOffset.from("0"), group);
             initializedGroups.add(group);
@@ -269,6 +273,10 @@ public class RedisHotStore {
                 if (id != null) {
                     redis.opsForStream().delete(stream, id);
                 }
+                if (groupExists(group)) {
+                    initializedGroups.add(group);
+                    return;
+                }
                 redis.opsForStream().createGroup(stream, ReadOffset.from("0"), group);
                 initializedGroups.add(group);
             } catch (Exception retryEx) {
@@ -281,6 +289,15 @@ public class RedisHotStore {
         }
     }
 
+    private boolean groupExists(String group) {
+        try {
+            var groups = redis.opsForStream().groups(stream);
+            return groups != null && groups.stream().anyMatch(info -> group.equals(info.groupName()));
+        } catch (Exception e) {
+            return false;
+        }
+    }
+
     private boolean isBusyGroupError(Throwable e) {
         while (e != null) {
             if (e.getMessage() != null && e.getMessage().contains("BUSYGROUP")) {