AppConfiguration.java 29 KB

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