| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111 |
- 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<QueuedEvent> events = hotStore.claimStale(group, consumer, pendingIdle, batch);
- if (events.isEmpty()) {
- events = hotStore.read(group, consumer, batch);
- }
- if (events.isEmpty()) return 0;
- List<String> 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());
- }
- }
- }
|