package com.adx.tencent.storage; import com.adx.tencent.storage.mapper.*; import com.adx.tencent.storage.model.*; import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.ibatis.session.ExecutorType; import org.apache.ibatis.session.SqlSession; import org.apache.ibatis.session.SqlSessionFactory; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.sql.Statement; import java.sql.Timestamp; import java.time.Instant; import java.util.List; /** * 对应 Go internal/storage/tidb.go 的 TiDBColdStore。 * 使用 MyBatis 操作 TiDB(MySQL 协议兼容)。 * 功能: * 1. migrate() - 建表 DDL * 2. saveBid / saveTracking / saveConversion / saveMediaCallback * 3. successfulCallbackExists / successfulCallbackExistsForEvent * 4. pendingMediaCallbacks */ public class TiDBColdStore { private static final Logger log = LoggerFactory.getLogger(TiDBColdStore.class); private static final int CONVERSION_BATCH_SIZE = 1000; private final SqlSessionFactory sqlSessionFactory; private final ObjectMapper objectMapper; public TiDBColdStore(SqlSessionFactory sqlSessionFactory, ObjectMapper objectMapper) { this.sqlSessionFactory = sqlSessionFactory; this.objectMapper = objectMapper; } // ─── DDL 迁移 ──────────────────────────────────────────────────────────── public void migrate() { String[] ddls = { """ CREATE TABLE IF NOT EXISTS tencent_ad_bid_events ( id BIGINT AUTO_INCREMENT PRIMARY KEY, qk VARCHAR(128) NOT NULL, media VARCHAR(64), media_trace_id VARCHAR(255), platform VARCHAR(16), req_id VARCHAR(512), bid_id VARCHAR(128), creative_id VARCHAR(128), imp_id VARCHAR(128), tag_id VARCHAR(128), price BIGINT UNSIGNED, show_urls JSON, click_urls JSON, landing_page TEXT, app_store_link TEXT, package_name VARCHAR(255), media_params JSON, created_at TIMESTAMP(3) NOT NULL, UNIQUE KEY uk_tencent_ad_bid_events_qk (qk), KEY idx_tencent_ad_bid_events_media_trace (media, media_trace_id), KEY idx_tencent_ad_bid_events_platform_tag (platform, tag_id), KEY idx_tencent_ad_bid_events_created_at (created_at) ) """, """ CREATE TABLE IF NOT EXISTS tencent_tracking_reports ( id BIGINT AUTO_INCREMENT PRIMARY KEY, qk VARCHAR(128), kind VARCHAR(32) NOT NULL, tracking_url TEXT NOT NULL, status INT NOT NULL, ok TINYINT(1) NOT NULL, created_at TIMESTAMP(3) NOT NULL, KEY idx_tencent_tracking_reports_qk_kind_created_at (qk, kind, created_at) ) """, """ CREATE TABLE IF NOT EXISTS tencent_baidu_conversions ( id BIGINT AUTO_INCREMENT PRIMARY KEY, dedupe_key VARCHAR(255), qk VARCHAR(128), media VARCHAR(64), tag_id VARCHAR(128), date_value VARCHAR(16), appsid VARCHAR(128), customer_name VARCHAR(128), device_id VARCHAR(255), conv DOUBLE, payment DOUBLE, gmv DOUBLE, act INT, tu VARCHAR(128), clk_time VARCHAR(32), created_at TIMESTAMP(3) NOT NULL, UNIQUE KEY uk_tencent_baidu_conversions_dedupe (dedupe_key), KEY idx_tencent_baidu_conversions_qk_act_created_at (qk, act, created_at), KEY idx_tencent_baidu_conversions_tag_id (tag_id), KEY idx_tencent_baidu_conversions_device_id (device_id) ) """, """ CREATE TABLE IF NOT EXISTS tencent_media_callbacks ( id BIGINT AUTO_INCREMENT PRIMARY KEY, dedupe_key VARCHAR(255), media VARCHAR(64) NOT NULL, qk VARCHAR(128), callback_url TEXT NOT NULL, event_type INT NOT NULL, event_time_ms BIGINT NOT NULL, purchase DOUBLE, status INT NOT NULL, ok TINYINT(1) NOT NULL, attempt INT NOT NULL DEFAULT 1, response_body TEXT, error_message TEXT, request_body TEXT, dispatch_status VARCHAR(32) NOT NULL DEFAULT 'SENT', tracking_version VARCHAR(16) NOT NULL DEFAULT 'v1', created_at TIMESTAMP(3) NOT NULL, KEY idx_tencent_media_callbacks_dedupe_ok_created_at (dedupe_key, ok, created_at), 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, baidu_act INT NOT NULL COMMENT '百度转化行为', tencent_action_type VARCHAR(128) NOT NULL COMMENT '腾讯回传行为', PRIMARY KEY (tag_id, baidu_act) ) COMMENT '竞价事件记录' """ }; try (SqlSession session = sqlSessionFactory.openSession(true)) { Statement stmt = session.getConnection().createStatement(); for (String ddl : ddls) { stmt.execute(ddl); } // 字段长度补丁(已存在的表可能是旧的 128) try { stmt.execute("ALTER TABLE tencent_ad_bid_events MODIFY COLUMN req_id VARCHAR(512)"); } catch (Exception ignored) {} // platform 字段补丁 try { stmt.execute("ALTER TABLE tencent_ad_bid_events ADD COLUMN platform VARCHAR(16) AFTER media_trace_id"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_ad_bid_events ADD KEY idx_tencent_ad_bid_events_platform_tag (platform, tag_id)"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_baidu_conversions ADD COLUMN gmv DOUBLE AFTER payment"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_baidu_conversions ADD COLUMN tag_id VARCHAR(128) AFTER media"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_baidu_conversions ADD KEY idx_tencent_baidu_conversions_tag_id (tag_id)"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_media_callbacks ADD COLUMN request_body TEXT AFTER error_message"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_media_callbacks ADD COLUMN dispatch_status VARCHAR(32) NOT NULL DEFAULT 'SENT' AFTER request_body"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_tracking_reports ADD KEY idx_tencent_tracking_reports_created_at_id (created_at, id)"); } catch (Exception ignored) {} try { stmt.execute("UPDATE tencent_media_callbacks SET dispatch_status = 'SENT' WHERE dispatch_status IS NULL OR dispatch_status = ''"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_media_callbacks ADD COLUMN tracking_version VARCHAR(16) NOT NULL DEFAULT 'v1' AFTER dispatch_status"); } catch (Exception ignored) {} try { stmt.execute("UPDATE tencent_media_callbacks SET tracking_version = 'v1' WHERE tracking_version IS NULL OR tracking_version = ''"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_tag_event ADD COLUMN baidu_act INT NULL COMMENT '百度转化行为' AFTER tag_id"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_tag_event CHANGE COLUMN event_type tencent_action_type VARCHAR(128) NOT NULL COMMENT '腾讯回传行为'"); } catch (Exception ignored) {} try { stmt.execute("UPDATE tencent_tag_event SET baidu_act = 0 WHERE baidu_act IS NULL"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_tag_event MODIFY COLUMN baidu_act INT NOT NULL"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_tag_event DROP PRIMARY KEY"); } catch (Exception ignored) {} try { stmt.execute("ALTER TABLE tencent_tag_event ADD PRIMARY KEY (tag_id, baidu_act)"); } catch (Exception ignored) {} stmt.close(); } catch (Exception e) { throw new RuntimeException("migrate failed", e); } } // ─── SaveBid ───────────────────────────────────────────────────────────── public void saveBid(BidRecord record) { if (record.getQk() == null || record.getQk().isBlank()) { throw new IllegalArgumentException("qk is required"); } if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now()); record.setMediaTraceId(record.getMediaTraceId()); BidRecord cold = record.toColdRecord(); String mediaParams = toJson(cold.getMediaParams()); try (SqlSession session = sqlSessionFactory.openSession(true)) { BidEventMapper mapper = session.getMapper(BidEventMapper.class); mapper.insertOrUpdate( cold.getQk(), cold.getMedia(), cold.getMediaTraceId(), cold.getPlatform(), cold.getAccountId(), cold.getEventType(), cold.getTagId(), cold.getPrice() != null ? cold.getPrice() : 0L, cold.getPackageName(), mediaParams, Timestamp.from(cold.getCreatedAt()) ); } } public BidRecord getBidByQk(String qk) { if (qk == null || qk.isBlank()) return null; try (SqlSession session = sqlSessionFactory.openSession(true)) { BidEventMapper mapper = session.getMapper(BidEventMapper.class); return mapper.selectByQk(qk); } } public List getBidsByQks(List qks) { if (qks == null || qks.isEmpty()) return List.of(); try (SqlSession session = sqlSessionFactory.openSession(true)) { BidEventMapper mapper = session.getMapper(BidEventMapper.class); return mapper.selectByQks(qks); } } // ─── SaveTracking ───────────────────────────────────────────────────────── public void saveTracking(TrackingRecord record) { if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now()); try (SqlSession session = sqlSessionFactory.openSession(true)) { TrackingReportMapper mapper = session.getMapper(TrackingReportMapper.class); mapper.insert( record.getQk(), record.getKind(), record.getUrl(), record.getStatus(), record.isOk() ? 1 : 0, Timestamp.from(record.getCreatedAt()) ); } } // ─── SaveConversion ────────────────────────────────────────────────────── public void saveConversion(ConversionRecord record) { if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now()); try (SqlSession session = sqlSessionFactory.openSession(true)) { ConversionMapper mapper = session.getMapper(ConversionMapper.class); mapper.insertOrUpdate( record.getDedupeKey(), record.getQk(), record.getMedia(), record.getTagId(), record.getDate(), record.getAppSid(), record.getCustomerName(), record.getDeviceId(), record.getConv(), record.getPayment(), record.getGmv(), record.getAct(), record.getTu(), record.getClkTime(), Timestamp.from(record.getCreatedAt()) ); } } public void saveConversions(List records) { if (records == null || records.isEmpty()) return; long startedAt = System.currentTimeMillis(); Instant now = Instant.now(); for (ConversionRecord record : records) { if (record.getCreatedAt() == null) record.setCreatedAt(now); } try (SqlSession session = sqlSessionFactory.openSession(false)) { ConversionMapper mapper = session.getMapper(ConversionMapper.class); for (int start = 0; start < records.size(); start += CONVERSION_BATCH_SIZE) { int end = Math.min(start + CONVERSION_BATCH_SIZE, records.size()); mapper.batchInsertOrUpdate(records.subList(start, end)); } session.commit(); long costMs = System.currentTimeMillis() - startedAt; int batchCount = (records.size() + CONVERSION_BATCH_SIZE - 1) / CONVERSION_BATCH_SIZE; log.info("[TiDBColdStore] saveConversions done: size={}, batches={}, batchSize={}, cost={}ms", records.size(), batchCount, CONVERSION_BATCH_SIZE, costMs); } } // ─── SaveMediaCallback ─────────────────────────────────────────────────── public void saveMediaCallback(MediaCallbackRecord record) { if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now()); int attempt = record.getAttempt() <= 0 ? 1 : record.getAttempt(); if (record.getDispatchStatus() == null || record.getDispatchStatus().isBlank()) { record.setDispatchStatus(MediaCallbackRecord.DISPATCH_STATUS_SENT); } if (record.getTrackingVersion() == null || record.getTrackingVersion().isBlank()) { record.setTrackingVersion("v1"); } try (SqlSession session = sqlSessionFactory.openSession(true)) { MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class); mapper.insert( record.getDedupeKey(), record.getMedia(), record.getQk(), record.getCallbackUrl(), record.getEventType(), record.getEventTimeMs(), record.getPurchase(), record.getStatus(), record.isOk() ? 1 : 0, attempt, record.getResponseBody(), record.getErrorMessage(), record.getRequestBody(), record.getDispatchStatus(), record.getTrackingVersion(), Timestamp.from(record.getCreatedAt()) ); } } public void saveMediaCallbacks(List records) { if (records == null || records.isEmpty()) return; for (MediaCallbackRecord record : records) { if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now()); if (record.getAttempt() <= 0) record.setAttempt(1); if (record.getDispatchStatus() == null || record.getDispatchStatus().isBlank()) { record.setDispatchStatus(MediaCallbackRecord.DISPATCH_STATUS_SENT); } if (record.getTrackingVersion() == null || record.getTrackingVersion().isBlank()) { record.setTrackingVersion("v1"); } } try (SqlSession session = sqlSessionFactory.openSession(ExecutorType.BATCH, false)) { MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class); mapper.batchInsert(records); session.commit(); } } // ─── 查询:成功回调是否存在 ──────────────────────────────────────────────── public boolean successfulCallbackExists(String dedupeKey) { if (dedupeKey == null || dedupeKey.isBlank()) return false; try (SqlSession session = sqlSessionFactory.openSession(true)) { MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class); int count = mapper.countSuccessfulByDedupeKey(dedupeKey); return count > 0; } } public boolean successfulCallbackExistsForEvent(String dedupeKey, String legacyDedupeKey, int eventType) { if (dedupeKey == null || dedupeKey.isBlank()) return false; try (SqlSession session = sqlSessionFactory.openSession(true)) { MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class); int count = mapper.countSuccessfulForEvent(dedupeKey, legacyDedupeKey, eventType); return count > 0; } } public List successfulCallbackKeys(List dedupeKeys) { if (dedupeKeys == null || dedupeKeys.isEmpty()) return List.of(); try (SqlSession session = sqlSessionFactory.openSession(true)) { MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class); return mapper.selectSuccessfulDedupeKeys(dedupeKeys); } } public List successfulCallbackKeysWithLegacy(List dedupeKeys, List legacyDedupeKeys, List eventTypes) { if (dedupeKeys == null || dedupeKeys.isEmpty()) return List.of(); try (SqlSession session = sqlSessionFactory.openSession(true)) { MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class); return mapper.selectSuccessfulKeysWithLegacy(dedupeKeys, legacyDedupeKeys, eventTypes); } } // ─── 查询:待重试的回调 ──────────────────────────────────────────────────── public List pendingMediaCallbacks(String media, int limit) { if (limit <= 0) limit = 100; try (SqlSession session = sqlSessionFactory.openSession(true)) { MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class); return mapper.selectPendingCallbacks(media, limit); } } public MediaCallbackRecord getMediaCallbackById(long id) { if (id <= 0) return null; try (SqlSession session = sqlSessionFactory.openSession(true)) { MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class); return mapper.selectById(id); } } public void updateMediaCallback(MediaCallbackRecord record) { if (record == null || record.getId() == null || record.getId() <= 0) { throw new IllegalArgumentException("media callback id is required"); } int attempt = record.getAttempt() <= 0 ? 1 : record.getAttempt(); record.setAttempt(attempt); if (record.getDispatchStatus() == null || record.getDispatchStatus().isBlank()) { record.setDispatchStatus(MediaCallbackRecord.DISPATCH_STATUS_SENT); } if (record.getTrackingVersion() == null || record.getTrackingVersion().isBlank()) { record.setTrackingVersion("v1"); } try (SqlSession session = sqlSessionFactory.openSession(true)) { MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class); mapper.updateById(record); } } // ─── 查询:广告位回传方式配置 ────────────────────────────────────────────── public List listAllTagEvents() { try (SqlSession session = sqlSessionFactory.openSession(true)) { TagEventMapper mapper = session.getMapper(TagEventMapper.class); return mapper.selectAll(); } } /** * 根据广告位ID查询回传方式配置(DB 兜底查询)。 */ public TagEventRecord getTagEvent(String tagId, int baiduAct) { if (tagId == null || tagId.isBlank()) return null; try (SqlSession session = sqlSessionFactory.openSession(true)) { TagEventMapper mapper = session.getMapper(TagEventMapper.class); return mapper.selectByTagIdAndAct(tagId, baiduAct); } } // ─── 工具方法 ───────────────────────────────────────────────────────────── private String toJson(Object value) { if (value == null) return "null"; try { return objectMapper.writeValueAsString(value); } catch (Exception e) { throw new RuntimeException("JSON serialize failed", e); } } }