yumeng 1 semana atrás
pai
commit
df29af8f5b

+ 7 - 2
src/main/java/com/adx/tencent/AppConfiguration.java

@@ -405,14 +405,19 @@ public class AppConfiguration {
             }
 
             Runnable job = () -> {
+                long startNs = System.nanoTime();
                 try {
                     taskLog.info("[ConversionBackfill] executing for offsets=1..7");
-                    ConversionSyncService.SyncResult r = conversionSyncRunner.runForOffsets(List.of(-1, -2, -3, -4, -5, -6, -7));
+                    ConversionSyncService.SyncResult r = conversionSyncRunner.runBackfillForOffsetsParallel(
+                            List.of(-1, -2, -3, -4, -5, -6, -7), 3);
+                    long costMs = java.util.concurrent.TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNs);
                     taskLog.info("[ConversionBackfill] done: fetched={} matched={} sent={} failed={} skipped={} alreadySent={}",
                             r.fetched.get(), r.matched.get(), r.sent.get(), r.failed.get(),
                             r.skipped.get(), r.alreadySent.get());
+                    taskLog.info("[ConversionBackfill] total cost={}ms", costMs);
                 } catch (Exception e) {
-                    taskLog.error("[ConversionBackfill] error: {}", e.getMessage(), e);
+                    long costMs = java.util.concurrent.TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNs);
+                    taskLog.error("[ConversionBackfill] error after {}ms: {}", costMs, e.getMessage(), e);
                 }
             };
 

+ 1 - 1
src/main/java/com/adx/tencent/conversionsync/ConversionBackfillJobService.java

