ConversionBackfillJobService.java 2.4 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667
  1. package com.adx.tencent.conversionsync;
  2. import org.slf4j.Logger;
  3. import org.slf4j.LoggerFactory;
  4. import java.time.Instant;
  5. import java.util.List;
  6. import java.util.Map;
  7. import java.util.UUID;
  8. import java.util.concurrent.*;
  9. import java.util.concurrent.atomic.AtomicReference;
  10. /**
  11. * 后台补拉任务服务。
  12. * 用于把 backfill 从同步 HTTP 请求改成异步 job。
  13. */
  14. public class ConversionBackfillJobService {
  15. private static final Logger log = LoggerFactory.getLogger(ConversionBackfillJobService.class);
  16. private final ConversionSyncRunner runner;
  17. private final ExecutorService executor;
  18. private final Map<String, JobState> jobs = new ConcurrentHashMap<>();
  19. public ConversionBackfillJobService(ConversionSyncRunner runner) {
  20. this.runner = runner;
  21. this.executor = Executors.newFixedThreadPool(2, r -> {
  22. Thread t = new Thread(r, "conversion-backfill-job");
  23. t.setDaemon(true);
  24. return t;
  25. });
  26. }
  27. public JobHandle submit(List<Integer> offsets) {
  28. String jobId = UUID.randomUUID().toString().replace("-", "");
  29. JobState state = new JobState(jobId, "QUEUED", Instant.now(), null, null);
  30. jobs.put(jobId, state);
  31. executor.submit(() -> {
  32. update(jobId, "RUNNING", null, null);
  33. try {
  34. ConversionSyncService.SyncResult result = runner.runBackfillForOffsetsParallel(offsets, 3);
  35. update(jobId, "SUCCEEDED", result, null);
  36. } catch (Exception e) {
  37. log.error("[ConversionBackfill] job failed: jobId={}, error={}", jobId, e.getMessage(), e);
  38. update(jobId, "FAILED", null, e.getMessage());
  39. }
  40. });
  41. return new JobHandle(jobId, state.status(), state.createdAt());
  42. }
  43. public JobState getJob(String jobId) {
  44. return jobs.get(jobId);
  45. }
  46. private void update(String jobId, String status, ConversionSyncService.SyncResult result, String error) {
  47. JobState current = jobs.get(jobId);
  48. if (current == null) return;
  49. jobs.put(jobId, new JobState(jobId, status, current.createdAt(), result, error));
  50. }
  51. public record JobHandle(String jobId, String status, Instant createdAt) {}
  52. public record JobState(String jobId, String status, Instant createdAt,
  53. ConversionSyncService.SyncResult result, String errorMessage) {}
  54. }