ConversionSyncRunner.java 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109
  1. package com.adx.tencent.conversionsync;
  2. import com.adx.tencent.baidu.model.ConversionQuery;
  3. import org.slf4j.Logger;
  4. import org.slf4j.LoggerFactory;
  5. import java.time.LocalDate;
  6. import java.time.format.DateTimeFormatter;
  7. import java.util.ArrayList;
  8. import java.util.List;
  9. import java.util.concurrent.*;
  10. /**
  11. * 转化同步调度器。
  12. * 定期调用 ConversionSyncService.syncTencentConversions()。
  13. */
  14. public class ConversionSyncRunner {
  15. private static final Logger log = LoggerFactory.getLogger(ConversionSyncRunner.class);
  16. private static final DateTimeFormatter DATE_FMT = DateTimeFormatter.ofPattern("yyyyMMdd");
  17. private final ConversionSyncService syncer;
  18. private final int dateOffsetDays;
  19. private final int pageSize;
  20. private final List<Integer> acts;
  21. private final ExecutorService dateExecutor;
  22. public ConversionSyncRunner(ConversionSyncService syncer,
  23. int dateOffsetDays, int pageSize, List<Integer> acts) {
  24. this.syncer = syncer;
  25. this.dateOffsetDays = dateOffsetDays;
  26. this.pageSize = pageSize > 0 ? pageSize : 1;
  27. this.acts = acts;
  28. this.dateExecutor = Executors.newFixedThreadPool(3, r -> {
  29. Thread t = new Thread(r, "conv-backfill-date");
  30. t.setDaemon(true);
  31. return t;
  32. });
  33. }
  34. public ConversionSyncService.SyncResult runOnce() throws Exception {
  35. return runForOffsets(List.of(dateOffsetDays));
  36. }
  37. public ConversionSyncService.SyncResult runForOffsets(List<Integer> offsets) throws Exception {
  38. return runForOffsetsParallel(offsets, 1);
  39. }
  40. public ConversionSyncService.SyncResult runForOffsetsParallel(List<Integer> offsets, int maxConcurrency) throws Exception {
  41. ConversionSyncService.SyncResult total = new ConversionSyncService.SyncResult();
  42. if (offsets == null || offsets.isEmpty()) {
  43. return total;
  44. }
  45. List<Integer> normalized = new ArrayList<>(offsets.size());
  46. for (Integer offset : offsets) {
  47. if (offset != null) {
  48. normalized.add(offset);
  49. }
  50. }
  51. int concurrency = Math.max(1, Math.min(maxConcurrency, normalized.size()));
  52. ExecutorService executor = concurrency == 1 ? null : this.dateExecutor;
  53. if (executor == null) {
  54. for (Integer offset : normalized) {
  55. String date = LocalDate.now().plusDays(offset).format(DATE_FMT);
  56. ConversionSyncService.SyncResult result = runForDate(date);
  57. merge(total, result);
  58. }
  59. return total;
  60. }
  61. List<CompletableFuture<ConversionSyncService.SyncResult>> futures = new ArrayList<>(normalized.size());
  62. for (Integer offset : normalized) {
  63. String date = LocalDate.now().plusDays(offset).format(DATE_FMT);
  64. futures.add(CompletableFuture.supplyAsync(() -> {
  65. try {
  66. return runForDate(date);
  67. } catch (Exception e) {
  68. throw new CompletionException(e);
  69. }
  70. }, executor));
  71. }
  72. CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
  73. for (CompletableFuture<ConversionSyncService.SyncResult> future : futures) {
  74. merge(total, future.join());
  75. }
  76. return total;
  77. }
  78. private ConversionSyncService.SyncResult runForDate(String date) throws Exception {
  79. ConversionQuery query = new ConversionQuery();
  80. query.setDate(date);
  81. query.setPageSize(pageSize);
  82. query.setActs(acts);
  83. return syncer.syncTencentConversions(query);
  84. }
  85. private static void merge(ConversionSyncService.SyncResult total, ConversionSyncService.SyncResult part) {
  86. if (total == null || part == null) return;
  87. total.fetched.addAndGet(part.fetched.get());
  88. total.matched.addAndGet(part.matched.get());
  89. total.sent.addAndGet(part.sent.get());
  90. total.skipped.addAndGet(part.skipped.get());
  91. total.failed.addAndGet(part.failed.get());
  92. total.alreadySent.addAndGet(part.alreadySent.get());
  93. }
  94. }