@@ -39,7 +39,7 @@ public class ConversionBackfillJobService {
         executor.submit(() -> {
             update(jobId, "RUNNING", null, null);
             try {
-                ConversionSyncService.SyncResult result = runner.runForOffsetsParallel(offsets, 3);
+                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);

+ 15 - 5
src/main/java/com/adx/tencent/conversionsync/ConversionSyncRunner.java

@@ -43,10 +43,19 @@ public class ConversionSyncRunner {
     }
 
     public ConversionSyncService.SyncResult runForOffsets(List<Integer> offsets) throws Exception {
-        return runForOffsetsParallel(offsets, 1);
+        return runForOffsetsParallel(offsets, 1, ConversionSyncService.SyncMode.REALTIME);
     }
 
     public ConversionSyncService.SyncResult runForOffsetsParallel(List<Integer> offsets, int maxConcurrency) throws Exception {
+        return runForOffsetsParallel(offsets, maxConcurrency, ConversionSyncService.SyncMode.REALTIME);
+    }
+
+    public ConversionSyncService.SyncResult runBackfillForOffsetsParallel(List<Integer> offsets, int maxConcurrency) throws Exception {
+        return runForOffsetsParallel(offsets, maxConcurrency, ConversionSyncService.SyncMode.BACKFILL);
+    }
+
+    public ConversionSyncService.SyncResult runForOffsetsParallel(List<Integer> offsets, int maxConcurrency,
+                                                                  ConversionSyncService.SyncMode mode) throws Exception {
         ConversionSyncService.SyncResult total = new ConversionSyncService.SyncResult();
         if (offsets == null || offsets.isEmpty()) {
             return total;
@@ -64,7 +73,7 @@ public class ConversionSyncRunner {
         if (executor == null) {
             for (Integer offset : normalized) {
                 String date = LocalDate.now().plusDays(offset).format(DATE_FMT);
-                ConversionSyncService.SyncResult result = runForDate(date);
+                ConversionSyncService.SyncResult result = runForDate(date, mode);
                 merge(total, result);
             }
             return total;
@@ -75,7 +84,7 @@ public class ConversionSyncRunner {
             String date = LocalDate.now().plusDays(offset).format(DATE_FMT);
             futures.add(CompletableFuture.supplyAsync(() -> {
                 try {
-                    return runForDate(date);
+                    return runForDate(date, mode);
                 } catch (Exception e) {
                     throw new CompletionException(e);
                 }
@@ -89,12 +98,13 @@ public class ConversionSyncRunner {
         return total;
     }
 
-    private ConversionSyncService.SyncResult runForDate(String date) throws Exception {
+    private ConversionSyncService.SyncResult runForDate(String date,
+                                                       ConversionSyncService.SyncMode mode) throws Exception {
         ConversionQuery query = new ConversionQuery();
         query.setDate(date);
         query.setPageSize(pageSize);
         query.setActs(acts);
-        return syncer.syncTencentConversions(query);
+        return syncer.syncTencentConversions(query, mode);
     }
 
     private static void merge(ConversionSyncService.SyncResult total, ConversionSyncService.SyncResult part) {

+ 48 - 14
src/main/java/com/adx/tencent/conversionsync/ConversionSyncService.java

@@ -25,14 +25,21 @@ import java.util.concurrent.atomic.AtomicInteger;
 public class ConversionSyncService {
 
     private static final Logger log = LoggerFactory.getLogger(ConversionSyncService.class);
-    private static final int CONCURRENCY = 10;
+    private static final int REALTIME_CONCURRENCY = 10;
+    private static final int BACKFILL_CONCURRENCY = 16;
 
     private final ConversionClient baiduClient;
     private final RedisHotStore hotStore;
     private final TencentClient tencentClient;
     private final TiDBColdStore coldStore;
     private final TagEventResolver tagEventResolver;
-    private final ExecutorService executor;
+    private final ExecutorService realtimeExecutor;
+    private final ExecutorService backfillExecutor;
+
+    public enum SyncMode {
+        REALTIME,
+        BACKFILL
+    }
 
     public ConversionSyncService(ConversionClient baiduClient, RedisHotStore hotStore,
                                   TencentClient tencentClient, TiDBColdStore coldStore,
@@ -42,29 +49,41 @@ public class ConversionSyncService {
         this.tencentClient = tencentClient;
         this.coldStore = coldStore;
         this.tagEventResolver = tagEventResolver;
-        this.executor = Executors.newFixedThreadPool(CONCURRENCY,
+        this.realtimeExecutor = Executors.newFixedThreadPool(REALTIME_CONCURRENCY,
                 r -> { Thread t = new Thread(r, "conv-sync-worker"); t.setDaemon(true); return t; });
+        this.backfillExecutor = Executors.newFixedThreadPool(BACKFILL_CONCURRENCY,
+                r -> { Thread t = new Thread(r, "conv-backfill-worker"); t.setDaemon(true); return t; });
     }
 
     /**
      * 分页拉取百度转化数据并回传腾讯。
      */
     public SyncResult syncTencentConversions(ConversionQuery query) throws Exception {
+        return syncTencentConversions(query, SyncMode.REALTIME);
+    }
+
+    public SyncResult syncTencentConversions(ConversionQuery query, SyncMode mode) throws Exception {
         if (query.getPageSize() <= 0) query.setPageSize(1);
 
         SyncResult result = new SyncResult();
         int maxPages = 100;
         for (int page = query.getPageSize(); ; page++) {
             query.setPageSize(page);
-            log.info("[ConversionSync] fetching page={}", page);
+            long pageStartNs = System.nanoTime();
+            log.info("[ConversionSync] fetching date={}, page={}, mode={}", query.getDate(), page, mode);
             ConversionResponse response = baiduClient.queryPayments(query);
             List<PaymentInfo> data = response.getData();
             int dataSize = data != null ? data.size() : 0;
-            log.info("[ConversionSync] page={} returned {} records, responsePageSize={}",
-                    page, dataSize, response.getPageSize());
+            log.info("[ConversionSync] date={}, page={} returned {} records, responsePageSize={}, mode={}",
+                    query.getDate(), page, dataSize, response.getPageSize(), mode);
             if (data != null) {
-                processPayments(data, result);
+                processPayments(data, result, mode);
             }
+            long pageCostMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - pageStartNs);
+            log.info("[ConversionSync] date={}, page={} processed, cost={}ms, totalFetched={}, matched={}, sent={}, failed={}, skipped={}, alreadySent={}, mode={}",
+                    query.getDate(), page, pageCostMs,
+                    result.fetched.get(), result.matched.get(), result.sent.get(), result.failed.get(),
+                    result.skipped.get(), result.alreadySent.get(), mode);
             if (response.getPageSize() <= page) break;
             if (page >= maxPages) {
                 log.warn("[ConversionSync] reached max page limit ({}), stopping", maxPages);
@@ -74,7 +93,8 @@ public class ConversionSyncService {
         return result;
     }
 
-    private void processPayments(List<PaymentInfo> payments, SyncResult result) throws Exception {
+    private void processPayments(List<PaymentInfo> payments, SyncResult result, SyncMode mode) throws Exception {
+        ExecutorService executor = mode == SyncMode.BACKFILL ? backfillExecutor : realtimeExecutor;
         List<CompletableFuture<Void>> futures = new ArrayList<>(payments.size());
         for (PaymentInfo payment : payments) {
             futures.add(CompletableFuture.runAsync(() -> processSinglePayment(payment, result), executor));
@@ -90,14 +110,18 @@ public class ConversionSyncService {
         try {
             bid = hotStore.findBidByQk(payment.getQk());
         } catch (Exception e) {
-            log.error("[ConversionSync] findBidByQk error: qk={}, {}", payment.getQk(), e.getMessage());
-            result.failed.incrementAndGet();
-            return;
+            log.warn("[ConversionSync] hot bid lookup failed, fallback to cold store: qk={}, {}", payment.getQk(), e.getMessage());
+            bid = coldStore != null ? coldStore.getBidByQk(payment.getQk()) : null;
         }
 
         if (bid == null) {
             if (coldStore != null) {
-                try { coldStore.saveConversion(conversionRecord(payment, null, "")); } catch (Exception ignored) {}
+                try {
+                    coldStore.saveConversion(conversionRecord(payment, null, ""));
+                } catch (Exception e) {
+                    log.error("[ConversionSync] save conversion failed: qk={}, act={}, error={}",
+                            payment.getQk(), payment.getAct(), e.getMessage(), e);
+                }
             }
             result.skipped.incrementAndGet();
             return;
@@ -107,7 +131,12 @@ public class ConversionSyncService {
         // bid 不是腾讯,或没有设备标识
         if (!"tencent".equals(bid.getMedia())) {
             if (coldStore != null) {
-                try { coldStore.saveConversion(conversionRecord(payment, bid, bid.getMedia())); } catch (Exception ignored) {}
+                try {
+                    coldStore.saveConversion(conversionRecord(payment, bid, bid.getMedia()));
+                } catch (Exception e) {
+                    log.error("[ConversionSync] save conversion failed: qk={}, act={}, error={}",
+                            payment.getQk(), payment.getAct(), e.getMessage(), e);
+                }
             }
             result.skipped.incrementAndGet();
             return;
@@ -117,7 +146,12 @@ public class ConversionSyncService {
         log.info("[ConversionSync] 匹配到腾讯转化 | qk={} | act={} | deviceId={} | date={} | payment={}",
                 payment.getQk(), payment.getAct(), payment.getDeviceId(), payment.getDate(), payment.getPayment());
         if (coldStore != null) {
-            try { coldStore.saveConversion(conversionRecord(payment, bid, bid.getMedia())); } catch (Exception ignored) {}
+            try {
+                coldStore.saveConversion(conversionRecord(payment, bid, bid.getMedia()));
+            } catch (Exception e) {
+                log.error("[ConversionSync] save conversion failed: qk={}, act={}, error={}",
+                        payment.getQk(), payment.getAct(), e.getMessage(), e);
+            }
         }
 
         String actionType = null;

+ 1 - 1
src/main/java/com/adx/tencent/httpapi/AdminController.java

@@ -64,7 +64,7 @@ public class AdminController {
     }
 
     @GetMapping("/conversions/backfill/{jobId}")
-    public Object getBackfillJob(@PathVariable String jobId) {
+    public Object getBackfillJob(@PathVariable("jobId") String jobId) {
         if (backfillJobService == null) {
             throw new ResponseStatusException(HttpStatus.SERVICE_UNAVAILABLE,
                     "conversion backfill job service is not configured");