ColdWorker.java 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111
  1. package com.adx.tencent.worker;
  2. import com.adx.tencent.storage.RedisHotStore;
  3. import com.adx.tencent.storage.TiDBColdStore;
  4. import com.adx.tencent.storage.model.BidRecord;
  5. import com.adx.tencent.storage.model.QueuedEvent;
  6. import com.adx.tencent.storage.model.TrackingRecord;
  7. import com.fasterxml.jackson.databind.ObjectMapper;
  8. import org.slf4j.Logger;
  9. import org.slf4j.LoggerFactory;
  10. import java.time.Duration;
  11. import java.util.ArrayList;
  12. import java.util.List;
  13. /**
  14. * 对应 Go internal/worker/worker.go 的 Worker。
  15. * 从 Redis Stream 读取事件,逐条写入 TiDB(冷路径)。
  16. *
  17. * 循环逻辑:
  18. * 1. EnsureGroup(幂等创建消费者组)
  19. * 2. ClaimStale(领取超时未 ACK 的消息)
  20. * 3. 若无 stale,则 Read 新消息
  21. * 4. 逐条处理(bid → saveBid,tracking → saveTracking)
  22. * 5. Ack + Trim
  23. */
  24. public class ColdWorker {
  25. private static final Logger log = LoggerFactory.getLogger(ColdWorker.class);
  26. private final RedisHotStore hotStore;
  27. private final TiDBColdStore coldStore;
  28. private final String group;
  29. private final String consumer;
  30. private final int batch;
  31. private final Duration pendingIdle;
  32. private final long streamMaxLen;
  33. private final ObjectMapper objectMapper;
  34. public ColdWorker(RedisHotStore hotStore, TiDBColdStore coldStore,
  35. String group, String consumer, int batch,
  36. Duration pendingIdle, long streamMaxLen,
  37. ObjectMapper objectMapper) {
  38. this.hotStore = hotStore;
  39. this.coldStore = coldStore;
  40. this.group = (group == null || group.isBlank()) ? "cold-writers" : group;
  41. this.consumer = (consumer == null || consumer.isBlank()) ? "worker-1" : consumer;
  42. this.batch = batch <= 0 ? 100 : batch;
  43. this.pendingIdle = (pendingIdle == null || pendingIdle.isZero()) ? Duration.ofMinutes(2) : pendingIdle;
  44. this.streamMaxLen = streamMaxLen;
  45. this.objectMapper = objectMapper;
  46. }
  47. /**
  48. * 单次运行:对应 Go Worker.RunOnce()。
  49. * 返回本次处理的事件数量。
  50. *
  51. * 失败语义(严格对齐 Go):
  52. * • 任一 handleEvent 报错 → 不 Ack 任何 event(包括已成功的),整批原子性。
  53. * • 异常必须向上抛出,由调用方(BackgroundTasks) sleep 冷却,
  54. * 否则 TiDB 持续写入失败时会陷入 busy loop。
  55. */
  56. public int runOnce() throws Exception {
  57. hotStore.ensureGroup(group);
  58. List<QueuedEvent> events = hotStore.claimStale(group, consumer, pendingIdle, batch);
  59. if (events.isEmpty()) {
  60. events = hotStore.read(group, consumer, batch);
  61. }
  62. if (events.isEmpty()) return 0;
  63. List<String> acked = new ArrayList<>();
  64. Exception failed = null;
  65. for (QueuedEvent event : events) {
  66. try {
  67. handleEvent(event);
  68. } catch (Exception e) {
  69. log.error("handle event {}: {}", event.getId(), e.getMessage(), e);
  70. failed = e;
  71. break;
  72. }
  73. acked.add(event.getId());
  74. }
  75. // 对齐 Go:失败时直接抛出,不 Ack 任何 event(包括已成功的)。
  76. // 失败的事件会在 pendingIdle 后被 ClaimStale 重新领取,TiDB 写入是幂等的。
  77. if (failed != null) {
  78. throw failed;
  79. }
  80. hotStore.ack(group, acked);
  81. if (streamMaxLen > 0) {
  82. hotStore.trim(streamMaxLen);
  83. }
  84. return events.size();
  85. }
  86. private void handleEvent(QueuedEvent event) throws Exception {
  87. switch (event.getType()) {
  88. case "bid" -> {
  89. BidRecord record = objectMapper.readValue(event.getPayload(), BidRecord.class);
  90. coldStore.saveBid(record);
  91. }
  92. case "tracking" -> {
  93. TrackingRecord record = objectMapper.readValue(event.getPayload(), TrackingRecord.class);
  94. coldStore.saveTracking(record);
  95. }
  96. default -> throw new IllegalArgumentException("unknown event type: " + event.getType());
  97. }
  98. }
  99. }