ConversionSyncService.java 10 KB

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