|
@@ -13,6 +13,8 @@ import org.springframework.data.redis.connection.stream.RecordId;
|
|
|
import org.springframework.data.redis.connection.stream.StreamOffset;
|
|
import org.springframework.data.redis.connection.stream.StreamOffset;
|
|
|
import org.springframework.data.redis.connection.stream.StreamReadOptions;
|
|
import org.springframework.data.redis.connection.stream.StreamReadOptions;
|
|
|
import org.springframework.data.redis.core.StringRedisTemplate;
|
|
import org.springframework.data.redis.core.StringRedisTemplate;
|
|
|
|
|
+import org.slf4j.Logger;
|
|
|
|
|
+import org.slf4j.LoggerFactory;
|
|
|
|
|
|
|
|
import java.nio.charset.StandardCharsets;
|
|
import java.nio.charset.StandardCharsets;
|
|
|
import java.time.Duration;
|
|
import java.time.Duration;
|
|
@@ -25,15 +27,19 @@ import java.util.LinkedHashSet;
|
|
|
import java.util.List;
|
|
import java.util.List;
|
|
|
import java.util.Map;
|
|
import java.util.Map;
|
|
|
import java.util.Set;
|
|
import java.util.Set;
|
|
|
|
|
+import java.util.concurrent.ConcurrentHashMap;
|
|
|
|
|
|
|
|
public class HonorHotStore {
|
|
public class HonorHotStore {
|
|
|
|
|
|
|
|
|
|
+ private static final Logger log = LoggerFactory.getLogger(HonorHotStore.class);
|
|
|
|
|
+
|
|
|
private final StringRedisTemplate redis;
|
|
private final StringRedisTemplate redis;
|
|
|
private final ObjectMapper objectMapper;
|
|
private final ObjectMapper objectMapper;
|
|
|
private final HonorColdStore coldStore;
|
|
private final HonorColdStore coldStore;
|
|
|
private final String prefix;
|
|
private final String prefix;
|
|
|
private final String stream;
|
|
private final String stream;
|
|
|
private final Duration bidTtl;
|
|
private final Duration bidTtl;
|
|
|
|
|
+ private final Set<String> initializedGroups = ConcurrentHashMap.newKeySet();
|
|
|
|
|
|
|
|
public HonorHotStore(StringRedisTemplate redis, ObjectMapper objectMapper,
|
|
public HonorHotStore(StringRedisTemplate redis, ObjectMapper objectMapper,
|
|
|
HonorColdStore coldStore, String prefix, String stream, Duration bidTtl) {
|
|
HonorColdStore coldStore, String prefix, String stream, Duration bidTtl) {
|
|
@@ -151,17 +157,28 @@ public class HonorHotStore {
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
public void ensureGroup(String group) {
|
|
public void ensureGroup(String group) {
|
|
|
|
|
+ if (group == null || group.isBlank()) return;
|
|
|
|
|
+ if (initializedGroups.contains(group)) return;
|
|
|
try {
|
|
try {
|
|
|
redis.opsForStream().createGroup(stream, ReadOffset.from("0"), group);
|
|
redis.opsForStream().createGroup(stream, ReadOffset.from("0"), group);
|
|
|
|
|
+ initializedGroups.add(group);
|
|
|
} catch (Exception e) {
|
|
} catch (Exception e) {
|
|
|
- if (isBusyGroupError(e)) return;
|
|
|
|
|
|
|
+ if (isBusyGroupError(e)) {
|
|
|
|
|
+ initializedGroups.add(group);
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
try {
|
|
try {
|
|
|
var id = redis.opsForStream().add(org.springframework.data.redis.connection.stream.StreamRecords
|
|
var id = redis.opsForStream().add(org.springframework.data.redis.connection.stream.StreamRecords
|
|
|
.newRecord().ofMap(Collections.singletonMap("_init", "1")).withStreamKey(stream));
|
|
.newRecord().ofMap(Collections.singletonMap("_init", "1")).withStreamKey(stream));
|
|
|
if (id != null) redis.opsForStream().delete(stream, id);
|
|
if (id != null) redis.opsForStream().delete(stream, id);
|
|
|
redis.opsForStream().createGroup(stream, ReadOffset.from("0"), group);
|
|
redis.opsForStream().createGroup(stream, ReadOffset.from("0"), group);
|
|
|
|
|
+ initializedGroups.add(group);
|
|
|
} catch (Exception retryEx) {
|
|
} catch (Exception retryEx) {
|
|
|
- if (!isBusyGroupError(retryEx)) throw new RuntimeException("ensure honor group failed", retryEx);
|
|
|
|
|
|
|
+ if (isBusyGroupError(retryEx)) {
|
|
|
|
|
+ initializedGroups.add(group);
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+ throw new RuntimeException("ensure honor group failed", retryEx);
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
@@ -176,6 +193,8 @@ public class HonorHotStore {
|
|
|
return toQueuedEvents(records);
|
|
return toQueuedEvents(records);
|
|
|
} catch (Exception e) {
|
|
} catch (Exception e) {
|
|
|
if (e.getMessage() != null && e.getMessage().contains("NOGROUP")) return Collections.emptyList();
|
|
if (e.getMessage() != null && e.getMessage().contains("NOGROUP")) return Collections.emptyList();
|
|
|
|
|
+ log.error("[HonorHotStore] read failed | stream={} | group={} | consumer={} | count={} | error={}",
|
|
|
|
|
+ stream, group, consumer, count, rootMessage(e), e);
|
|
|
throw new RuntimeException("read honor stream failed", e);
|
|
throw new RuntimeException("read honor stream failed", e);
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
@@ -243,4 +262,16 @@ public class HonorHotStore {
|
|
|
}
|
|
}
|
|
|
return false;
|
|
return false;
|
|
|
}
|
|
}
|
|
|
|
|
+
|
|
|
|
|
+ private static String rootMessage(Throwable e) {
|
|
|
|
|
+ Throwable cur = e;
|
|
|
|
|
+ String msg = null;
|
|
|
|
|
+ while (cur != null) {
|
|
|
|
|
+ if (cur.getMessage() != null && !cur.getMessage().isBlank()) {
|
|
|
|
|
+ msg = cur.getMessage();
|
|
|
|
|
+ }
|
|
|
|
|
+ cur = cur.getCause();
|
|
|
|
|
+ }
|
|
|
|
|
+ return msg == null ? "" : msg;
|
|
|
|
|
+ }
|
|
|
}
|
|
}
|