apply-lobster-profile-migration-all-tenants.js 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380
  1. /**
  2. * Apply lobster profile stats / summary projection / segment seed / menu remove to all active tenant DBs.
  3. *
  4. * Usage:
  5. * node apply-lobster-profile-migration-all-tenants.js
  6. * node apply-lobster-profile-migration-all-tenants.js --tenant-id=33
  7. */
  8. const fs = require('fs')
  9. const path = require('path')
  10. const mysql = require('mysql2/promise')
  11. const MASTER = {
  12. host: process.env.DB_HOST || 'cq-cdb-8fjmemkb.sql.tencentcdb.com',
  13. port: Number(process.env.DB_PORT || 27220),
  14. user: process.env.DB_USER || 'root',
  15. password: process.env.DB_PASSWORD || 'Ylrz_1q2w3e4r5t6y',
  16. database: process.env.DB_NAME || 'ylrz_saas'
  17. }
  18. const SQL_DIR = path.join(__dirname, '../sql')
  19. const PROFILE_DDL = fs.readFileSync(path.join(SQL_DIR, 'lobster_user_profile.sql'), 'utf8')
  20. const STATS_DDL = fs.readFileSync(path.join(SQL_DIR, 'lobster_profile_stats.sql'), 'utf8')
  21. const SEGMENT_DDL = fs.readFileSync(path.join(__dirname, '../fs-agent/src/main/resources/db/migration/tenant/V20260625_02__lobster_user_segment.sql'), 'utf8')
  22. const SEGMENT_SQL = fs.readFileSync(path.join(SQL_DIR, 'V20260722_01__lobster_segment_threshold_config.sql'), 'utf8')
  23. const MENU_SQL = fs.readFileSync(path.join(SQL_DIR, 'V20260711_01__lobster_remove_conversation_summary_menu.sql'), 'utf8')
  24. const PROFILE_COLUMNS = [
  25. ['channel_type', "VARCHAR(20) NOT NULL DEFAULT 'QW'"],
  26. ['latest_stage', "VARCHAR(64) DEFAULT NULL"],
  27. ['latest_customer_attitude', "VARCHAR(32) DEFAULT NULL"],
  28. ['latest_summary_text', "VARCHAR(512) DEFAULT NULL"],
  29. ['profile_aggregate_preview', "VARCHAR(420) DEFAULT NULL"],
  30. ['latest_summary_time', "DATETIME DEFAULT NULL"]
  31. ]
  32. const SUMMARY_COLUMNS = [
  33. ['external_user_id', "varchar(128) DEFAULT NULL COMMENT '\u5916\u90e8\u7528\u6237ID'"],
  34. ['channel_type', "varchar(20) NOT NULL DEFAULT 'QW' COMMENT '\u6e20\u9053\u7c7b\u578b QW/WX/IM'"],
  35. ['summary_text', "text DEFAULT NULL COMMENT '\u6458\u8981\u6587\u672c\u5185\u5bb9'"],
  36. ['customer_attitude', "varchar(32) DEFAULT NULL COMMENT '\u5ba2\u6237\u6001\u5ea6'"],
  37. ['stage', "varchar(64) DEFAULT NULL COMMENT '\u5bf9\u8bdd\u9636\u6bb5'"]
  38. ]
  39. const PROFILE_INDEXES = [
  40. ['idx_profile_latest_stage', '(company_id, latest_stage)'],
  41. ['idx_profile_latest_attitude', '(company_id, latest_customer_attitude)'],
  42. ['idx_profile_latest_summary_time', '(company_id, latest_summary_time)']
  43. ]
  44. function parseArgs() {
  45. const tenantIds = []
  46. for (const arg of process.argv.slice(2)) {
  47. if (arg.startsWith('--tenant-id=')) {
  48. tenantIds.push(Number(arg.split('=')[1]))
  49. }
  50. }
  51. return tenantIds
  52. }
  53. function parseJdbcUrl(jdbcUrl, dbName, user, password) {
  54. if (jdbcUrl && jdbcUrl.startsWith('jdbc:mysql://')) {
  55. const body = jdbcUrl.slice('jdbc:mysql://'.length)
  56. const slash = body.indexOf('/')
  57. const hostPort = slash >= 0 ? body.slice(0, slash) : body
  58. const rest = slash >= 0 ? body.slice(slash + 1) : dbName
  59. const q = rest.indexOf('?')
  60. const database = q >= 0 ? rest.slice(0, q) : rest
  61. const colon = hostPort.indexOf(':')
  62. return {
  63. host: colon >= 0 ? hostPort.slice(0, colon) : hostPort,
  64. port: colon >= 0 ? Number(hostPort.slice(colon + 1)) : MASTER.port,
  65. database,
  66. user: user || MASTER.user,
  67. password: password || MASTER.password
  68. }
  69. }
  70. return {
  71. host: MASTER.host,
  72. port: MASTER.port,
  73. database: dbName,
  74. user: user || MASTER.user,
  75. password: password || MASTER.password
  76. }
  77. }
  78. async function tableExists(conn, table) {
  79. const [rows] = await conn.query(
  80. 'SELECT COUNT(*) AS cnt FROM information_schema.TABLES WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ?',
  81. [table]
  82. )
  83. return rows[0].cnt > 0
  84. }
  85. async function columnExists(conn, table, column) {
  86. const [rows] = await conn.query(
  87. 'SELECT COUNT(*) AS cnt FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND COLUMN_NAME = ?',
  88. [table, column]
  89. )
  90. return rows[0].cnt > 0
  91. }
  92. async function indexExists(conn, table, indexName) {
  93. const [rows] = await conn.query(
  94. 'SELECT COUNT(*) AS cnt FROM information_schema.STATISTICS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND INDEX_NAME = ?',
  95. [table, indexName]
  96. )
  97. return rows[0].cnt > 0
  98. }
  99. async function addColumnIfMissing(conn, table, column, definition) {
  100. if (await columnExists(conn, table, column)) return false
  101. await conn.query(`ALTER TABLE \`${table}\` ADD COLUMN \`${column}\` ${definition}`)
  102. return true
  103. }
  104. async function addIndexIfMissing(conn, table, indexName, columns) {
  105. if (await indexExists(conn, table, indexName)) return false
  106. await conn.query(`ALTER TABLE \`${table}\` ADD INDEX \`${indexName}\` ${columns}`)
  107. return true
  108. }
  109. async function ensureProfileSchema(conn) {
  110. if (!(await tableExists(conn, 'lobster_user_profile'))) {
  111. await conn.query(PROFILE_DDL)
  112. console.log(' created lobster_user_profile')
  113. }
  114. await conn.query(STATS_DDL)
  115. for (const [col, def] of PROFILE_COLUMNS) {
  116. if (await addColumnIfMissing(conn, 'lobster_user_profile', col, def)) {
  117. console.log(` added lobster_user_profile.${col}`)
  118. }
  119. }
  120. for (const [idx, cols] of PROFILE_INDEXES) {
  121. if (await addIndexIfMissing(conn, 'lobster_user_profile', idx, cols)) {
  122. console.log(` added index ${idx}`)
  123. }
  124. }
  125. }
  126. async function ensureSummaryColumns(conn) {
  127. if (!(await tableExists(conn, 'lobster_conversation_summary'))) return false
  128. for (const [col, def] of SUMMARY_COLUMNS) {
  129. await addColumnIfMissing(conn, 'lobster_conversation_summary', col, def)
  130. }
  131. return true
  132. }
  133. async function backfillSummaryProjection(conn) {
  134. if (!(await tableExists(conn, 'lobster_user_profile'))) return
  135. if (!(await tableExists(conn, 'lobster_conversation_summary'))) return
  136. if (!(await columnExists(conn, 'lobster_conversation_summary', 'external_user_id'))) {
  137. console.log(' skip summary backfill: external_user_id missing')
  138. return
  139. }
  140. const hasStage = await columnExists(conn, 'lobster_conversation_summary', 'stage')
  141. const hasAttitude = await columnExists(conn, 'lobster_conversation_summary', 'customer_attitude')
  142. const hasSummaryText = await columnExists(conn, 'lobster_conversation_summary', 'summary_text')
  143. const hasSummaryContent = await columnExists(conn, 'lobster_conversation_summary', 'summary_content')
  144. const hasDelFlag = await columnExists(conn, 'lobster_conversation_summary', 'del_flag')
  145. const stageExpr = hasStage ? 's.stage' : 'NULL'
  146. const attitudeExpr = hasAttitude ? 's.customer_attitude' : 'NULL'
  147. let summaryExpr = 'NULL'
  148. if (hasSummaryText && hasSummaryContent) {
  149. summaryExpr = "COALESCE(NULLIF(TRIM(s.summary_text), ''), NULLIF(TRIM(s.summary_content), ''))"
  150. } else if (hasSummaryText) {
  151. summaryExpr = "NULLIF(TRIM(s.summary_text), '')"
  152. } else if (hasSummaryContent) {
  153. summaryExpr = "NULLIF(TRIM(s.summary_content), '')"
  154. }
  155. const hasProfileChannel = await columnExists(conn, 'lobster_user_profile', 'channel_type')
  156. const joinChannel = hasProfileChannel ? ' AND p.channel_type = src.channel_type' : ''
  157. const channelTypeExpr = hasProfileChannel ? 's.channel_type' : "'QW'"
  158. const delFilter = hasDelFlag ? '(del_flag IS NULL OR del_flag = 0)' : '1=1'
  159. const sql = `
  160. UPDATE lobster_user_profile p
  161. INNER JOIN (
  162. SELECT s.company_id, s.external_user_id, ${channelTypeExpr} AS channel_type,
  163. ${stageExpr} AS stage,
  164. ${attitudeExpr} AS customer_attitude,
  165. ${summaryExpr} AS summary_text,
  166. COALESCE(s.update_time, s.create_time) AS summary_time
  167. FROM lobster_conversation_summary s
  168. INNER JOIN (
  169. SELECT company_id, external_user_id, channel_type,
  170. MAX(COALESCE(update_time, create_time)) AS max_time
  171. FROM lobster_conversation_summary
  172. WHERE ${delFilter}
  173. AND external_user_id IS NOT NULL
  174. AND company_id IS NOT NULL
  175. GROUP BY company_id, external_user_id, channel_type
  176. ) latest
  177. ON latest.company_id = s.company_id
  178. AND latest.external_user_id = s.external_user_id
  179. AND latest.channel_type = s.channel_type
  180. AND latest.max_time = COALESCE(s.update_time, s.create_time)
  181. WHERE ${delFilter}
  182. ) src
  183. ON p.company_id = src.company_id
  184. AND p.external_user_id = src.external_user_id${joinChannel}
  185. SET p.latest_stage = COALESCE(src.stage, p.latest_stage),
  186. p.latest_customer_attitude = COALESCE(src.customer_attitude, p.latest_customer_attitude),
  187. p.latest_summary_text = COALESCE(LEFT(src.summary_text, 512), p.latest_summary_text),
  188. p.profile_aggregate_preview = CASE
  189. WHEN src.stage IS NOT NULL AND src.summary_text IS NOT NULL AND TRIM(src.summary_text) != ''
  190. THEN LEFT(CONCAT('\u3010\u9636\u6bb5\u3011', src.stage, '\uff5c', TRIM(src.summary_text)), 420)
  191. WHEN src.stage IS NOT NULL
  192. THEN CONCAT('\u3010\u9636\u6bb5\u3011', src.stage)
  193. WHEN src.summary_text IS NOT NULL AND TRIM(src.summary_text) != ''
  194. THEN LEFT(TRIM(src.summary_text), 420)
  195. ELSE p.profile_aggregate_preview
  196. END,
  197. p.latest_summary_time = COALESCE(src.summary_time, p.latest_summary_time)
  198. WHERE p.deleted = 0
  199. AND (p.latest_summary_text IS NULL OR p.latest_summary_time IS NULL)`
  200. const [result] = await conn.query(sql)
  201. console.log(` summary backfill affectedRows=${result.affectedRows || 0}`)
  202. }
  203. async function backfillStatsCache(conn) {
  204. if (!(await tableExists(conn, 'lobster_user_profile'))) return
  205. const companySql = `
  206. INSERT INTO lobster_profile_stats
  207. (company_id, total_profiles, enriched_count, interaction_sum, avg_conversations, active_today, refresh_time)
  208. SELECT
  209. p.company_id,
  210. COUNT(1),
  211. IFNULL(SUM(CASE WHEN p.value_score > 10 OR p.total_purchase > 0
  212. OR (p.internal_tags IS NOT NULL AND p.internal_tags != '') THEN 1 ELSE 0 END), 0),
  213. IFNULL(SUM(p.interaction_count), 0),
  214. IFNULL(ROUND(AVG(p.interaction_count), 2), 0),
  215. IFNULL(SUM(CASE WHEN DATE(p.last_active_time) = CURDATE() THEN 1 ELSE 0 END), 0),
  216. NOW()
  217. FROM lobster_user_profile p
  218. WHERE p.deleted = 0
  219. GROUP BY p.company_id
  220. ON DUPLICATE KEY UPDATE
  221. total_profiles = VALUES(total_profiles),
  222. enriched_count = VALUES(enriched_count),
  223. interaction_sum = VALUES(interaction_sum),
  224. avg_conversations = VALUES(avg_conversations),
  225. active_today = VALUES(active_today),
  226. refresh_time = NOW()`
  227. const tenantSql = `
  228. INSERT INTO lobster_profile_stats
  229. (company_id, total_profiles, enriched_count, interaction_sum, avg_conversations, active_today, refresh_time)
  230. SELECT
  231. 0,
  232. COUNT(1),
  233. IFNULL(SUM(CASE WHEN p.value_score > 10 OR p.total_purchase > 0
  234. OR (p.internal_tags IS NOT NULL AND p.internal_tags != '') THEN 1 ELSE 0 END), 0),
  235. IFNULL(SUM(p.interaction_count), 0),
  236. IFNULL(ROUND(AVG(p.interaction_count), 2), 0),
  237. IFNULL(SUM(CASE WHEN DATE(p.last_active_time) = CURDATE() THEN 1 ELSE 0 END), 0),
  238. NOW()
  239. FROM lobster_user_profile p
  240. WHERE p.deleted = 0
  241. ON DUPLICATE KEY UPDATE
  242. total_profiles = VALUES(total_profiles),
  243. enriched_count = VALUES(enriched_count),
  244. interaction_sum = VALUES(interaction_sum),
  245. avg_conversations = VALUES(avg_conversations),
  246. active_today = VALUES(active_today),
  247. refresh_time = NOW()`
  248. await conn.query(companySql)
  249. await conn.query(tenantSql)
  250. console.log(' stats cache backfill done')
  251. }
  252. async function ensureUserSegmentTable(conn) {
  253. if (!(await tableExists(conn, 'lobster_user_segment'))) {
  254. await conn.query(SEGMENT_DDL)
  255. console.log(' created lobster_user_segment')
  256. }
  257. }
  258. async function applySegmentSeed(conn) {
  259. await ensureUserSegmentTable(conn)
  260. await conn.query(SEGMENT_SQL)
  261. console.log(' segment seed applied')
  262. }
  263. async function removeSummaryMenu(conn) {
  264. const statements = MENU_SQL.split(';').map(s => s.trim()).filter(Boolean)
  265. for (const sql of statements) {
  266. await conn.query(sql)
  267. }
  268. console.log(' conversation-summary menu removed')
  269. }
  270. async function verifyTenant(conn) {
  271. return {
  272. lobster_profile_stats: await tableExists(conn, 'lobster_profile_stats'),
  273. latest_stage: await columnExists(conn, 'lobster_user_profile', 'latest_stage'),
  274. latest_summary_text: await columnExists(conn, 'lobster_user_profile', 'latest_summary_text'),
  275. interaction_sum: await columnExists(conn, 'lobster_profile_stats', 'interaction_sum'),
  276. builtinSegmentRows: (await tableExists(conn, 'lobster_user_segment'))
  277. ? (await conn.query("SELECT COUNT(*) AS cnt FROM lobster_user_segment WHERE deleted = 0 AND segment_code IN ('NEW','ACTIVE','SLEEP','DORMANT','CHURN')"))[0][0].cnt
  278. : 0,
  279. statsCacheRows: (await tableExists(conn, 'lobster_profile_stats'))
  280. ? (await conn.query('SELECT COUNT(*) AS cnt FROM lobster_profile_stats'))[0][0].cnt
  281. : 0
  282. }
  283. }
  284. async function applyTenant(tenant) {
  285. const conn = await mysql.createConnection({
  286. host: tenant.host,
  287. port: tenant.port,
  288. user: tenant.user,
  289. password: tenant.password,
  290. database: tenant.database,
  291. multipleStatements: true,
  292. charset: 'utf8mb4'
  293. })
  294. try {
  295. console.log(`\n[RUN] tenantId=${tenant.id} db=${tenant.database}`)
  296. await ensureProfileSchema(conn)
  297. await ensureSummaryColumns(conn)
  298. await backfillSummaryProjection(conn)
  299. await backfillStatsCache(conn)
  300. await applySegmentSeed(conn)
  301. await removeSummaryMenu(conn)
  302. const checks = await verifyTenant(conn)
  303. console.log(` verify: ${JSON.stringify(checks)}`)
  304. } finally {
  305. await conn.end()
  306. }
  307. }
  308. async function main() {
  309. const filterIds = parseArgs()
  310. const masterConn = await mysql.createConnection({ ...MASTER, charset: 'utf8mb4' })
  311. const tenants = []
  312. try {
  313. const [rows] = await masterConn.query(
  314. 'SELECT id, db_name, db_url, db_account, db_pwd, status FROM tenant_info ORDER BY id'
  315. )
  316. for (const row of rows) {
  317. if (row.status !== 1) continue
  318. if (filterIds.length > 0 && !filterIds.includes(Number(row.id))) continue
  319. tenants.push({
  320. id: row.id,
  321. status: row.status,
  322. ...parseJdbcUrl(row.db_url, row.db_name, row.db_account, row.db_pwd)
  323. })
  324. }
  325. } finally {
  326. await masterConn.end()
  327. }
  328. console.log(`Active tenants to migrate: ${tenants.length}`)
  329. let ok = 0
  330. let fail = 0
  331. for (const tenant of tenants) {
  332. try {
  333. await applyTenant(tenant)
  334. ok++
  335. } catch (err) {
  336. fail++
  337. console.error(`[FAIL] tenantId=${tenant.id} db=${tenant.database}: ${err.message}`)
  338. }
  339. }
  340. console.log(`\nDone. ok=${ok} fail=${fail}`)
  341. if (fail > 0) process.exit(1)
  342. }
  343. main().catch(err => {
  344. console.error(err)
  345. process.exit(1)
  346. })