yumeng hai 3 semanas
pai
achega
bf5a786203

+ 50 - 0
src/main/java/com/adx/tencent/AppConfiguration.java

@@ -9,6 +9,7 @@ import com.adx.tencent.conversionsync.ConversionSyncService;
 import com.adx.tencent.conversionsync.RetryService;
 import com.adx.tencent.httpapi.MediaPlacement;
 import com.adx.tencent.leader.LeaderElection;
+import com.adx.tencent.tagsync.TagEventSyncService;
 import com.adx.tencent.tencent.TencentClient;
 import com.adx.tencent.storage.RedisHotStore;
 import com.adx.tencent.storage.RedisLock;
@@ -34,6 +35,7 @@ import org.springframework.stereotype.Component;
 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;
@@ -203,6 +205,16 @@ public class AppConfiguration {
         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);
+    }
+
     // ─── ColdWorker ─────────────────────────────────────────────────────────
 
     @Bean
@@ -254,6 +266,7 @@ public class AppConfiguration {
         @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);
@@ -323,6 +336,28 @@ public class AppConfiguration {
                 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) {
@@ -359,6 +394,21 @@ public class AppConfiguration {
             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()) {

+ 14 - 0
src/main/java/com/adx/tencent/config/AppProperties.java

@@ -86,6 +86,11 @@ public class AppProperties {
     private Duration callbackRetryInterval = Duration.ofSeconds(300);
     private int callbackRetryLimit = 100;
 
+    // Tag Event Sync(广告位回传方式 -> Redis)
+    private boolean tagEventSyncEnabled = false;
+    private Duration tagEventSyncInterval = Duration.ofMinutes(5);
+    private String tagEventSyncRedisPrefix = "";
+
     private static String defaultHostname() {
         try {
             return java.net.InetAddress.getLocalHost().getHostName();
@@ -242,4 +247,13 @@ public class AppProperties {
 
     public int getCallbackRetryLimit() { return callbackRetryLimit; }
     public void setCallbackRetryLimit(int v) { this.callbackRetryLimit = v; }
+
+    public boolean isTagEventSyncEnabled() { return tagEventSyncEnabled; }
+    public void setTagEventSyncEnabled(boolean v) { this.tagEventSyncEnabled = v; }
+
+    public Duration getTagEventSyncInterval() { return tagEventSyncInterval; }
+    public void setTagEventSyncInterval(Duration v) { this.tagEventSyncInterval = v; }
+
+    public String getTagEventSyncRedisPrefix() { return tagEventSyncRedisPrefix; }
+    public void setTagEventSyncRedisPrefix(String v) { this.tagEventSyncRedisPrefix = v; }
 }

+ 15 - 0
src/main/java/com/adx/tencent/storage/TiDBColdStore.java

@@ -113,6 +113,12 @@ public class TiDBColdStore {
                 KEY idx_tencent_media_callbacks_qk_created_at (qk, created_at),
                 KEY idx_tencent_media_callbacks_media_event_created_at (media, event_type, created_at)
             )
+            """,
+            """
+            CREATE TABLE IF NOT EXISTS tencent_tag_event (
+                tag_id VARCHAR(32) NOT NULL PRIMARY KEY,
+                event_type VARCHAR(128) NOT NULL COMMENT '回传方式'
+            ) COMMENT '竞价事件记录'
             """
         };
         try (SqlSession session = sqlSessionFactory.openSession(true)) {
@@ -241,6 +247,15 @@ public class TiDBColdStore {
         }
     }
 
+    // ─── 查询:广告位回传方式配置 ──────────────────────────────────────────────
+
+    public List<TagEventRecord> listAllTagEvents() {
+        try (SqlSession session = sqlSessionFactory.openSession(true)) {
+            TagEventMapper mapper = session.getMapper(TagEventMapper.class);
+            return mapper.selectAll();
+        }
+    }
+
     // ─── 工具方法 ─────────────────────────────────────────────────────────────
 
     private String toJson(Object value) {

+ 15 - 0
src/main/java/com/adx/tencent/storage/mapper/TagEventMapper.java

@@ -0,0 +1,15 @@
+package com.adx.tencent.storage.mapper;
+
+import com.adx.tencent.storage.model.TagEventRecord;
+import org.apache.ibatis.annotations.Mapper;
+
+import java.util.List;
+
+@Mapper
+public interface TagEventMapper {
+
+    /**
+     * 查询全部广告位回传方式配置。
+     */
+    List<TagEventRecord> selectAll();
+}

+ 18 - 0
src/main/java/com/adx/tencent/storage/model/TagEventRecord.java

@@ -0,0 +1,18 @@
+package com.adx.tencent.storage.model;
+
+/**
+ * 广告位回传方式记录,对应 tencent_tag_event 表。
+ *   tag_id     广告位ID
+ *   event_type 回传方式
+ */
+public class TagEventRecord {
+
+    private String tagId;
+    private String eventType;
+
+    public String getTagId() { return tagId; }
+    public void setTagId(String v) { this.tagId = v; }
+
+    public String getEventType() { return eventType; }
+    public void setEventType(String v) { this.eventType = v; }
+}

+ 71 - 0
src/main/java/com/adx/tencent/tagsync/TagEventSyncService.java

@@ -0,0 +1,71 @@
+package com.adx.tencent.tagsync;
+
+import com.adx.tencent.storage.TiDBColdStore;
+import com.adx.tencent.storage.model.TagEventRecord;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.data.redis.connection.RedisConnection;
+import org.springframework.data.redis.connection.RedisStringCommands;
+import org.springframework.data.redis.core.StringRedisTemplate;
+import org.springframework.data.redis.core.types.Expiration;
+
+import java.nio.charset.StandardCharsets;
+import java.time.Duration;
+import java.util.List;
+
+/**
+ * 广告位回传方式同步服务。
+ * 定时从 tencent_tag_event 表读取全部记录,同步到 Redis:
+ *   key   = {prefix}{tag_id}(广告位ID)
+ *   value = event_type(回传方式)
+ * 为避免表中删除的记录在 Redis 残留,写入时带 TTL(大于同步间隔)。
+ */
+public class TagEventSyncService {
+
+    private static final Logger log = LoggerFactory.getLogger(TagEventSyncService.class);
+
+    private final TiDBColdStore coldStore;
+    private final StringRedisTemplate redis;
+    private final String keyPrefix;
+    private final Duration ttl;
+
+    public TagEventSyncService(TiDBColdStore coldStore, StringRedisTemplate redis,
+                               String keyPrefix, Duration ttl) {
+        this.coldStore = coldStore;
+        this.redis = redis;
+        this.keyPrefix = keyPrefix == null ? "" : keyPrefix;
+        this.ttl = ttl;
+    }
+
+    /**
+     * 执行一次同步,返回写入 Redis 的记录数。
+     */
+    public int syncOnce() {
+        List<TagEventRecord> rows = coldStore.listAllTagEvents();
+        if (rows == null || rows.isEmpty()) {
+            log.info("[TagEventSync] 无数据可同步");
+            return 0;
+        }
+        long ttlMs = (ttl != null && !ttl.isZero()) ? ttl.toMillis() : 0;
+
+        redis.executePipelined((RedisConnection conn) -> {
+            for (TagEventRecord row : rows) {
+                if (row.getTagId() == null || row.getTagId().isBlank()) continue;
+                byte[] key = (keyPrefix + row.getTagId()).getBytes(StandardCharsets.UTF_8);
+                byte[] val = (row.getEventType() == null ? "" : row.getEventType())
+                        .getBytes(StandardCharsets.UTF_8);
+                if (ttlMs > 0) {
+                    conn.stringCommands().set(key, val,
+                            Expiration.milliseconds(ttlMs),
+                            RedisStringCommands.SetOption.UPSERT);
+                } else {
+                    conn.stringCommands().set(key, val);
+                }
+            }
+            return null;
+        });
+
+        log.info("[TagEventSync] 同步完成,共 {} 条广告位回传方式写入 Redis", rows.size());
+        return rows.size();
+    }
+}

+ 4 - 0
src/main/resources/application-dev.yml

@@ -48,3 +48,7 @@ adx:
 
   # --- Callback Retry ---
   callback-retry-enabled: true
+
+  # --- Tag Event Sync(广告位回传方式 -> Redis)---
+  tag-event-sync-enabled: true
+  tag-event-sync-interval: 300s

+ 4 - 0
src/main/resources/application-prod.yml

@@ -53,3 +53,7 @@ adx:
   # --- Callback Retry ---
   callback-retry-enabled: true
   callback-retry-interval: 300s
+
+  # --- Tag Event Sync(广告位回传方式 -> Redis)---
+  tag-event-sync-enabled: true
+  tag-event-sync-interval: 300s

+ 4 - 0
src/main/resources/application-test.yml

@@ -52,3 +52,7 @@ adx:
   # --- Callback Retry ---
   callback-retry-enabled: true
   callback-retry-interval: 300s
+
+  # --- Tag Event Sync(广告位回传方式 -> Redis)---
+  tag-event-sync-enabled: true
+  tag-event-sync-interval: 300s

+ 13 - 0
src/main/resources/mapper/TagEventMapper.xml

@@ -0,0 +1,13 @@
+<?xml version="1.0" encoding="UTF-8" ?>
+<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN"
+        "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
+<mapper namespace="com.adx.tencent.storage.mapper.TagEventMapper">
+    <resultMap id="tagEventResultMap" type="com.adx.tencent.storage.model.TagEventRecord">
+        <result property="tagId" column="tag_id"/>
+        <result property="eventType" column="event_type"/>
+    </resultMap>
+
+    <select id="selectAll" resultMap="tagEventResultMap">
+        SELECT tag_id, event_type FROM tencent_tag_event
+    </select>
+</mapper>