TiDBColdStore.java 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304
  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.SqlSession;
  6. import org.apache.ibatis.session.SqlSessionFactory;
  7. import java.sql.Statement;
  8. import java.sql.Timestamp;
  9. import java.time.Instant;
  10. import java.util.List;
  11. /**
  12. * 对应 Go internal/storage/tidb.go 的 TiDBColdStore。
  13. * 使用 MyBatis 操作 TiDB(MySQL 协议兼容)。
  14. * 功能:
  15. * 1. migrate() - 建表 DDL
  16. * 2. saveBid / saveTracking / saveConversion / saveMediaCallback
  17. * 3. successfulCallbackExists / successfulCallbackExistsForEvent
  18. * 4. pendingMediaCallbacks
  19. */
  20. public class TiDBColdStore {
  21. private final SqlSessionFactory sqlSessionFactory;
  22. private final ObjectMapper objectMapper;
  23. public TiDBColdStore(SqlSessionFactory sqlSessionFactory, ObjectMapper objectMapper) {
  24. this.sqlSessionFactory = sqlSessionFactory;
  25. this.objectMapper = objectMapper;
  26. }
  27. // ─── DDL 迁移 ────────────────────────────────────────────────────────────
  28. public void migrate() {
  29. String[] ddls = {
  30. """
  31. CREATE TABLE IF NOT EXISTS tencent_ad_bid_events (
  32. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  33. qk VARCHAR(128) NOT NULL,
  34. media VARCHAR(64),
  35. media_trace_id VARCHAR(255),
  36. platform VARCHAR(16),
  37. req_id VARCHAR(512),
  38. bid_id VARCHAR(128),
  39. creative_id VARCHAR(128),
  40. imp_id VARCHAR(128),
  41. tag_id VARCHAR(128),
  42. price BIGINT UNSIGNED,
  43. show_urls JSON,
  44. click_urls JSON,
  45. landing_page TEXT,
  46. app_store_link TEXT,
  47. package_name VARCHAR(255),
  48. media_params JSON,
  49. created_at TIMESTAMP(3) NOT NULL,
  50. UNIQUE KEY uk_tencent_ad_bid_events_qk (qk),
  51. KEY idx_tencent_ad_bid_events_media_trace (media, media_trace_id),
  52. KEY idx_tencent_ad_bid_events_platform_tag (platform, tag_id),
  53. KEY idx_tencent_ad_bid_events_created_at (created_at)
  54. )
  55. """,
  56. """
  57. CREATE TABLE IF NOT EXISTS tencent_tracking_reports (
  58. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  59. qk VARCHAR(128),
  60. kind VARCHAR(32) NOT NULL,
  61. tracking_url TEXT NOT NULL,
  62. status INT NOT NULL,
  63. ok TINYINT(1) NOT NULL,
  64. created_at TIMESTAMP(3) NOT NULL,
  65. KEY idx_tencent_tracking_reports_qk_kind_created_at (qk, kind, created_at)
  66. )
  67. """,
  68. """
  69. CREATE TABLE IF NOT EXISTS tencent_baidu_conversions (
  70. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  71. dedupe_key VARCHAR(255),
  72. qk VARCHAR(128),
  73. media VARCHAR(64),
  74. date_value VARCHAR(16),
  75. appsid VARCHAR(128),
  76. customer_name VARCHAR(128),
  77. device_id VARCHAR(255),
  78. conv DOUBLE,
  79. payment DOUBLE,
  80. gmv DOUBLE,
  81. act INT,
  82. tu VARCHAR(128),
  83. clk_time VARCHAR(32),
  84. created_at TIMESTAMP(3) NOT NULL,
  85. UNIQUE KEY uk_tencent_baidu_conversions_dedupe (dedupe_key),
  86. KEY idx_tencent_baidu_conversions_qk_act_created_at (qk, act, created_at),
  87. KEY idx_tencent_baidu_conversions_device_id (device_id)
  88. )
  89. """,
  90. """
  91. CREATE TABLE IF NOT EXISTS tencent_media_callbacks (
  92. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  93. dedupe_key VARCHAR(255),
  94. media VARCHAR(64) NOT NULL,
  95. qk VARCHAR(128),
  96. callback_url TEXT NOT NULL,
  97. event_type INT NOT NULL,
  98. event_time_ms BIGINT NOT NULL,
  99. purchase DOUBLE,
  100. status INT NOT NULL,
  101. ok TINYINT(1) NOT NULL,
  102. attempt INT NOT NULL DEFAULT 1,
  103. response_body TEXT,
  104. error_message TEXT,
  105. created_at TIMESTAMP(3) NOT NULL,
  106. KEY idx_tencent_media_callbacks_dedupe_ok_created_at (dedupe_key, ok, created_at),
  107. KEY idx_tencent_media_callbacks_qk_created_at (qk, created_at),
  108. KEY idx_tencent_media_callbacks_media_event_created_at (media, event_type, created_at)
  109. )
  110. """,
  111. """
  112. CREATE TABLE IF NOT EXISTS tencent_tag_event (
  113. tag_id VARCHAR(32) NOT NULL,
  114. baidu_act INT NOT NULL COMMENT '百度转化行为',
  115. tencent_action_type VARCHAR(128) NOT NULL COMMENT '腾讯回传行为',
  116. PRIMARY KEY (tag_id, baidu_act)
  117. ) COMMENT '竞价事件记录'
  118. """
  119. };
  120. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  121. Statement stmt = session.getConnection().createStatement();
  122. for (String ddl : ddls) {
  123. stmt.execute(ddl);
  124. }
  125. // 字段长度补丁(已存在的表可能是旧的 128)
  126. try {
  127. stmt.execute("ALTER TABLE tencent_ad_bid_events MODIFY COLUMN req_id VARCHAR(512)");
  128. } catch (Exception ignored) {}
  129. // platform 字段补丁
  130. try {
  131. stmt.execute("ALTER TABLE tencent_ad_bid_events ADD COLUMN platform VARCHAR(16) AFTER media_trace_id");
  132. } catch (Exception ignored) {}
  133. try {
  134. stmt.execute("ALTER TABLE tencent_ad_bid_events ADD KEY idx_tencent_ad_bid_events_platform_tag (platform, tag_id)");
  135. } catch (Exception ignored) {}
  136. try {
  137. stmt.execute("ALTER TABLE tencent_baidu_conversions ADD COLUMN gmv DOUBLE AFTER payment");
  138. } catch (Exception ignored) {}
  139. try {
  140. stmt.execute("ALTER TABLE tencent_tag_event ADD COLUMN baidu_act INT NULL COMMENT '百度转化行为' AFTER tag_id");
  141. } catch (Exception ignored) {}
  142. try {
  143. stmt.execute("ALTER TABLE tencent_tag_event CHANGE COLUMN event_type tencent_action_type VARCHAR(128) NOT NULL COMMENT '腾讯回传行为'");
  144. } catch (Exception ignored) {}
  145. try {
  146. stmt.execute("UPDATE tencent_tag_event SET baidu_act = 0 WHERE baidu_act IS NULL");
  147. } catch (Exception ignored) {}
  148. try {
  149. stmt.execute("ALTER TABLE tencent_tag_event MODIFY COLUMN baidu_act INT NOT NULL");
  150. } catch (Exception ignored) {}
  151. try {
  152. stmt.execute("ALTER TABLE tencent_tag_event DROP PRIMARY KEY");
  153. } catch (Exception ignored) {}
  154. try {
  155. stmt.execute("ALTER TABLE tencent_tag_event ADD PRIMARY KEY (tag_id, baidu_act)");
  156. } catch (Exception ignored) {}
  157. stmt.close();
  158. } catch (Exception e) {
  159. throw new RuntimeException("migrate failed", e);
  160. }
  161. }
  162. // ─── SaveBid ─────────────────────────────────────────────────────────────
  163. public void saveBid(BidRecord record) {
  164. if (record.getQk() == null || record.getQk().isBlank()) {
  165. throw new IllegalArgumentException("qk is required");
  166. }
  167. if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now());
  168. String showUrls = toJson(record.getShowUrls());
  169. String clickUrls = toJson(record.getClickUrls());
  170. String mediaParams = toJson(record.getMediaParams());
  171. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  172. BidEventMapper mapper = session.getMapper(BidEventMapper.class);
  173. mapper.insertOrUpdate(
  174. record.getQk(), record.getMedia(), record.getMediaTraceId(),
  175. record.getPlatform(), record.getAdId(), record.getAccountId(),
  176. record.getEventType(),
  177. record.getReqId(), record.getBidId(), record.getCreativeId(),
  178. record.getImpId(), record.getTagId(), record.getPrice() != null ? record.getPrice() : 0L,
  179. showUrls, clickUrls, record.getLandingPage(),
  180. record.getAppStoreLink(), record.getPackageName(),
  181. mediaParams, Timestamp.from(record.getCreatedAt())
  182. );
  183. }
  184. }
  185. // ─── SaveTracking ─────────────────────────────────────────────────────────
  186. public void saveTracking(TrackingRecord record) {
  187. if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now());
  188. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  189. TrackingReportMapper mapper = session.getMapper(TrackingReportMapper.class);
  190. mapper.insert(
  191. record.getQk(), record.getKind(), record.getUrl(),
  192. record.getStatus(), record.isOk() ? 1 : 0,
  193. Timestamp.from(record.getCreatedAt())
  194. );
  195. }
  196. }
  197. // ─── SaveConversion ──────────────────────────────────────────────────────
  198. public void saveConversion(ConversionRecord record) {
  199. if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now());
  200. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  201. ConversionMapper mapper = session.getMapper(ConversionMapper.class);
  202. mapper.insertOrUpdate(
  203. record.getDedupeKey(), record.getQk(), record.getMedia(),
  204. record.getDate(), record.getAppSid(), record.getCustomerName(),
  205. record.getDeviceId(), record.getConv(), record.getPayment(), record.getGmv(),
  206. record.getAct(), record.getTu(), record.getClkTime(),
  207. Timestamp.from(record.getCreatedAt())
  208. );
  209. }
  210. }
  211. // ─── SaveMediaCallback ───────────────────────────────────────────────────
  212. public void saveMediaCallback(MediaCallbackRecord record) {
  213. if (record.getCreatedAt() == null) record.setCreatedAt(Instant.now());
  214. int attempt = record.getAttempt() <= 0 ? 1 : record.getAttempt();
  215. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  216. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  217. mapper.insert(
  218. record.getDedupeKey(), record.getMedia(), record.getQk(),
  219. record.getCallbackUrl(), record.getEventType(), record.getEventTimeMs(),
  220. record.getPurchase(), record.getStatus(), record.isOk() ? 1 : 0,
  221. attempt, record.getResponseBody(), record.getErrorMessage(),
  222. record.getRequestBody(), Timestamp.from(record.getCreatedAt())
  223. );
  224. }
  225. }
  226. // ─── 查询:成功回调是否存在 ────────────────────────────────────────────────
  227. public boolean successfulCallbackExists(String dedupeKey) {
  228. if (dedupeKey == null || dedupeKey.isBlank()) return false;
  229. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  230. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  231. int count = mapper.countSuccessfulByDedupeKey(dedupeKey);
  232. return count > 0;
  233. }
  234. }
  235. public boolean successfulCallbackExistsForEvent(String dedupeKey, String legacyDedupeKey, int eventType) {
  236. if (dedupeKey == null || dedupeKey.isBlank()) return false;
  237. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  238. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  239. int count = mapper.countSuccessfulForEvent(dedupeKey, legacyDedupeKey, eventType);
  240. return count > 0;
  241. }
  242. }
  243. // ─── 查询:待重试的回调 ────────────────────────────────────────────────────
  244. public List<MediaCallbackRecord> pendingMediaCallbacks(String media, int limit) {
  245. if (limit <= 0) limit = 100;
  246. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  247. MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
  248. return mapper.selectPendingCallbacks(media, limit);
  249. }
  250. }
  251. // ─── 查询:广告位回传方式配置 ──────────────────────────────────────────────
  252. public List<TagEventRecord> listAllTagEvents() {
  253. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  254. TagEventMapper mapper = session.getMapper(TagEventMapper.class);
  255. return mapper.selectAll();
  256. }
  257. }
  258. /**
  259. * 根据广告位ID查询回传方式配置(DB 兜底查询)。
  260. */
  261. public TagEventRecord getTagEvent(String tagId, int baiduAct) {
  262. if (tagId == null || tagId.isBlank()) return null;
  263. try (SqlSession session = sqlSessionFactory.openSession(true)) {
  264. TagEventMapper mapper = session.getMapper(TagEventMapper.class);
  265. return mapper.selectByTagIdAndAct(tagId, baiduAct);
  266. }
  267. }
  268. // ─── 工具方法 ─────────────────────────────────────────────────────────────
  269. private String toJson(Object value) {
  270. if (value == null) return "null";
  271. try {
  272. return objectMapper.writeValueAsString(value);
  273. } catch (Exception e) {
  274. throw new RuntimeException("JSON serialize failed", e);
  275. }
  276. }
  277. }