|
|
@@ -0,0 +1,67 @@
|
|
|
+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<String, JobState> 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<Integer> 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.runForOffsetsParallel(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) {}
|
|
|
+}
|