TiDBColdStore.java 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442
  1. package com.adx.tencent.storage;
  2. import com.adx.tencent.storage.mapper.*;
  3. import com.adx.tencent.storage.model.*;
  4. import com.fasterxml.jackson.databind.ObjectMapper;
  5. import org.apache.ibatis.session.ExecutorType;
  6. import org.apache.ibatis.session.SqlSession;
  7. import org.apache.ibatis.session.SqlSessionFactory;
  8. import org.slf4j.Logger;
  9. import org.slf4j.LoggerFactory;
  10. import java.sql.Statement;
  11. import java.sql.Timestamp;
  12. import java.time.Instant;
  13. import java.util.List;
  14. /**
  15. * 对应 Go internal/storage/tidb.go 的 TiDBColdStore。
  16. * 使用 MyBatis 操作 TiDB(MySQL 协议兼容)。
  17. * 功能:
  18. * 1. migrate() - 建表 DDL
  19. * 2. saveBid / saveTracking / saveConversion / saveMediaCallback
  20. * 3. successfulCallbackExists / successfulCallbackExistsForEvent
  21. * 4. pendingMediaCallbacks
  22. */
  23. public class TiDBColdStore {
  24. private static final Logger log = LoggerFactory.getLogger(TiDBColdStore.class);
  25. private static final int CONVERSION_BATCH_SIZE = 1000;
  26. private final SqlSessionFactory sqlSessionFactory;
  27. private final ObjectMapper objectMapper;
  28. public TiDBColdStore(SqlSessionFactory sqlSessionFactory, ObjectMapper objectMapper) {
  29. this.sqlSessionFactory = sqlSessionFactory;
  30. this.objectMapper = objectMapper;
  31. }
  32. // ─── DDL 迁移 ────────────────────────────────────────────────────────────
  33. public void migrate() {
  34. String[] ddls = {
  35. """
  36. CREATE TABLE IF NOT EXISTS tencent_ad_bid_events (
  37. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  38. qk VARCHAR(128) NOT NULL,
  39. media VARCHAR(64),
  40. media_trace_id VARCHAR(255),
  41. platform VARCHAR(16),
  42. req_id VARCHAR(512),
  43. bid_id VARCHAR(128),
  44. creative_id VARCHAR(128),
  45. imp_id VARCHAR(128),
  46. tag_id VARCHAR(128),
  47. price BIGINT UNSIGNED,
  48. show_urls JSON,
  49. click_urls JSON,
  50. landing_page TEXT,
  51. app_store_link TEXT,
  52. package_name VARCHAR(255),
  53. media_params JSON,
  54. created_at TIMESTAMP(3) NOT NULL,
  55. UNIQUE KEY uk_tencent_ad_bid_events_qk (qk),
  56. KEY idx_tencent_ad_bid_events_media_trace (media, media_trace_id),
  57. KEY idx_tencent_ad_bid_events_platform_tag (platform, tag_id),
  58. KEY idx_tencent_ad_bid_events_created_at (created_at)
  59. )
  60. """,
  61. """
  62. CREATE TABLE IF NOT EXISTS tencent_tracking_reports (
  63. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  64. qk VARCHAR(128),
  65. kind VARCHAR(32) NOT NULL,
  66. tracking_url TEXT NOT NULL,
  67. status INT NOT NULL,
  68. ok TINYINT(1) NOT NULL,
  69. created_at TIMESTAMP(3) NOT NULL,
  70. KEY idx_tencent_tracking_reports_qk_kind_created_at (qk, kind, created_at)
  71. )
  72. """,
  73. """
  74. CREATE TABLE IF NOT EXISTS tencent_baidu_conversions (
  75. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  76. dedupe_key VARCHAR(255),
  77. qk VARCHAR(128),
  78. media VARCHAR(64),
  79. tag_id VARCHAR(128),
  80. date_value VARCHAR(16),
  81. appsid VARCHAR(128),
  82. customer_name VARCHAR(128),
  83. device_id VARCHAR(255),
  84. conv DOUBLE,
  85. payment DOUBLE,
  86. gmv DOUBLE,
  87. act INT,
  88. tu VARCHAR(128),
  89. clk_time VARCHAR(32),
  90. created_at TIMESTAMP(3) NOT NULL,
  91. UNIQUE KEY uk_tencent_baidu_conversions_dedupe (dedupe_key),
  92. KEY idx_tencent_baidu_conversions_qk_act_created_at (qk, act, created_at),
  93. KEY idx_tencent_baidu_conversions_tag_id (tag_id),
  94. KEY idx_tencent_baidu_conversions_device_id (device_id)
  95. )
  96. """,
  97. """
  98. CREATE TABLE IF NOT EXISTS tencent_media_callbacks (
  99. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  100. dedupe_key VARCHAR(255),
  101. media VARCHAR(64) NOT NULL,
  102. qk VARCHAR(128),
  103. callback_url TEXT NOT NULL,
  104. event_type INT NOT NULL,
  105. event_time_ms BIGINT NOT NULL,
  106. purchase DOUBLE,
  107. status INT NOT NULL,
  108. ok TINYINT(1) NOT NULL,
  109. attempt INT NOT NULL DEFAULT 1,
  110. response_body TEXT,
  111. error_message TEXT,
  112. request_body TEXT,
  113. dispatch_status VARCHAR(32) NOT NULL DEFAULT 'SENT',
  114. tracking_version VARCHAR(16) NOT NULL DEFAULT 'v1',
  115. created_at TIMESTAMP(3) NOT NULL,
  116. KEY idx_tencent_media_callbacks_dedupe_ok_created_at (dedupe_key, ok, created_at),
  117. KEY idx_tencent_media_callbacks_qk_created_at (qk, created_at),
  118. KEY idx_tencent_media_callbacks_media_event_created_at (media, event_type, created_at)
  119. )
  120. """,
  121. """
  122. CREATE TABLE IF NOT EXISTS tencent_tag_event (
  123. tag_id VARCHAR(32) NOT NULL,
  124. baidu_act INT NOT NULL COMMENT '百度转化行为',
  125. tencent_action_type VARCHAR(128) NOT NULL COMMENT '腾讯回传行为',
  126. PRIMARY KEY (tag_id, baidu_act)
  127. ) COMMENT '竞价事件记录'
  128. """
  129. };
  130. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  131. Statement stmt = session.getConnection().createStatement();
  132. for (String ddl : ddls) {
  133. stmt.execute(ddl);
  134. }
  135. // 字段长度补丁(已存在的表可能是旧的 128)
  136. try {
  137. stmt.execute("ALTER TABLE tencent_ad_bid_events MODIFY COLUMN req_id VARCHAR(512)");
  138. } catch (Exception ignored) {}
  139. // platform 字段补丁
  140. try {
  141. stmt.execute("ALTER TABLE tencent_ad_bid_events ADD COLUMN platform VARCHAR(16) AFTER media_trace_id");
  142. } catch (Exception ignored) {}
  143. try {
  144. stmt.execute("ALTER TABLE tencent_ad_bid_events ADD KEY idx_tencent_ad_bid_events_platform_tag (platform, tag_id)");
  145. } catch (Exception ignored) {}
  146. try {
  147. stmt.execute("ALTER TABLE tencent_baidu_conversions ADD COLUMN gmv DOUBLE AFTER payment");
  148. } catch (Exception ignored) {}
  149. try {
  150. stmt.execute("ALTER TABLE tencent_baidu_conversions ADD COLUMN tag_id VARCHAR(128) AFTER media");
  151. } catch (Exception ignored) {}
  152. try {
  153. stmt.execute("ALTER TABLE tencent_baidu_conversions ADD KEY idx_tencent_baidu_conversions_tag_id (tag_id)");
  154. } catch (Exception ignored) {}
  155. try {
  156. stmt.execute("ALTER TABLE tencent_media_callbacks ADD COLUMN request_body TEXT AFTER error_message");
  157. } catch (Exception ignored) {}
  158. try {
  159. stmt.execute("ALTER TABLE tencent_media_callbacks ADD COLUMN dispatch_status VARCHAR(32) NOT NULL DEFAULT 'SENT' AFTER request_body");
  160. } catch (Exception ignored) {}
  161. try {
  162. stmt.execute("ALTER TABLE tencent_tracking_reports ADD KEY idx_tencent_tracking_reports_created_at_id (created_at, id)");
  163. } catch (Exception ignored) {}
  164. try {
  165. stmt.execute("UPDATE tencent_media_callbacks SET dispatch_status = 'SENT' WHERE dispatch_status IS NULL OR dispatch_status = ''");
  166. } catch (Exception ignored) {}
  167. try {
  168. stmt.execute("ALTER TABLE tencent_media_callbacks ADD COLUMN tracking_version VARCHAR(16) NOT NULL DEFAULT 'v1' AFTER dispatch_status");
  169. } catch (Exception ignored) {}
  170. try {
  171. stmt.execute("UPDATE tencent_media_callbacks SET tracking_version = 'v1' WHERE tracking_version IS NULL OR tracking_version = ''");
  172. } catch (Exception ignored) {}
  173. try {
  174. stmt.execute("ALTER TABLE tencent_tag_event ADD COLUMN baidu_act INT NULL COMMENT '百度转化行为' AFTER tag_id");
  175. } catch (Exception ignored) {}
  176. try {
  177. stmt.execute("ALTER TABLE tencent_tag_event CHANGE COLUMN event_type tencent_action_type VARCHAR(128) NOT NULL COMMENT '腾讯回传行为'");
  178. } catch (Exception ignored) {}
  179. try {
  180. stmt.execute("UPDATE tencent_tag_event SET baidu_act = 0 WHERE baidu_act IS NULL");
  181. } catch (Exception ignored) {}
  182. try {
  183. stmt.execute("ALTER TABLE tencent_tag_event MODIFY COLUMN baidu_act INT NOT NULL");
  184. } catch (Exception ignored) {}
  185. try {
  186. stmt.execute("ALTER TABLE tencent_tag_event DROP PRIMARY KEY");
  187. } catch (Exception ignored) {}
  188. try {
  189. stmt.execute("ALTER TABLE tencent_tag_event ADD PRIMARY KEY (tag_id, baidu_act)");
  190. } catch (Exception ignored) {}
  191. stmt.close();
  192. } catch (Exception e) {
  193. throw new RuntimeException("migrate failed", e);
  194. }
  195. }
  196. // ─── SaveBid ─────────────────────────────────────────────────────────────
  197. public void saveBid(BidRecord record) {
  198. if (record.getQk() == null || record.getQk().isBlank()) {
  199. throw new IllegalArgumentException("qk is required");
  200. }
  201. if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now());
  202. record.setMediaTraceId(record.getMediaTraceId());
  203. BidRecord cold = record.toColdRecord();
  204. String mediaParams = toJson(cold.getMediaParams());
  205. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  206. BidEventMapper mapper = session.getMapper(BidEventMapper.class);
  207. mapper.insertOrUpdate(
  208. cold.getQk(), cold.getMedia(), cold.getMediaTraceId(),
  209. cold.getPlatform(), cold.getAccountId(), cold.getEventType(),
  210. cold.getTagId(), cold.getPrice() != null ? cold.getPrice() : 0L,
  211. cold.getPackageName(), mediaParams, Timestamp.from(cold.getCreatedAt())
  212. );
  213. }
  214. }
  215. public BidRecord getBidByQk(String qk) {
  216. if (qk == null || qk.isBlank()) return null;
  217. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  218. BidEventMapper mapper = session.getMapper(BidEventMapper.class);
  219. return mapper.selectByQk(qk);
  220. }
  221. }
  222. public List<BidRecord> getBidsByQks(List<String> qks) {
  223. if (qks == null || qks.isEmpty()) return List.of();
  224. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  225. BidEventMapper mapper = session.getMapper(BidEventMapper.class);
  226. return mapper.selectByQks(qks);
  227. }
  228. }
  229. // ─── SaveTracking ─────────────────────────────────────────────────────────
  230. public void saveTracking(TrackingRecord record) {
  231. if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now());
  232. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  233. TrackingReportMapper mapper = session.getMapper(TrackingReportMapper.class);
  234. mapper.insert(
  235. record.getQk(), record.getKind(), record.getUrl(),
  236. record.getStatus(), record.isOk() ? 1 : 0,
  237. Timestamp.from(record.getCreatedAt())
  238. );
  239. }
  240. }
  241. // ─── SaveConversion ──────────────────────────────────────────────────────
  242. public void saveConversion(ConversionRecord record) {
  243. if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now());
  244. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  245. ConversionMapper mapper = session.getMapper(ConversionMapper.class);
  246. mapper.insertOrUpdate(
  247. record.getDedupeKey(), record.getQk(), record.getMedia(), record.getTagId(),
  248. record.getDate(), record.getAppSid(), record.getCustomerName(),
  249. record.getDeviceId(), record.getConv(), record.getPayment(), record.getGmv(),
  250. record.getAct(), record.getTu(), record.getClkTime(),
  251. Timestamp.from(record.getCreatedAt())
  252. );
  253. }
  254. }
  255. public void saveConversions(List<ConversionRecord> records) {
  256. if (records == null || records.isEmpty()) return;
  257. long startedAt = System.currentTimeMillis();
  258. Instant now = Instant.now();
  259. for (ConversionRecord record : records) {
  260. if (record.getCreatedAt() == null) record.setCreatedAt(now);
  261. }
  262. try (SqlSession session = sqlSessionFactory.openSession(false)) {
  263. ConversionMapper mapper = session.getMapper(ConversionMapper.class);
  264. for (int start = 0; start < records.size(); start += CONVERSION_BATCH_SIZE) {
  265. int end = Math.min(start + CONVERSION_BATCH_SIZE, records.size());
  266. mapper.batchInsertOrUpdate(records.subList(start, end));
  267. }
  268. session.commit();
  269. long costMs = System.currentTimeMillis() - startedAt;
  270. int batchCount = (records.size() + CONVERSION_BATCH_SIZE - 1) / CONVERSION_BATCH_SIZE;
  271. log.info("[TiDBColdStore] saveConversions done: size={}, batches={}, batchSize={}, cost={}ms",
  272. records.size(), batchCount, CONVERSION_BATCH_SIZE, costMs);
  273. }
  274. }
  275. // ─── SaveMediaCallback ───────────────────────────────────────────────────
  276. public void saveMediaCallback(MediaCallbackRecord record) {
  277. if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now());
  278. int attempt = record.getAttempt() <= 0 ? 1 : record.getAttempt();
  279. if (record.getDispatchStatus() == null || record.getDispatchStatus().isBlank()) {
  280. record.setDispatchStatus(MediaCallbackRecord.DISPATCH_STATUS_SENT);
  281. }
  282. if (record.getTrackingVersion() == null || record.getTrackingVersion().isBlank()) {
  283. record.setTrackingVersion("v1");
  284. }
  285. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  286. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  287. mapper.insert(
  288. record.getDedupeKey(), record.getMedia(), record.getQk(),
  289. record.getCallbackUrl(), record.getEventType(), record.getEventTimeMs(),
  290. record.getPurchase(), record.getStatus(), record.isOk() ? 1 : 0,
  291. attempt, record.getResponseBody(), record.getErrorMessage(),
  292. record.getRequestBody(), record.getDispatchStatus(), record.getTrackingVersion(),
  293. Timestamp.from(record.getCreatedAt())
  294. );
  295. }
  296. }
  297. public void saveMediaCallbacks(List<MediaCallbackRecord> records) {
  298. if (records == null || records.isEmpty()) return;
  299. for (MediaCallbackRecord record : records) {
  300. if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now());
  301. if (record.getAttempt() <= 0) record.setAttempt(1);
  302. if (record.getDispatchStatus() == null || record.getDispatchStatus().isBlank()) {
  303. record.setDispatchStatus(MediaCallbackRecord.DISPATCH_STATUS_SENT);
  304. }
  305. if (record.getTrackingVersion() == null || record.getTrackingVersion().isBlank()) {
  306. record.setTrackingVersion("v1");
  307. }
  308. }
  309. try (SqlSession session = sqlSessionFactory.openSession(ExecutorType.BATCH, false)) {
  310. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  311. mapper.batchInsert(records);
  312. session.commit();
  313. }
  314. }
  315. // ─── 查询:成功回调是否存在 ────────────────────────────────────────────────
  316. public boolean successfulCallbackExists(String dedupeKey) {
  317. if (dedupeKey == null || dedupeKey.isBlank()) return false;
  318. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  319. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  320. int count = mapper.countSuccessfulByDedupeKey(dedupeKey);
  321. return count > 0;
  322. }
  323. }
  324. public boolean successfulCallbackExistsForEvent(String dedupeKey, String legacyDedupeKey, int eventType) {
  325. if (dedupeKey == null || dedupeKey.isBlank()) return false;
  326. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  327. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  328. int count = mapper.countSuccessfulForEvent(dedupeKey, legacyDedupeKey, eventType);
  329. return count > 0;
  330. }
  331. }
  332. public List<String> successfulCallbackKeys(List<String> dedupeKeys) {
  333. if (dedupeKeys == null || dedupeKeys.isEmpty()) return List.of();
  334. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  335. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  336. return mapper.selectSuccessfulDedupeKeys(dedupeKeys);
  337. }
  338. }
  339. public List<String> successfulCallbackKeysWithLegacy(List<String> dedupeKeys,
  340. List<String> legacyDedupeKeys,
  341. List<Integer> eventTypes) {
  342. if (dedupeKeys == null || dedupeKeys.isEmpty()) return List.of();
  343. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  344. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  345. return mapper.selectSuccessfulKeysWithLegacy(dedupeKeys, legacyDedupeKeys, eventTypes);
  346. }
  347. }
  348. // ─── 查询:待重试的回调 ────────────────────────────────────────────────────
  349. public List<MediaCallbackRecord> pendingMediaCallbacks(String media, int limit) {
  350. if (limit <= 0) limit = 100;
  351. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  352. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  353. return mapper.selectPendingCallbacks(media, limit);
  354. }
  355. }
  356. public MediaCallbackRecord getMediaCallbackById(long id) {
  357. if (id <= 0) return null;
  358. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  359. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  360. return mapper.selectById(id);
  361. }
  362. }
  363. public void updateMediaCallback(MediaCallbackRecord record) {
  364. if (record == null || record.getId() == null || record.getId() <= 0) {
  365. throw new IllegalArgumentException("media callback id is required");
  366. }
  367. int attempt = record.getAttempt() <= 0 ? 1 : record.getAttempt();
  368. record.setAttempt(attempt);
  369. if (record.getDispatchStatus() == null || record.getDispatchStatus().isBlank()) {
  370. record.setDispatchStatus(MediaCallbackRecord.DISPATCH_STATUS_SENT);
  371. }
  372. if (record.getTrackingVersion() == null || record.getTrackingVersion().isBlank()) {
  373. record.setTrackingVersion("v1");
  374. }
  375. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  376. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  377. mapper.updateById(record);
  378. }
  379. }
  380. // ─── 查询:广告位回传方式配置 ──────────────────────────────────────────────
  381. public List<TagEventRecord> listAllTagEvents() {
  382. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  383. TagEventMapper mapper = session.getMapper(TagEventMapper.class);
  384. return mapper.selectAll();
  385. }
  386. }
  387. /**
  388. * 根据广告位ID查询回传方式配置(DB 兜底查询)。
  389. */
  390. public TagEventRecord getTagEvent(String tagId, int baiduAct) {
  391. if (tagId == null || tagId.isBlank()) return null;
  392. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  393. TagEventMapper mapper = session.getMapper(TagEventMapper.class);
  394. return mapper.selectByTagIdAndAct(tagId, baiduAct);
  395. }
  396. }
  397. // ─── 工具方法 ─────────────────────────────────────────────────────────────
  398. private String toJson(Object value) {
  399. if (value == null) return "null";
  400. try {
  401. return objectMapper.writeValueAsString(value);
  402. } catch (Exception e) {
  403. throw new RuntimeException("JSON serialize failed", e);
  404. }
  405. }
  406. }