AppConfiguration.java 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490
  1. package com.adx.tencent;
  2. import com.adx.tencent.baidu.AdxClient;
  3. import com.adx.tencent.baidu.AuctionPriceEncoder;
  4. import com.adx.tencent.baidu.ConversionClient;
  5. import com.adx.tencent.config.AppProperties;
  6. import com.adx.tencent.conversionsync.ConversionSyncRunner;
  7. import com.adx.tencent.conversionsync.ConversionSyncService;
  8. import com.adx.tencent.conversionsync.ConversionBackfillJobService;
  9. import com.adx.tencent.conversionsync.RetryService;
  10. import com.adx.tencent.httpapi.MediaPlacement;
  11. import com.adx.tencent.leader.LeaderElection;
  12. import com.adx.tencent.tagsync.TagEventResolver;
  13. import com.adx.tencent.tagsync.TagEventSyncService;
  14. import com.adx.tencent.tencent.TencentClient;
  15. import com.adx.tencent.storage.RedisHotStore;
  16. import com.adx.tencent.storage.RedisLock;
  17. import com.adx.tencent.storage.TiDBColdStore;
  18. import com.adx.tencent.worker.ColdWorker;
  19. import com.fasterxml.jackson.databind.ObjectMapper;
  20. import org.apache.ibatis.session.SqlSessionFactory;
  21. import org.mybatis.spring.SqlSessionFactoryBean;
  22. import org.slf4j.Logger;
  23. import org.slf4j.LoggerFactory;
  24. import org.springframework.beans.factory.annotation.Autowired;
  25. import org.springframework.boot.ApplicationArguments;
  26. import org.springframework.boot.ApplicationRunner;
  27. import org.springframework.context.annotation.Bean;
  28. import org.springframework.context.annotation.Configuration;
  29. import org.springframework.core.io.support.PathMatchingResourcePatternResolver;
  30. import org.springframework.data.redis.connection.RedisConnectionFactory;
  31. import org.springframework.data.redis.connection.RedisStandaloneConfiguration;
  32. import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory;
  33. import org.springframework.data.redis.core.StringRedisTemplate;
  34. import org.springframework.stereotype.Component;
  35. import org.springframework.scheduling.annotation.Scheduled;
  36. import org.springframework.lang.Nullable;
  37. import javax.sql.DataSource;
  38. import java.time.Duration;
  39. import java.util.HashMap;
  40. import java.util.List;
  41. import java.util.Map;
  42. import java.util.concurrent.ExecutorService;
  43. import java.util.concurrent.Executors;
  44. import java.util.concurrent.atomic.AtomicBoolean;
  45. /**
  46. * 腾讯-百度回传项目 Bean 组装配置类。
  47. * 负责所有核心 Bean 的创建和后台任务(Worker、Leader 选举、定时同步等)。
  48. */
  49. @Configuration
  50. public class AppConfiguration {
  51. private static final Logger log = LoggerFactory.getLogger(AppConfiguration.class);
  52. @Autowired private AppProperties props;
  53. // ─── Redis ───────────────────────────────────────────────────────────────
  54. @Bean
  55. public RedisConnectionFactory redisConnectionFactory() {
  56. String addr = props.getRedisAddr();
  57. String[] parts = addr.split(":", 2);
  58. String host = parts[0];
  59. int port = parts.length > 1 ? Integer.parseInt(parts[1]) : 6379;
  60. RedisStandaloneConfiguration cfg = new RedisStandaloneConfiguration(host, port);
  61. cfg.setDatabase(props.getRedisDb());
  62. if (props.getRedisUsername() != null && !props.getRedisUsername().isBlank()) {
  63. cfg.setUsername(props.getRedisUsername());
  64. }
  65. if (props.getRedisPassword() != null && !props.getRedisPassword().isBlank()) {
  66. cfg.setPassword(props.getRedisPassword());
  67. }
  68. return new LettuceConnectionFactory(cfg);
  69. }
  70. @Bean
  71. public StringRedisTemplate stringRedisTemplate(RedisConnectionFactory factory) {
  72. return new StringRedisTemplate(factory);
  73. }
  74. @Bean
  75. public RedisHotStore redisHotStore(StringRedisTemplate redis, ObjectMapper objectMapper) {
  76. return new RedisHotStore(redis, objectMapper,
  77. tiDBColdStore(objectMapper),
  78. "adx:tencent:",
  79. props.getRedisStream(),
  80. props.getBidTtl());
  81. }
  82. @Bean
  83. public RedisLock redisLock(StringRedisTemplate redis) {
  84. return new RedisLock(redis);
  85. }
  86. @Bean
  87. public LeaderElection leaderElection(RedisLock redisLock) {
  88. return new LeaderElection(redisLock);
  89. }
  90. // ─── TiDB(可选)────────────────────────────────────────────────────────
  91. @Bean
  92. public TiDBColdStore tiDBColdStore(ObjectMapper objectMapper) {
  93. String url = props.getTidbUrl();
  94. if (url == null || url.isBlank()) {
  95. log.info("TiDB URL is empty; cold storage disabled");
  96. return null;
  97. }
  98. try {
  99. DataSource ds = createDataSource(url, props.getTidbUsername(), props.getTidbPassword());
  100. SqlSessionFactory sqlSessionFactory = createSqlSessionFactory(ds);
  101. TiDBColdStore store = new TiDBColdStore(sqlSessionFactory, objectMapper);
  102. store.migrate();
  103. log.info("TiDB cold store initialized and migrated");
  104. return store;
  105. } catch (Exception e) {
  106. log.warn("Failed to open TiDB: {}; cold storage disabled", e.getMessage(), e);
  107. return null;
  108. }
  109. }
  110. private SqlSessionFactory createSqlSessionFactory(DataSource dataSource) throws Exception {
  111. SqlSessionFactoryBean factoryBean = new SqlSessionFactoryBean();
  112. factoryBean.setDataSource(dataSource);
  113. factoryBean.setMapperLocations(
  114. new PathMatchingResourcePatternResolver().getResources("classpath:mapper/*.xml"));
  115. return factoryBean.getObject();
  116. }
  117. private DataSource createDataSource(String jdbcUrl, String username, String password) {
  118. org.springframework.jdbc.datasource.DriverManagerDataSource ds =
  119. new org.springframework.jdbc.datasource.DriverManagerDataSource();
  120. ds.setDriverClassName("com.mysql.cj.jdbc.Driver");
  121. ds.setUrl(jdbcUrl);
  122. if (username != null && !username.isBlank()) ds.setUsername(username);
  123. if (password != null && !password.isBlank()) ds.setPassword(password);
  124. return ds;
  125. }
  126. // ─── Baidu 客户端 ────────────────────────────────────────────────────────
  127. @Bean
  128. public AuctionPriceEncoder auctionPriceEncoder() {
  129. return AuctionPriceEncoder.create(
  130. props.getBaiduAuctionPriceEncryption(),
  131. props.getBaiduAuctionPriceEKey(),
  132. props.getBaiduAuctionPriceIKey());
  133. }
  134. @Bean
  135. public AdxClient adxClient(AuctionPriceEncoder encoder, ObjectMapper objectMapper) {
  136. return new AdxClient(props.getBaiduAdxEndpoint(), encoder, objectMapper);
  137. }
  138. @Bean
  139. public ConversionClient conversionClient(ObjectMapper objectMapper) {
  140. return new ConversionClient(
  141. props.getBaiduConversionBaseUrl(),
  142. props.getBaiduCustomerName(),
  143. props.getBaiduSecret(),
  144. objectMapper);
  145. }
  146. // ─── Tencent 客户端 + Sync 服务 ──────────────────────────────────────────
  147. @Bean
  148. public TencentClient tencentClient(ObjectMapper objectMapper, StringRedisTemplate redisTemplate) {
  149. // token 从 Redis 动态获取,支持外部刷新
  150. String tokenRedisKey = props.getTencentAccessTokenRedisKey();
  151. java.util.function.Supplier<String> tokenSupplier;
  152. if (tokenRedisKey != null && !tokenRedisKey.isBlank()) {
  153. tokenSupplier = () -> redisTemplate.opsForValue().get(tokenRedisKey);
  154. } else {
  155. // 兜底:使用配置文件中的静态 token
  156. String staticToken = props.getTencentAccessToken();
  157. tokenSupplier = () -> staticToken;
  158. }
  159. return new TencentClient(
  160. tokenSupplier,
  161. props.getTencentActMap(),
  162. objectMapper);
  163. }
  164. @Bean
  165. public TagEventResolver tagEventResolver(@Nullable TiDBColdStore coldStore,
  166. StringRedisTemplate redisTemplate) {
  167. if (coldStore == null) return null;
  168. return new TagEventResolver(redisTemplate, coldStore, props.getTagEventSyncRedisPrefix());
  169. }
  170. @Bean
  171. public ConversionSyncService conversionSyncService(ConversionClient baiduClient,
  172. RedisHotStore hotStore,
  173. TencentClient tencentClient,
  174. @Nullable TiDBColdStore coldStore,
  175. @Nullable TagEventResolver tagEventResolver) {
  176. if (coldStore == null) return null;
  177. return new ConversionSyncService(baiduClient, hotStore, tencentClient, coldStore, tagEventResolver);
  178. }
  179. @Bean
  180. public ConversionSyncRunner conversionSyncRunner(@Nullable ConversionSyncService syncer) {
  181. if (syncer == null) return null;
  182. return new ConversionSyncRunner(syncer,
  183. props.getConversionSyncDateOffsetDays(),
  184. props.getConversionSyncPageSize(),
  185. props.getConversionSyncActs());
  186. }
  187. @Bean
  188. public ConversionBackfillJobService conversionBackfillJobService(@Nullable ConversionSyncRunner syncer) {
  189. if (syncer == null) return null;
  190. return new ConversionBackfillJobService(syncer);
  191. }
  192. @Bean
  193. public RetryService retryService(@Nullable TiDBColdStore coldStore, TencentClient tencentClient) {
  194. if (coldStore == null) return null;
  195. return new RetryService(coldStore, tencentClient, props.getCallbackRetryLimit());
  196. }
  197. @Bean
  198. public TagEventSyncService tagEventSyncService(@Nullable TiDBColdStore coldStore,
  199. StringRedisTemplate redisTemplate) {
  200. if (coldStore == null) return null;
  201. // TTL 取同步间隔的 3 倍,避免表中已删除的 tag 在 Redis 长期残留
  202. Duration ttl = props.getTagEventSyncInterval().multipliedBy(3);
  203. return new TagEventSyncService(coldStore, redisTemplate,
  204. props.getTagEventSyncRedisPrefix(), ttl);
  205. }
  206. // ─── ColdWorker ─────────────────────────────────────────────────────────
  207. @Bean
  208. public ColdWorker coldWorker(RedisHotStore hotStore, @Nullable TiDBColdStore coldStore,
  209. ObjectMapper objectMapper) {
  210. if (coldStore == null) return null;
  211. return new ColdWorker(hotStore, coldStore,
  212. props.getWorkerGroup(), props.getWorkerConsumer(),
  213. props.getWorkerBatch(), null,
  214. props.getRedisStreamMaxLen(), objectMapper);
  215. }
  216. // ─── MediaPlacement(腾讯)───────────────────────────────────────────────
  217. @Bean
  218. public MediaPlacement tencentPlacement() {
  219. MediaPlacement p = new MediaPlacement();
  220. p.setBaiduMediaId(props.getBaiduMediaId());
  221. p.setBaiduAppId(props.getBaiduAppId());
  222. p.setBaiduTagId(props.getBaiduTagId());
  223. p.setBidFloor(props.getBaiduBidFloor());
  224. p.setActionTypes(List.of(0, 1, 2));
  225. Map<String, MediaPlacement.BaiduAppPlacement> platforms = new HashMap<>();
  226. if (props.getBaiduAndroidAppId() != null && !props.getBaiduAndroidAppId().isBlank()) {
  227. MediaPlacement.BaiduAppPlacement android = new MediaPlacement.BaiduAppPlacement();
  228. android.setAppId(props.getBaiduAndroidAppId());
  229. android.setTagIds(props.getBaiduAndroidTagIds());
  230. platforms.put("android", android);
  231. }
  232. if (props.getBaiduIosAppId() != null && !props.getBaiduIosAppId().isBlank()) {
  233. MediaPlacement.BaiduAppPlacement ios = new MediaPlacement.BaiduAppPlacement();
  234. ios.setAppId(props.getBaiduIosAppId());
  235. ios.setTagIds(props.getBaiduIosTagIds());
  236. platforms.put("ios", ios);
  237. }
  238. if (!platforms.isEmpty()) p.setPlatforms(platforms);
  239. return p;
  240. }
  241. // ─── 后台任务调度 ────────────────────────────────────────────────────────
  242. @Component
  243. static class BackgroundTasks implements ApplicationRunner {
  244. private static final Logger taskLog = LoggerFactory.getLogger(BackgroundTasks.class);
  245. @Autowired private AppProperties props;
  246. @Autowired(required = false) private ColdWorker coldWorker;
  247. @Autowired(required = false) private ConversionSyncRunner conversionSyncRunner;
  248. @Autowired(required = false) private RetryService retryService;
  249. @Autowired(required = false) private TagEventSyncService tagEventSyncService;
  250. @Autowired(required = false) private LeaderElection leaderElection;
  251. private final AtomicBoolean stopped = new AtomicBoolean(false);
  252. private final ExecutorService executor = Executors.newCachedThreadPool(r -> {
  253. Thread t = new Thread(r);
  254. t.setDaemon(true);
  255. return t;
  256. });
  257. @Override
  258. public void run(ApplicationArguments args) {
  259. // 冷路径 Worker
  260. if (coldWorker != null) {
  261. executor.submit(() -> {
  262. while (!stopped.get()) {
  263. try {
  264. coldWorker.runOnce();
  265. } catch (Exception e) {
  266. taskLog.error("cold worker: {}", e.getMessage(), e);
  267. sleep(1000);
  268. }
  269. }
  270. });
  271. taskLog.info("Cold worker started");
  272. }
  273. // 转化同步
  274. if (conversionSyncRunner != null && props.isConversionSyncEnabled()) {
  275. if (props.isSkipLeaderElection()) {
  276. executor.submit(() -> runConversionSync(() -> stopped.get()));
  277. taskLog.info("Conversion sync started (skip leader election)");
  278. } else if (leaderElection != null) {
  279. executor.submit(() ->
  280. leaderElection.run(
  281. "adx:lock:tencent:conversion-sync",
  282. props.getTaskLockTtl(), props.getTaskLockRenewInterval(), props.getTaskLockRetryInterval(),
  283. stopped::get,
  284. (jobStop) -> runConversionSync(jobStop),
  285. e -> taskLog.error("tencent conversion sync leader: {}", e.getMessage(), e)
  286. )
  287. );
  288. taskLog.info("Conversion sync leader election started");
  289. }
  290. } else {
  291. taskLog.warn("Conversion sync NOT started: runner={}, enabled={}",
  292. conversionSyncRunner != null, props.isConversionSyncEnabled());
  293. }
  294. // 回调重试
  295. if (retryService != null && props.isCallbackRetryEnabled()) {
  296. if (props.isSkipLeaderElection()) {
  297. executor.submit(() -> runCallbackRetry(() -> stopped.get()));
  298. taskLog.info("Callback retry started (skip leader election)");
  299. } else if (leaderElection != null) {
  300. executor.submit(() ->
  301. leaderElection.run(
  302. "adx:lock:tencent:callback-retry",
  303. props.getTaskLockTtl(), props.getTaskLockRenewInterval(), props.getTaskLockRetryInterval(),
  304. stopped::get,
  305. (jobStop) -> runCallbackRetry(jobStop),
  306. e -> taskLog.error("tencent callback retry leader: {}", e.getMessage(), e)
  307. )
  308. );
  309. taskLog.info("Callback retry leader election started");
  310. }
  311. } else {
  312. taskLog.warn("Callback retry NOT started: retryService={}, enabled={}",
  313. retryService != null, props.isCallbackRetryEnabled());
  314. }
  315. // 广告位回传方式同步
  316. if (tagEventSyncService != null && props.isTagEventSyncEnabled()) {
  317. if (props.isSkipLeaderElection()) {
  318. executor.submit(() -> runTagEventSync(() -> stopped.get()));
  319. taskLog.info("Tag event sync started (skip leader election)");
  320. } else if (leaderElection != null) {
  321. executor.submit(() ->
  322. leaderElection.run(
  323. "adx:lock:tencent:tag-event-sync",
  324. props.getTaskLockTtl(), props.getTaskLockRenewInterval(), props.getTaskLockRetryInterval(),
  325. stopped::get,
  326. (jobStop) -> runTagEventSync(jobStop),
  327. e -> taskLog.error("tencent tag event sync leader: {}", e.getMessage(), e)
  328. )
  329. );
  330. taskLog.info("Tag event sync leader election started");
  331. }
  332. } else {
  333. taskLog.warn("Tag event sync NOT started: service={}, enabled={}",
  334. tagEventSyncService != null, props.isTagEventSyncEnabled());
  335. }
  336. }
  337. private void runConversionSync(LeaderElection.StopSignal jobStop) {
  338. long intervalMs = props.getConversionSyncInterval().toMillis();
  339. taskLog.info("[ConversionSync] task started, interval={}ms", intervalMs);
  340. while (!jobStop.isStopped()) {
  341. try {
  342. taskLog.info("[ConversionSync] executing...");
  343. ConversionSyncService.SyncResult r = conversionSyncRunner.runOnce();
  344. taskLog.info("[ConversionSync] done: fetched={} sent={} failed={}", r.fetched.get(), r.sent.get(), r.failed.get());
  345. } catch (Exception e) {
  346. taskLog.error("[ConversionSync] error: {}", e.getMessage(), e);
  347. }
  348. taskLog.info("[ConversionSync] sleeping {}ms until next run", intervalMs);
  349. sleepResponsive(intervalMs, jobStop);
  350. }
  351. taskLog.info("[ConversionSync] task stopped");
  352. }
  353. @Scheduled(cron = "0 17 1,7,10,16,20,23 * * ?")
  354. public void runDelayedConversionSync() {
  355. if (conversionSyncRunner == null || !props.isConversionSyncEnabled()) {
  356. return;
  357. }
  358. if (stopped.get()) {
  359. return;
  360. }
  361. Runnable job = () -> {
  362. long startNs = System.nanoTime();
  363. try {
  364. taskLog.info("[ConversionBackfill] executing for offsets=1..7");
  365. ConversionSyncService.SyncResult r = conversionSyncRunner.runBackfillForOffsetsParallel(
  366. List.of(-1, -2, -3, -4, -5, -6, -7), 3);
  367. long costMs = java.util.concurrent.TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNs);
  368. taskLog.info("[ConversionBackfill] done: fetched={} matched={} sent={} failed={} skipped={} alreadySent={}",
  369. r.fetched.get(), r.matched.get(), r.sent.get(), r.failed.get(),
  370. r.skipped.get(), r.alreadySent.get());
  371. taskLog.info("[ConversionBackfill] total cost={}ms", costMs);
  372. } catch (Exception e) {
  373. long costMs = java.util.concurrent.TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNs);
  374. taskLog.error("[ConversionBackfill] error after {}ms: {}", costMs, e.getMessage(), e);
  375. }
  376. };
  377. if (props.isSkipLeaderElection() || leaderElection == null) {
  378. job.run();
  379. return;
  380. }
  381. executor.submit(() ->
  382. leaderElection.runOnce(
  383. "adx:lock:tencent:conversion-backfill",
  384. props.getTaskLockTtl(), props.getTaskLockRenewInterval(),
  385. stopped::get,
  386. stop -> job.run(),
  387. e -> taskLog.error("tencent conversion backfill leader: {}", e.getMessage(), e)
  388. )
  389. );
  390. }
  391. private void runCallbackRetry(LeaderElection.StopSignal jobStop) {
  392. long intervalMs = props.getCallbackRetryInterval().toMillis();
  393. taskLog.info("[CallbackRetry] task started, interval={}ms, limit={}", intervalMs, props.getCallbackRetryLimit());
  394. while (!jobStop.isStopped()) {
  395. try {
  396. taskLog.info("[CallbackRetry] executing...");
  397. RetryService.RetryResult r = retryService.retryTencentCallbacks(props.getCallbackRetryLimit());
  398. taskLog.info("[CallbackRetry] done: fetched={} sent={} failed={}", r.fetched, r.sent, r.failed);
  399. } catch (Exception e) {
  400. taskLog.error("[CallbackRetry] error: {}", e.getMessage(), e);
  401. }
  402. taskLog.info("[CallbackRetry] sleeping {}ms until next run", intervalMs);
  403. sleepResponsive(intervalMs, jobStop);
  404. }
  405. taskLog.info("[CallbackRetry] task stopped");
  406. }
  407. private void runTagEventSync(LeaderElection.StopSignal jobStop) {
  408. long intervalMs = props.getTagEventSyncInterval().toMillis();
  409. taskLog.info("[TagEventSync] task started, interval={}ms", intervalMs);
  410. while (!jobStop.isStopped()) {
  411. try {
  412. int n = tagEventSyncService.syncOnce();
  413. taskLog.info("[TagEventSync] done: synced={}", n);
  414. } catch (Exception e) {
  415. taskLog.error("[TagEventSync] error: {}", e.getMessage(), e);
  416. }
  417. sleepResponsive(intervalMs, jobStop);
  418. }
  419. taskLog.info("[TagEventSync] task stopped");
  420. }
  421. private static void sleepResponsive(long ms, LeaderElection.StopSignal jobStop) {
  422. long deadline = System.currentTimeMillis() + ms;
  423. while (!jobStop.isStopped()) {
  424. long remaining = deadline - System.currentTimeMillis();
  425. if (remaining <= 0) break;
  426. try {
  427. Thread.sleep(Math.min(remaining, 200));
  428. } catch (InterruptedException e) {
  429. Thread.currentThread().interrupt();
  430. break;
  431. }
  432. }
  433. }
  434. private static void sleep(long ms) {
  435. try { Thread.sleep(ms); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
  436. }
  437. }
  438. }