package com.adx.tencent.conversionsync; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.time.Instant; import java.util.List; import java.util.Map; import java.util.UUID; import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicReference; /** * 后台补拉任务服务。 * 用于把 backfill 从同步 HTTP 请求改成异步 job。 */ public class ConversionBackfillJobService { private static final Logger log = LoggerFactory.getLogger(ConversionBackfillJobService.class); private final ConversionSyncRunner runner; private final ExecutorService executor; private final Map jobs = new ConcurrentHashMap<>(); public ConversionBackfillJobService(ConversionSyncRunner runner) { this.runner = runner; this.executor = Executors.newFixedThreadPool(2, r -> { Thread t = new Thread(r, "conversion-backfill-job"); t.setDaemon(true); return t; }); } public JobHandle submit(List offsets) { String jobId = UUID.randomUUID().toString().replace("-", ""); JobState state = new JobState(jobId, "QUEUED", Instant.now(), null, null); jobs.put(jobId, state); executor.submit(() -> { update(jobId, "RUNNING", null, null); try { ConversionSyncService.SyncResult result = runner.runBackfillForOffsetsParallel(offsets, 3); update(jobId, "SUCCEEDED", result, null); } catch (Exception e) { log.error("[ConversionBackfill] job failed: jobId={}, error={}", jobId, e.getMessage(), e); update(jobId, "FAILED", null, e.getMessage()); } }); return new JobHandle(jobId, state.status(), state.createdAt()); } public JobState getJob(String jobId) { return jobs.get(jobId); } private void update(String jobId, String status, ConversionSyncService.SyncResult result, String error) { JobState current = jobs.get(jobId); if (current == null) return; jobs.put(jobId, new JobState(jobId, status, current.createdAt(), result, error)); } public record JobHandle(String jobId, String status, Instant createdAt) {} public record JobState(String jobId, String status, Instant createdAt, ConversionSyncService.SyncResult result, String errorMessage) {} }