ConversionSyncService.java 8.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211
  1. package com.adx.tencent.conversionsync;
  2. import com.adx.tencent.baidu.ConversionClient;
  3. import com.adx.tencent.baidu.model.ConversionQuery;
  4. import com.adx.tencent.baidu.model.ConversionResponse;
  5. import com.adx.tencent.baidu.model.PaymentInfo;
  6. import com.adx.tencent.tencent.TencentClient;
  7. import com.adx.tencent.storage.RedisHotStore;
  8. import com.adx.tencent.storage.TiDBColdStore;
  9. import com.adx.tencent.storage.model.*;
  10. import org.slf4j.Logger;
  11. import org.slf4j.LoggerFactory;
  12. import java.time.Instant;
  13. import java.util.ArrayList;
  14. import java.util.List;
  15. import java.util.concurrent.*;
  16. import java.util.concurrent.atomic.AtomicInteger;
  17. /**
  18. * 转化同步核心逻辑。
  19. * 定时从百度拉取转化数据,匹配腾讯 bid,通过 callback URL 回传腾讯。
  20. */
  21. public class ConversionSyncService {
  22. private static final Logger log = LoggerFactory.getLogger(ConversionSyncService.class);
  23. private static final int CONCURRENCY = 10;
  24. private final ConversionClient baiduClient;
  25. private final RedisHotStore hotStore;
  26. private final TencentClient tencentClient;
  27. private final TiDBColdStore coldStore;
  28. private final ExecutorService executor;
  29. public ConversionSyncService(ConversionClient baiduClient, RedisHotStore hotStore,
  30. TencentClient tencentClient, TiDBColdStore coldStore) {
  31. this.baiduClient = baiduClient;
  32. this.hotStore = hotStore;
  33. this.tencentClient = tencentClient;
  34. this.coldStore = coldStore;
  35. this.executor = Executors.newFixedThreadPool(CONCURRENCY,
  36. r -> { Thread t = new Thread(r, "conv-sync-worker"); t.setDaemon(true); return t; });
  37. }
  38. /**
  39. * 分页拉取百度转化数据并回传腾讯。
  40. */
  41. public SyncResult syncTencentConversions(ConversionQuery query) throws Exception {
  42. if (query.getPageSize() <= 0) query.setPageSize(1);
  43. SyncResult result = new SyncResult();
  44. int maxPages = 100;
  45. for (int page = query.getPageSize(); ; page++) {
  46. query.setPageSize(page);
  47. log.info("[ConversionSync] fetching page={}", page);
  48. ConversionResponse response = baiduClient.queryPayments(query);
  49. List<PaymentInfo> data = response.getData();
  50. int dataSize = data != null ? data.size() : 0;
  51. log.info("[ConversionSync] page={} returned {} records, responsePageSize={}",
  52. page, dataSize, response.getPageSize());
  53. if (data != null) {
  54. processPayments(data, result);
  55. }
  56. if (response.getPageSize() <= page) break;
  57. if (page >= maxPages) {
  58. log.warn("[ConversionSync] reached max page limit ({}), stopping", maxPages);
  59. break;
  60. }
  61. }
  62. return result;
  63. }
  64. private void processPayments(List<PaymentInfo> payments, SyncResult result) throws Exception {
  65. List<CompletableFuture<Void>> futures = new ArrayList<>(payments.size());
  66. for (PaymentInfo payment : payments) {
  67. futures.add(CompletableFuture.runAsync(() -> processSinglePayment(payment, result), executor));
  68. }
  69. CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
  70. }
  71. private void processSinglePayment(PaymentInfo payment, SyncResult result) {
  72. result.fetched.incrementAndGet();
  73. // 用 qk 查找 bid
  74. BidRecord bid;
  75. try {
  76. bid = hotStore.findBidByQk(payment.getQk());
  77. } catch (Exception e) {
  78. log.error("[ConversionSync] findBidByQk error: qk={}, {}", payment.getQk(), e.getMessage());
  79. result.failed.incrementAndGet();
  80. return;
  81. }
  82. if (bid == null) {
  83. if (coldStore != null) {
  84. try { coldStore.saveConversion(conversionRecord(payment, "")); } catch (Exception ignored) {}
  85. }
  86. result.skipped.incrementAndGet();
  87. return;
  88. }
  89. // bid 不是腾讯,或没有设备标识
  90. if (!"tencent".equals(bid.getMedia())) {
  91. if (coldStore != null) {
  92. try { coldStore.saveConversion(conversionRecord(payment, bid.getMedia())); } catch (Exception ignored) {}
  93. }
  94. result.skipped.incrementAndGet();
  95. return;
  96. }
  97. result.matched.incrementAndGet();
  98. log.info("[ConversionSync] 匹配到腾讯转化 | qk={} | act={} | deviceId={} | date={} | payment={}",
  99. payment.getQk(), payment.getAct(), payment.getDeviceId(), payment.getDate(), payment.getPayment());
  100. if (coldStore != null) {
  101. try { coldStore.saveConversion(conversionRecord(payment, bid.getMedia())); } catch (Exception ignored) {}
  102. }
  103. // 幂等检查
  104. String callbackKey = callbackDedupeKey(payment, bid.getMedia());
  105. if (coldStore != null) {
  106. try {
  107. String conversionKey = conversionDedupeKey(payment, bid.getMedia());
  108. boolean exists = coldStore.successfulCallbackExistsForEvent(
  109. callbackKey, conversionKey, payment.getAct());
  110. if (exists) {
  111. result.alreadySent.incrementAndGet();
  112. return;
  113. }
  114. } catch (Exception e) {
  115. log.error("[ConversionSync] dedup check error: key={}, {}", callbackKey, e.getMessage());
  116. }
  117. }
  118. // 发送到腾讯
  119. String platform = resolvePlatform(bid);
  120. MediaCallbackRecord callbackRecord = tencentClient.sendConversionEvent(bid, payment, platform);
  121. callbackRecord.setDedupeKey(callbackKey);
  122. if (callbackRecord.getStatus() == 0 && callbackRecord.getErrorMessage() != null) {
  123. // 网络级错误
  124. result.failed.incrementAndGet();
  125. if (coldStore != null) {
  126. try { coldStore.saveMediaCallback(callbackRecord); } catch (Exception ignored) {}
  127. }
  128. return;
  129. }
  130. // HTTP 请求完成
  131. result.sent.incrementAndGet();
  132. if (coldStore != null) {
  133. try { coldStore.saveMediaCallback(callbackRecord); } catch (Exception ignored) {}
  134. }
  135. }
  136. // ─── 辅助方法 ─────────────────────────────────────────────────────────────
  137. /**
  138. * 从 BidRecord 中提取平台信息。
  139. * 优先读取 platform 字段,其次从 mediaParams 中提取,最后根据设备标识推断。
  140. */
  141. private static String resolvePlatform(BidRecord bid) {
  142. // 优先使用 BidRecord 的 platform 字段
  143. if (bid.getPlatform() != null && !bid.getPlatform().isBlank()) {
  144. return bid.getPlatform().toLowerCase();
  145. }
  146. if (bid.getMediaParams() == null) return "android";
  147. String platform = bid.getMediaParams().get("baidu_platform");
  148. if (platform != null && !platform.isBlank()) return platform.toLowerCase();
  149. // 有 idfa 则为 ios
  150. String idfa = bid.getMediaParams().get("idfa");
  151. if (idfa != null && !idfa.isBlank()) return "ios";
  152. return "android";
  153. }
  154. private static ConversionRecord conversionRecord(PaymentInfo p, String media) {
  155. ConversionRecord r = new ConversionRecord();
  156. r.setDedupeKey(conversionDedupeKey(p, media));
  157. r.setQk(p.getQk());
  158. r.setMedia(media);
  159. r.setDate(p.getDate());
  160. r.setAppSid(p.getAppSid());
  161. r.setCustomerName(p.getCustomerName());
  162. r.setDeviceId(p.getDeviceId());
  163. r.setConv(p.getConv());
  164. r.setPayment(p.getPayment());
  165. r.setAct(p.getAct());
  166. r.setTu(p.getTu());
  167. r.setClkTime(p.getClkTime());
  168. r.setCreatedAt(Instant.now());
  169. return r;
  170. }
  171. static String conversionDedupeKey(PaymentInfo p, String media) {
  172. return String.format("%s:%s:%s:%s:%d", media, p.getQk(), p.getDate(), p.getDeviceId(), p.getAct());
  173. }
  174. static String callbackDedupeKey(PaymentInfo p, String media) {
  175. return conversionDedupeKey(p, media) + ":" + p.getAct();
  176. }
  177. // ─── Result DTO ─────────────────────────────────────────────────────────
  178. public static class SyncResult {
  179. public final AtomicInteger fetched = new AtomicInteger();
  180. public final AtomicInteger matched = new AtomicInteger();
  181. public final AtomicInteger sent = new AtomicInteger();
  182. public final AtomicInteger skipped = new AtomicInteger();
  183. public final AtomicInteger failed = new AtomicInteger();
  184. public final AtomicInteger alreadySent = new AtomicInteger();
  185. }
  186. }