package com.adx.tencent; import com.adx.tencent.baidu.AdxClient; import com.adx.tencent.baidu.AuctionPriceEncoder; import com.adx.tencent.baidu.ConversionClient; import com.adx.tencent.config.AppProperties; import com.adx.tencent.conversionsync.ConversionSyncRunner; import com.adx.tencent.conversionsync.ConversionSyncService; import com.adx.tencent.conversionsync.ConversionBackfillJobService; import com.adx.tencent.conversionsync.RetryService; import com.adx.tencent.httpapi.MediaPlacement; import com.adx.tencent.honor.HonorClient; import com.adx.tencent.honor.HonorColdStore; import com.adx.tencent.honor.HonorColdWorker; import com.adx.tencent.honor.HonorHotStore; import com.adx.tencent.honor.HonorPlacement; import com.adx.tencent.honor.HonorRetryService; import com.adx.tencent.honor.HonorTagEventResolver; import com.adx.tencent.honor.HonorTagEventSyncService; import com.adx.tencent.leader.LeaderElection; import com.adx.tencent.tagsync.TagEventResolver; import com.adx.tencent.tagsync.TagEventSyncService; import com.adx.tencent.tencent.TencentClient; import com.adx.tencent.storage.RedisHotStore; import com.adx.tencent.storage.RedisLock; import com.adx.tencent.storage.TiDBColdStore; import com.adx.tencent.worker.ColdWorker; import com.fasterxml.jackson.databind.ObjectMapper; import com.zaxxer.hikari.HikariConfig; import com.zaxxer.hikari.HikariDataSource; import io.lettuce.core.ClientOptions; import io.lettuce.core.SocketOptions; import org.apache.ibatis.session.SqlSessionFactory; import org.mybatis.spring.SqlSessionFactoryBean; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.ApplicationArguments; import org.springframework.boot.ApplicationRunner; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.core.io.support.PathMatchingResourcePatternResolver; import org.springframework.data.redis.connection.RedisConnectionFactory; import org.springframework.data.redis.connection.RedisStandaloneConfiguration; import org.springframework.data.redis.connection.lettuce.LettuceClientConfiguration; import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.stereotype.Component; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.lang.Nullable; import javax.sql.DataSource; import java.time.Duration; import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.atomic.AtomicBoolean; /** * 腾讯-百度回传项目 Bean 组装配置类。 * 负责所有核心 Bean 的创建和后台任务(Worker、Leader 选举、定时同步等)。 */ @Configuration public class AppConfiguration { private static final Logger log = LoggerFactory.getLogger(AppConfiguration.class); @Autowired private AppProperties props; // ─── Redis ─────────────────────────────────────────────────────────────── @Bean public RedisConnectionFactory redisConnectionFactory() { String addr = props.getRedisAddr(); String[] parts = addr.split(":", 2); String host = parts[0]; int port = parts.length > 1 ? Integer.parseInt(parts[1]) : 6379; RedisStandaloneConfiguration cfg = new RedisStandaloneConfiguration(host, port); cfg.setDatabase(props.getRedisDb()); if (props.getRedisUsername() != null && !props.getRedisUsername().isBlank()) { cfg.setUsername(props.getRedisUsername()); } if (props.getRedisPassword() != null && !props.getRedisPassword().isBlank()) { cfg.setPassword(props.getRedisPassword()); } LettuceClientConfiguration clientCfg = LettuceClientConfiguration.builder() .clientOptions(ClientOptions.builder() .socketOptions(SocketOptions.builder() .connectTimeout(props.getRedisConnectTimeout()) .build()) .build()) .commandTimeout(props.getRedisCommandTimeout()) .shutdownTimeout(props.getRedisShutdownTimeout()) .build(); LettuceConnectionFactory factory = new LettuceConnectionFactory(cfg, clientCfg); factory.setValidateConnection(props.isRedisValidateConnection()); factory.setShareNativeConnection(false); return factory; } @Bean public StringRedisTemplate stringRedisTemplate(RedisConnectionFactory factory) { return new StringRedisTemplate(factory); } @Bean public RedisHotStore redisHotStore(StringRedisTemplate redis, ObjectMapper objectMapper) { return new RedisHotStore(redis, objectMapper, tiDBColdStore(objectMapper), "adx:tencent:", props.getRedisStream(), props.getBidTtl()); } @Bean public RedisLock redisLock(StringRedisTemplate redis) { return new RedisLock(redis); } @Bean public LeaderElection leaderElection(RedisLock redisLock) { return new LeaderElection(redisLock); } // ─── TiDB(可选)──────────────────────────────────────────────────────── @Bean public TiDBColdStore tiDBColdStore(ObjectMapper objectMapper) { String url = props.getTidbUrl(); if (url == null || url.isBlank()) { log.info("TiDB URL is empty; cold storage disabled"); return null; } try { DataSource ds = createDataSource(url, props.getTidbUsername(), props.getTidbPassword()); SqlSessionFactory sqlSessionFactory = createSqlSessionFactory(ds); TiDBColdStore store = new TiDBColdStore(sqlSessionFactory, objectMapper); store.migrate(); log.info("TiDB cold store initialized and migrated"); return store; } catch (Exception e) { log.warn("Failed to open TiDB: {}; cold storage disabled", e.getMessage(), e); return null; } } private SqlSessionFactory createSqlSessionFactory(DataSource dataSource) throws Exception { SqlSessionFactoryBean factoryBean = new SqlSessionFactoryBean(); factoryBean.setDataSource(dataSource); factoryBean.setMapperLocations( new PathMatchingResourcePatternResolver().getResources("classpath:mapper/*.xml")); return factoryBean.getObject(); } private DataSource createDataSource(String jdbcUrl, String username, String password) { HikariConfig config = new HikariConfig(); config.setDriverClassName("com.mysql.cj.jdbc.Driver"); config.setJdbcUrl(jdbcUrl); if (username != null && !username.isBlank()) config.setUsername(username); if (password != null && !password.isBlank()) config.setPassword(password); config.setPoolName("adx-tidb"); config.setMinimumIdle(props.getTidbMinimumIdle()); config.setMaximumPoolSize(props.getTidbMaximumPoolSize()); config.setConnectionTimeout(props.getTidbConnectionTimeout().toMillis()); config.setValidationTimeout(props.getTidbValidationTimeout().toMillis()); config.setIdleTimeout(props.getTidbIdleTimeout().toMillis()); config.setMaxLifetime(props.getTidbMaxLifetime().toMillis()); config.setKeepaliveTime(props.getTidbKeepaliveTime().toMillis()); return new HikariDataSource(config); } // ─── Baidu 客户端 ──────────────────────────────────────────────────────── @Bean public AuctionPriceEncoder auctionPriceEncoder() { return AuctionPriceEncoder.create( props.getBaiduAuctionPriceEncryption(), props.getBaiduAuctionPriceEKey(), props.getBaiduAuctionPriceIKey()); } @Bean public AdxClient adxClient(AuctionPriceEncoder encoder, ObjectMapper objectMapper) { return new AdxClient(props.getBaiduAdxEndpoint(), encoder, objectMapper); } @Bean public ConversionClient conversionClient(ObjectMapper objectMapper) { return new ConversionClient( props.getBaiduConversionBaseUrl(), props.getBaiduCustomerName(), props.getBaiduSecret(), objectMapper); } @Bean public HonorColdStore honorColdStore(ObjectMapper objectMapper) { String url = props.getTidbUrl(); if (url == null || url.isBlank()) return null; try { DataSource ds = createDataSource(url, props.getTidbUsername(), props.getTidbPassword()); SqlSessionFactory sqlSessionFactory = createSqlSessionFactory(ds); HonorColdStore store = new HonorColdStore(sqlSessionFactory, objectMapper); store.migrate(); return store; } catch (Exception e) { log.warn("Failed to open Honor TiDB: {}", e.getMessage(), e); return null; } } @Bean public HonorHotStore honorHotStore(StringRedisTemplate redis, ObjectMapper objectMapper, @Nullable HonorColdStore honorColdStore) { if (honorColdStore == null) return null; return new HonorHotStore(redis, objectMapper, honorColdStore, "adx:honor:", props.getHonorRedisStream(), props.getBidTtl()); } @Bean public HonorPlacement honorPlacement() { HonorPlacement p = new HonorPlacement(); p.setBaiduMediaId(props.getBaiduMediaId()); p.setBaiduAppId(props.getBaiduAppId()); p.setBaiduTagId(props.getBaiduTagId()); p.setBidFloor(props.getBaiduBidFloor()); p.setActionTypes(List.of(0, 1, 2)); Map platforms = new HashMap<>(); String honorAndroidAppId = props.getBaiduAndroidAppId(); List honorAndroidTagIds = props.getBaiduAndroidTagIds(); if (honorAndroidAppId != null && !honorAndroidAppId.isBlank()) { HonorPlacement.HonorAppPlacement android = new HonorPlacement.HonorAppPlacement(); android.setAppId(honorAndroidAppId); android.setTagIds(honorAndroidTagIds); platforms.put("android", android); } String honorIosAppId = props.getBaiduIosAppId(); List honorIosTagIds = props.getBaiduIosTagIds(); if (honorIosAppId != null && !honorIosAppId.isBlank()) { HonorPlacement.HonorAppPlacement ios = new HonorPlacement.HonorAppPlacement(); ios.setAppId(honorIosAppId); ios.setTagIds(honorIosTagIds); platforms.put("ios", ios); } if (!platforms.isEmpty()) p.setPlatforms(platforms); return p; } @Bean public HonorClient honorClient(ObjectMapper objectMapper) { return new HonorClient(props.getHonorConversionBaseUrl(), props.getHonorChannelId(), objectMapper); } @Bean public HonorRetryService honorRetryService(@Nullable HonorColdStore honorColdStore, HonorClient honorClient) { if (honorColdStore == null) return null; return new HonorRetryService(honorColdStore, honorClient, props.getHonorCallbackRetryLimit()); } @Bean public HonorColdWorker honorColdWorker(@Nullable HonorHotStore honorHotStore, @Nullable HonorColdStore honorColdStore, ObjectMapper objectMapper) { if (honorHotStore == null || honorColdStore == null) return null; return new HonorColdWorker(honorHotStore, honorColdStore, "honor-cold-writers", props.getWorkerConsumer() + "-honor", props.getWorkerBatch(), null, props.getRedisStreamMaxLen(), objectMapper); } // ─── Tencent 客户端 + Sync 服务 ────────────────────────────────────────── @Bean public TencentClient tencentClient(ObjectMapper objectMapper, StringRedisTemplate redisTemplate) { // token 从 Redis 动态获取,支持外部刷新 String tokenRedisKey = props.getTencentAccessTokenRedisKey(); java.util.function.Supplier tokenSupplier; if (tokenRedisKey != null && !tokenRedisKey.isBlank()) { tokenSupplier = () -> redisTemplate.opsForValue().get(tokenRedisKey); } else { // 兜底:使用配置文件中的静态 token String staticToken = props.getTencentAccessToken(); tokenSupplier = () -> staticToken; } return new TencentClient( tokenSupplier, props.getTencentActMap(), objectMapper); } @Bean public TagEventResolver tagEventResolver(@Nullable TiDBColdStore coldStore, StringRedisTemplate redisTemplate) { if (coldStore == null) return null; return new TagEventResolver(redisTemplate, coldStore, props.getTagEventSyncRedisPrefix()); } @Bean public HonorTagEventResolver honorTagEventResolver(@Nullable HonorColdStore honorColdStore, StringRedisTemplate redisTemplate) { if (honorColdStore == null) return null; return new HonorTagEventResolver(redisTemplate, honorColdStore, props.getHonorTagEventSyncRedisPrefix()); } @Bean public ConversionSyncService conversionSyncService(ConversionClient baiduClient, RedisHotStore hotStore, TencentClient tencentClient, @Nullable TiDBColdStore coldStore, @Nullable TagEventResolver tagEventResolver, @Nullable HonorTagEventResolver honorTagEventResolver, @Nullable HonorHotStore honorHotStore, @Nullable HonorColdStore honorColdStore, HonorClient honorClient) { if (coldStore == null) return null; return new ConversionSyncService(baiduClient, hotStore, tencentClient, coldStore, tagEventResolver, honorTagEventResolver, honorHotStore, honorColdStore, honorClient); } @Bean public ConversionSyncRunner conversionSyncRunner(@Nullable ConversionSyncService syncer) { if (syncer == null) return null; return new ConversionSyncRunner(syncer, props.getConversionSyncDateOffsetDays(), props.getConversionSyncPageSize(), props.getConversionSyncActs()); } @Bean public ConversionBackfillJobService conversionBackfillJobService(@Nullable ConversionSyncRunner syncer) { if (syncer == null) return null; return new ConversionBackfillJobService(syncer); } @Bean public RetryService retryService(@Nullable TiDBColdStore coldStore, TencentClient tencentClient) { if (coldStore == null) return null; return new RetryService(coldStore, tencentClient, props.getCallbackRetryLimit()); } @Bean public TagEventSyncService tagEventSyncService(@Nullable TiDBColdStore coldStore, StringRedisTemplate redisTemplate) { if (coldStore == null) return null; // TTL 取同步间隔的 3 倍,避免表中已删除的 tag 在 Redis 长期残留 Duration ttl = props.getTagEventSyncInterval().multipliedBy(3); return new TagEventSyncService(coldStore, redisTemplate, props.getTagEventSyncRedisPrefix(), ttl); } @Bean public HonorTagEventSyncService honorTagEventSyncService(@Nullable HonorColdStore honorColdStore, StringRedisTemplate redisTemplate) { if (honorColdStore == null) return null; Duration ttl = props.getHonorTagEventSyncInterval().multipliedBy(3); return new HonorTagEventSyncService(honorColdStore, redisTemplate, props.getHonorTagEventSyncRedisPrefix(), ttl); } // ─── ColdWorker ───────────────────────────────────────────────────────── @Bean public ColdWorker coldWorker(RedisHotStore hotStore, @Nullable TiDBColdStore coldStore, ObjectMapper objectMapper) { if (coldStore == null) return null; return new ColdWorker(hotStore, coldStore, props.getWorkerGroup(), props.getWorkerConsumer(), props.getWorkerBatch(), null, props.getRedisStreamMaxLen(), objectMapper); } // ─── MediaPlacement(腾讯)─────────────────────────────────────────────── @Bean public MediaPlacement tencentPlacement() { MediaPlacement p = new MediaPlacement(); p.setBaiduMediaId(props.getBaiduMediaId()); p.setBaiduAppId(props.getBaiduAppId()); p.setBaiduTagId(props.getBaiduTagId()); p.setBidFloor(props.getBaiduBidFloor()); p.setActionTypes(List.of(0, 1, 2)); Map platforms = new HashMap<>(); if (props.getBaiduAndroidAppId() != null && !props.getBaiduAndroidAppId().isBlank()) { MediaPlacement.BaiduAppPlacement android = new MediaPlacement.BaiduAppPlacement(); android.setAppId(props.getBaiduAndroidAppId()); android.setTagIds(props.getBaiduAndroidTagIds()); platforms.put("android", android); } if (props.getBaiduIosAppId() != null && !props.getBaiduIosAppId().isBlank()) { MediaPlacement.BaiduAppPlacement ios = new MediaPlacement.BaiduAppPlacement(); ios.setAppId(props.getBaiduIosAppId()); ios.setTagIds(props.getBaiduIosTagIds()); platforms.put("ios", ios); } if (!platforms.isEmpty()) p.setPlatforms(platforms); return p; } // ─── 后台任务调度 ──────────────────────────────────────────────────────── @Component static class BackgroundTasks implements ApplicationRunner { private static final Logger taskLog = LoggerFactory.getLogger(BackgroundTasks.class); @Autowired private AppProperties props; @Autowired(required = false) private ColdWorker coldWorker; @Autowired(required = false) private ConversionSyncRunner conversionSyncRunner; @Autowired(required = false) private RetryService retryService; @Autowired(required = false) private TagEventSyncService tagEventSyncService; @Autowired(required = false) private LeaderElection leaderElection; private final AtomicBoolean stopped = new AtomicBoolean(false); private final ExecutorService executor = Executors.newCachedThreadPool(r -> { Thread t = new Thread(r); t.setDaemon(true); return t; }); @Override public void run(ApplicationArguments args) { // 冷路径 Worker if (coldWorker != null) { executor.submit(() -> { while (!stopped.get()) { try { coldWorker.runOnce(); } catch (Exception e) { taskLog.error("cold worker: {}", e.getMessage(), e); sleep(1000); } } }); taskLog.info("Cold worker started"); } // 转化同步 if (conversionSyncRunner != null && props.isConversionSyncEnabled()) { if (props.isSkipLeaderElection()) { executor.submit(() -> runConversionSync(() -> stopped.get())); taskLog.info("Conversion sync started (skip leader election)"); } else if (leaderElection != null) { executor.submit(() -> leaderElection.run( "adx:lock:tencent:conversion-sync", props.getTaskLockTtl(), props.getTaskLockRenewInterval(), props.getTaskLockRetryInterval(), stopped::get, (jobStop) -> runConversionSync(jobStop), e -> taskLog.error("tencent conversion sync leader: {}", e.getMessage(), e) ) ); taskLog.info("Conversion sync leader election started"); } } else { taskLog.warn("Conversion sync NOT started: runner={}, enabled={}", conversionSyncRunner != null, props.isConversionSyncEnabled()); } // 回调重试 if (retryService != null && props.isCallbackRetryEnabled()) { if (props.isSkipLeaderElection()) { executor.submit(() -> runCallbackRetry(() -> stopped.get())); taskLog.info("Callback retry started (skip leader election)"); } else if (leaderElection != null) { executor.submit(() -> leaderElection.run( "adx:lock:tencent:callback-retry", props.getTaskLockTtl(), props.getTaskLockRenewInterval(), props.getTaskLockRetryInterval(), stopped::get, (jobStop) -> runCallbackRetry(jobStop), e -> taskLog.error("tencent callback retry leader: {}", e.getMessage(), e) ) ); taskLog.info("Callback retry leader election started"); } } else { taskLog.warn("Callback retry NOT started: retryService={}, enabled={}", retryService != null, props.isCallbackRetryEnabled()); } // 广告位回传方式同步 if (tagEventSyncService != null && props.isTagEventSyncEnabled()) { if (props.isSkipLeaderElection()) { executor.submit(() -> runTagEventSync(() -> stopped.get())); taskLog.info("Tag event sync started (skip leader election)"); } else if (leaderElection != null) { executor.submit(() -> leaderElection.run( "adx:lock:tencent:tag-event-sync", props.getTaskLockTtl(), props.getTaskLockRenewInterval(), props.getTaskLockRetryInterval(), stopped::get, (jobStop) -> runTagEventSync(jobStop), e -> taskLog.error("tencent tag event sync leader: {}", e.getMessage(), e) ) ); taskLog.info("Tag event sync leader election started"); } } else { taskLog.warn("Tag event sync NOT started: service={}, enabled={}", tagEventSyncService != null, props.isTagEventSyncEnabled()); } } private void runConversionSync(LeaderElection.StopSignal jobStop) { long intervalMs = props.getConversionSyncInterval().toMillis(); taskLog.info("[ConversionSync] task started, interval={}ms", intervalMs); while (!jobStop.isStopped()) { try { taskLog.info("[ConversionSync] executing..."); ConversionSyncService.SyncResult r = conversionSyncRunner.runOnce(); taskLog.info("[ConversionSync] done: fetched={} sent={} failed={}", r.fetched.get(), r.sent.get(), r.failed.get()); } catch (Exception e) { taskLog.error("[ConversionSync] error: {}", e.getMessage(), e); } taskLog.info("[ConversionSync] sleeping {}ms until next run", intervalMs); sleepResponsive(intervalMs, jobStop); } taskLog.info("[ConversionSync] task stopped"); } @Scheduled(cron = "0 17 1,7,10,16,20,23 * * ?") public void runDelayedConversionSync() { if (conversionSyncRunner == null || !props.isConversionSyncEnabled()) { return; } if (stopped.get()) { return; } Runnable job = () -> { long startNs = System.nanoTime(); try { taskLog.info("[ConversionBackfill] executing for offsets=1..5"); ConversionSyncService.SyncResult r = conversionSyncRunner.runBackfillForOffsetsParallel( List.of(-1, -2, -3, -4, -5), 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) { long costMs = java.util.concurrent.TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNs); taskLog.error("[ConversionBackfill] error after {}ms: {}", costMs, e.getMessage(), e); } }; if (props.isSkipLeaderElection() || leaderElection == null) { job.run(); return; } executor.submit(() -> leaderElection.runOnce( "adx:lock:tencent:conversion-backfill", props.getTaskLockTtl(), props.getTaskLockRenewInterval(), stopped::get, stop -> job.run(), e -> taskLog.error("tencent conversion backfill leader: {}", e.getMessage(), e) ) ); } private void runCallbackRetry(LeaderElection.StopSignal jobStop) { long intervalMs = props.getCallbackRetryInterval().toMillis(); taskLog.info("[CallbackRetry] task started, interval={}ms, limit={}", intervalMs, props.getCallbackRetryLimit()); while (!jobStop.isStopped()) { try { taskLog.info("[CallbackRetry] executing..."); RetryService.RetryResult r = retryService.retryTencentCallbacks(props.getCallbackRetryLimit()); taskLog.info("[CallbackRetry] done: fetched={} sent={} failed={}", r.fetched, r.sent, r.failed); } catch (Exception e) { taskLog.error("[CallbackRetry] error: {}", e.getMessage(), e); } taskLog.info("[CallbackRetry] sleeping {}ms until next run", intervalMs); sleepResponsive(intervalMs, jobStop); } taskLog.info("[CallbackRetry] task stopped"); } private void runTagEventSync(LeaderElection.StopSignal jobStop) { long intervalMs = props.getTagEventSyncInterval().toMillis(); taskLog.info("[TagEventSync] task started, interval={}ms", intervalMs); while (!jobStop.isStopped()) { try { int n = tagEventSyncService.syncOnce(); taskLog.info("[TagEventSync] done: synced={}", n); } catch (Exception e) { taskLog.error("[TagEventSync] error: {}", e.getMessage(), e); } sleepResponsive(intervalMs, jobStop); } taskLog.info("[TagEventSync] task stopped"); } private static void sleepResponsive(long ms, LeaderElection.StopSignal jobStop) { long deadline = System.currentTimeMillis() + ms; while (!jobStop.isStopped()) { long remaining = deadline - System.currentTimeMillis(); if (remaining <= 0) break; try { Thread.sleep(Math.min(remaining, 200)); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } } private static void sleep(long ms) { try { Thread.sleep(ms); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } } }