| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109 |
- package com.adx.tencent.conversionsync;
- import com.adx.tencent.baidu.model.ConversionQuery;
- import org.slf4j.Logger;
- import org.slf4j.LoggerFactory;
- import java.time.LocalDate;
- import java.time.format.DateTimeFormatter;
- import java.util.ArrayList;
- import java.util.List;
- import java.util.concurrent.*;
- /**
- * 转化同步调度器。
- * 定期调用 ConversionSyncService.syncTencentConversions()。
- */
- public class ConversionSyncRunner {
- private static final Logger log = LoggerFactory.getLogger(ConversionSyncRunner.class);
- private static final DateTimeFormatter DATE_FMT = DateTimeFormatter.ofPattern("yyyyMMdd");
- private final ConversionSyncService syncer;
- private final int dateOffsetDays;
- private final int pageSize;
- private final List<Integer> acts;
- private final ExecutorService dateExecutor;
- public ConversionSyncRunner(ConversionSyncService syncer,
- int dateOffsetDays, int pageSize, List<Integer> acts) {
- this.syncer = syncer;
- this.dateOffsetDays = dateOffsetDays;
- this.pageSize = pageSize > 0 ? pageSize : 1;
- this.acts = acts;
- this.dateExecutor = Executors.newFixedThreadPool(3, r -> {
- Thread t = new Thread(r, "conv-backfill-date");
- t.setDaemon(true);
- return t;
- });
- }
- public ConversionSyncService.SyncResult runOnce() throws Exception {
- return runForOffsets(List.of(dateOffsetDays));
- }
- public ConversionSyncService.SyncResult runForOffsets(List<Integer> offsets) throws Exception {
- return runForOffsetsParallel(offsets, 1);
- }
- public ConversionSyncService.SyncResult runForOffsetsParallel(List<Integer> offsets, int maxConcurrency) throws Exception {
- ConversionSyncService.SyncResult total = new ConversionSyncService.SyncResult();
- if (offsets == null || offsets.isEmpty()) {
- return total;
- }
- List<Integer> normalized = new ArrayList<>(offsets.size());
- for (Integer offset : offsets) {
- if (offset != null) {
- normalized.add(offset);
- }
- }
- int concurrency = Math.max(1, Math.min(maxConcurrency, normalized.size()));
- ExecutorService executor = concurrency == 1 ? null : this.dateExecutor;
- if (executor == null) {
- for (Integer offset : normalized) {
- String date = LocalDate.now().plusDays(offset).format(DATE_FMT);
- ConversionSyncService.SyncResult result = runForDate(date);
- merge(total, result);
- }
- return total;
- }
- List<CompletableFuture<ConversionSyncService.SyncResult>> futures = new ArrayList<>(normalized.size());
- for (Integer offset : normalized) {
- String date = LocalDate.now().plusDays(offset).format(DATE_FMT);
- futures.add(CompletableFuture.supplyAsync(() -> {
- try {
- return runForDate(date);
- } catch (Exception e) {
- throw new CompletionException(e);
- }
- }, executor));
- }
- CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
- for (CompletableFuture<ConversionSyncService.SyncResult> future : futures) {
- merge(total, future.join());
- }
- return total;
- }
- private ConversionSyncService.SyncResult runForDate(String date) throws Exception {
- ConversionQuery query = new ConversionQuery();
- query.setDate(date);
- query.setPageSize(pageSize);
- query.setActs(acts);
- return syncer.syncTencentConversions(query);
- }
- private static void merge(ConversionSyncService.SyncResult total, ConversionSyncService.SyncResult part) {
- if (total == null || part == null) return;
- total.fetched.addAndGet(part.fetched.get());
- total.matched.addAndGet(part.matched.get());
- total.sent.addAndGet(part.sent.get());
- total.skipped.addAndGet(part.skipped.get());
- total.failed.addAndGet(part.failed.get());
- total.alreadySent.addAndGet(part.alreadySent.get());
- }
- }
|