package com.adx.tencent.worker; import com.adx.tencent.storage.RedisHotStore; import com.adx.tencent.storage.TiDBColdStore; import com.adx.tencent.storage.model.BidRecord; import com.adx.tencent.storage.model.QueuedEvent; import com.adx.tencent.storage.model.TrackingRecord; import com.fasterxml.jackson.databind.ObjectMapper; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.time.Duration; import java.util.ArrayList; import java.util.List; /** * 对应 Go internal/worker/worker.go 的 Worker。 * 从 Redis Stream 读取事件,逐条写入 TiDB(冷路径)。 * * 循环逻辑: * 1. EnsureGroup(幂等创建消费者组) * 2. ClaimStale(领取超时未 ACK 的消息) * 3. 若无 stale,则 Read 新消息 * 4. 逐条处理(bid → saveBid,tracking → saveTracking) * 5. Ack + Trim */ public class ColdWorker { private static final Logger log = LoggerFactory.getLogger(ColdWorker.class); private final RedisHotStore hotStore; private final TiDBColdStore coldStore; private final String group; private final String consumer; private final int batch; private final Duration pendingIdle; private final long streamMaxLen; private final ObjectMapper objectMapper; public ColdWorker(RedisHotStore hotStore, TiDBColdStore coldStore, String group, String consumer, int batch, Duration pendingIdle, long streamMaxLen, ObjectMapper objectMapper) { this.hotStore = hotStore; this.coldStore = coldStore; this.group = (group == null || group.isBlank()) ? "cold-writers" : group; this.consumer = (consumer == null || consumer.isBlank()) ? "worker-1" : consumer; this.batch = batch <= 0 ? 100 : batch; this.pendingIdle = (pendingIdle == null || pendingIdle.isZero()) ? Duration.ofMinutes(2) : pendingIdle; this.streamMaxLen = streamMaxLen; this.objectMapper = objectMapper; } /** * 单次运行:对应 Go Worker.RunOnce()。 * 返回本次处理的事件数量。 * * 失败语义(严格对齐 Go): * • 任一 handleEvent 报错 → 不 Ack 任何 event(包括已成功的),整批原子性。 * • 异常必须向上抛出,由调用方(BackgroundTasks) sleep 冷却, * 否则 TiDB 持续写入失败时会陷入 busy loop。 */ public int runOnce() throws Exception { hotStore.ensureGroup(group); List events = hotStore.claimStale(group, consumer, pendingIdle, batch); if (events.isEmpty()) { events = hotStore.read(group, consumer, batch); } if (events.isEmpty()) return 0; List acked = new ArrayList<>(); Exception failed = null; for (QueuedEvent event : events) { try { handleEvent(event); } catch (Exception e) { log.error("handle event {}: {}", event.getId(), e.getMessage(), e); failed = e; break; } acked.add(event.getId()); } // 对齐 Go:失败时直接抛出,不 Ack 任何 event(包括已成功的)。 // 失败的事件会在 pendingIdle 后被 ClaimStale 重新领取,TiDB 写入是幂等的。 if (failed != null) { throw failed; } hotStore.ack(group, acked); if (streamMaxLen > 0) { hotStore.trim(streamMaxLen); } return events.size(); } private void handleEvent(QueuedEvent event) throws Exception { switch (event.getType()) { case "bid" -> { BidRecord record = objectMapper.readValue(event.getPayload(), BidRecord.class); coldStore.saveBid(record); } case "tracking" -> { TrackingRecord record = objectMapper.readValue(event.getPayload(), TrackingRecord.class); coldStore.saveTracking(record); } default -> throw new IllegalArgumentException("unknown event type: " + event.getType()); } } }