| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380 |
- /**
- * Apply lobster profile stats / summary projection / segment seed / menu remove to all active tenant DBs.
- *
- * Usage:
- * node apply-lobster-profile-migration-all-tenants.js
- * node apply-lobster-profile-migration-all-tenants.js --tenant-id=33
- */
- const fs = require('fs')
- const path = require('path')
- const mysql = require('mysql2/promise')
- const MASTER = {
- host: process.env.DB_HOST || 'cq-cdb-8fjmemkb.sql.tencentcdb.com',
- port: Number(process.env.DB_PORT || 27220),
- user: process.env.DB_USER || 'root',
- password: process.env.DB_PASSWORD || 'Ylrz_1q2w3e4r5t6y',
- database: process.env.DB_NAME || 'ylrz_saas'
- }
- const SQL_DIR = path.join(__dirname, '../sql')
- const PROFILE_DDL = fs.readFileSync(path.join(SQL_DIR, 'lobster_user_profile.sql'), 'utf8')
- const STATS_DDL = fs.readFileSync(path.join(SQL_DIR, 'lobster_profile_stats.sql'), 'utf8')
- const SEGMENT_DDL = fs.readFileSync(path.join(__dirname, '../fs-agent/src/main/resources/db/migration/tenant/V20260625_02__lobster_user_segment.sql'), 'utf8')
- const SEGMENT_SQL = fs.readFileSync(path.join(SQL_DIR, 'V20260722_01__lobster_segment_threshold_config.sql'), 'utf8')
- const MENU_SQL = fs.readFileSync(path.join(SQL_DIR, 'V20260711_01__lobster_remove_conversation_summary_menu.sql'), 'utf8')
- const PROFILE_COLUMNS = [
- ['channel_type', "VARCHAR(20) NOT NULL DEFAULT 'QW'"],
- ['latest_stage', "VARCHAR(64) DEFAULT NULL"],
- ['latest_customer_attitude', "VARCHAR(32) DEFAULT NULL"],
- ['latest_summary_text', "VARCHAR(512) DEFAULT NULL"],
- ['profile_aggregate_preview', "VARCHAR(420) DEFAULT NULL"],
- ['latest_summary_time', "DATETIME DEFAULT NULL"]
- ]
- const SUMMARY_COLUMNS = [
- ['external_user_id', "varchar(128) DEFAULT NULL COMMENT '\u5916\u90e8\u7528\u6237ID'"],
- ['channel_type', "varchar(20) NOT NULL DEFAULT 'QW' COMMENT '\u6e20\u9053\u7c7b\u578b QW/WX/IM'"],
- ['summary_text', "text DEFAULT NULL COMMENT '\u6458\u8981\u6587\u672c\u5185\u5bb9'"],
- ['customer_attitude', "varchar(32) DEFAULT NULL COMMENT '\u5ba2\u6237\u6001\u5ea6'"],
- ['stage', "varchar(64) DEFAULT NULL COMMENT '\u5bf9\u8bdd\u9636\u6bb5'"]
- ]
- const PROFILE_INDEXES = [
- ['idx_profile_latest_stage', '(company_id, latest_stage)'],
- ['idx_profile_latest_attitude', '(company_id, latest_customer_attitude)'],
- ['idx_profile_latest_summary_time', '(company_id, latest_summary_time)']
- ]
- function parseArgs() {
- const tenantIds = []
- for (const arg of process.argv.slice(2)) {
- if (arg.startsWith('--tenant-id=')) {
- tenantIds.push(Number(arg.split('=')[1]))
- }
- }
- return tenantIds
- }
- function parseJdbcUrl(jdbcUrl, dbName, user, password) {
- if (jdbcUrl && jdbcUrl.startsWith('jdbc:mysql://')) {
- const body = jdbcUrl.slice('jdbc:mysql://'.length)
- const slash = body.indexOf('/')
- const hostPort = slash >= 0 ? body.slice(0, slash) : body
- const rest = slash >= 0 ? body.slice(slash + 1) : dbName
- const q = rest.indexOf('?')
- const database = q >= 0 ? rest.slice(0, q) : rest
- const colon = hostPort.indexOf(':')
- return {
- host: colon >= 0 ? hostPort.slice(0, colon) : hostPort,
- port: colon >= 0 ? Number(hostPort.slice(colon + 1)) : MASTER.port,
- database,
- user: user || MASTER.user,
- password: password || MASTER.password
- }
- }
- return {
- host: MASTER.host,
- port: MASTER.port,
- database: dbName,
- user: user || MASTER.user,
- password: password || MASTER.password
- }
- }
- async function tableExists(conn, table) {
- const [rows] = await conn.query(
- 'SELECT COUNT(*) AS cnt FROM information_schema.TABLES WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ?',
- [table]
- )
- return rows[0].cnt > 0
- }
- async function columnExists(conn, table, column) {
- const [rows] = await conn.query(
- 'SELECT COUNT(*) AS cnt FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND COLUMN_NAME = ?',
- [table, column]
- )
- return rows[0].cnt > 0
- }
- async function indexExists(conn, table, indexName) {
- const [rows] = await conn.query(
- 'SELECT COUNT(*) AS cnt FROM information_schema.STATISTICS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND INDEX_NAME = ?',
- [table, indexName]
- )
- return rows[0].cnt > 0
- }
- async function addColumnIfMissing(conn, table, column, definition) {
- if (await columnExists(conn, table, column)) return false
- await conn.query(`ALTER TABLE \`${table}\` ADD COLUMN \`${column}\` ${definition}`)
- return true
- }
- async function addIndexIfMissing(conn, table, indexName, columns) {
- if (await indexExists(conn, table, indexName)) return false
- await conn.query(`ALTER TABLE \`${table}\` ADD INDEX \`${indexName}\` ${columns}`)
- return true
- }
- async function ensureProfileSchema(conn) {
- if (!(await tableExists(conn, 'lobster_user_profile'))) {
- await conn.query(PROFILE_DDL)
- console.log(' created lobster_user_profile')
- }
- await conn.query(STATS_DDL)
- for (const [col, def] of PROFILE_COLUMNS) {
- if (await addColumnIfMissing(conn, 'lobster_user_profile', col, def)) {
- console.log(` added lobster_user_profile.${col}`)
- }
- }
- for (const [idx, cols] of PROFILE_INDEXES) {
- if (await addIndexIfMissing(conn, 'lobster_user_profile', idx, cols)) {
- console.log(` added index ${idx}`)
- }
- }
- }
- async function ensureSummaryColumns(conn) {
- if (!(await tableExists(conn, 'lobster_conversation_summary'))) return false
- for (const [col, def] of SUMMARY_COLUMNS) {
- await addColumnIfMissing(conn, 'lobster_conversation_summary', col, def)
- }
- return true
- }
- async function backfillSummaryProjection(conn) {
- if (!(await tableExists(conn, 'lobster_user_profile'))) return
- if (!(await tableExists(conn, 'lobster_conversation_summary'))) return
- if (!(await columnExists(conn, 'lobster_conversation_summary', 'external_user_id'))) {
- console.log(' skip summary backfill: external_user_id missing')
- return
- }
- const hasStage = await columnExists(conn, 'lobster_conversation_summary', 'stage')
- const hasAttitude = await columnExists(conn, 'lobster_conversation_summary', 'customer_attitude')
- const hasSummaryText = await columnExists(conn, 'lobster_conversation_summary', 'summary_text')
- const hasSummaryContent = await columnExists(conn, 'lobster_conversation_summary', 'summary_content')
- const hasDelFlag = await columnExists(conn, 'lobster_conversation_summary', 'del_flag')
- const stageExpr = hasStage ? 's.stage' : 'NULL'
- const attitudeExpr = hasAttitude ? 's.customer_attitude' : 'NULL'
- let summaryExpr = 'NULL'
- if (hasSummaryText && hasSummaryContent) {
- summaryExpr = "COALESCE(NULLIF(TRIM(s.summary_text), ''), NULLIF(TRIM(s.summary_content), ''))"
- } else if (hasSummaryText) {
- summaryExpr = "NULLIF(TRIM(s.summary_text), '')"
- } else if (hasSummaryContent) {
- summaryExpr = "NULLIF(TRIM(s.summary_content), '')"
- }
- const hasProfileChannel = await columnExists(conn, 'lobster_user_profile', 'channel_type')
- const joinChannel = hasProfileChannel ? ' AND p.channel_type = src.channel_type' : ''
- const channelTypeExpr = hasProfileChannel ? 's.channel_type' : "'QW'"
- const delFilter = hasDelFlag ? '(del_flag IS NULL OR del_flag = 0)' : '1=1'
- const sql = `
- UPDATE lobster_user_profile p
- INNER JOIN (
- SELECT s.company_id, s.external_user_id, ${channelTypeExpr} AS channel_type,
- ${stageExpr} AS stage,
- ${attitudeExpr} AS customer_attitude,
- ${summaryExpr} AS summary_text,
- COALESCE(s.update_time, s.create_time) AS summary_time
- FROM lobster_conversation_summary s
- INNER JOIN (
- SELECT company_id, external_user_id, channel_type,
- MAX(COALESCE(update_time, create_time)) AS max_time
- FROM lobster_conversation_summary
- WHERE ${delFilter}
- AND external_user_id IS NOT NULL
- AND company_id IS NOT NULL
- GROUP BY company_id, external_user_id, channel_type
- ) latest
- ON latest.company_id = s.company_id
- AND latest.external_user_id = s.external_user_id
- AND latest.channel_type = s.channel_type
- AND latest.max_time = COALESCE(s.update_time, s.create_time)
- WHERE ${delFilter}
- ) src
- ON p.company_id = src.company_id
- AND p.external_user_id = src.external_user_id${joinChannel}
- SET p.latest_stage = COALESCE(src.stage, p.latest_stage),
- p.latest_customer_attitude = COALESCE(src.customer_attitude, p.latest_customer_attitude),
- p.latest_summary_text = COALESCE(LEFT(src.summary_text, 512), p.latest_summary_text),
- p.profile_aggregate_preview = CASE
- WHEN src.stage IS NOT NULL AND src.summary_text IS NOT NULL AND TRIM(src.summary_text) != ''
- THEN LEFT(CONCAT('\u3010\u9636\u6bb5\u3011', src.stage, '\uff5c', TRIM(src.summary_text)), 420)
- WHEN src.stage IS NOT NULL
- THEN CONCAT('\u3010\u9636\u6bb5\u3011', src.stage)
- WHEN src.summary_text IS NOT NULL AND TRIM(src.summary_text) != ''
- THEN LEFT(TRIM(src.summary_text), 420)
- ELSE p.profile_aggregate_preview
- END,
- p.latest_summary_time = COALESCE(src.summary_time, p.latest_summary_time)
- WHERE p.deleted = 0
- AND (p.latest_summary_text IS NULL OR p.latest_summary_time IS NULL)`
- const [result] = await conn.query(sql)
- console.log(` summary backfill affectedRows=${result.affectedRows || 0}`)
- }
- async function backfillStatsCache(conn) {
- if (!(await tableExists(conn, 'lobster_user_profile'))) return
- const companySql = `
- INSERT INTO lobster_profile_stats
- (company_id, total_profiles, enriched_count, interaction_sum, avg_conversations, active_today, refresh_time)
- SELECT
- p.company_id,
- COUNT(1),
- IFNULL(SUM(CASE WHEN p.value_score > 10 OR p.total_purchase > 0
- OR (p.internal_tags IS NOT NULL AND p.internal_tags != '') THEN 1 ELSE 0 END), 0),
- IFNULL(SUM(p.interaction_count), 0),
- IFNULL(ROUND(AVG(p.interaction_count), 2), 0),
- IFNULL(SUM(CASE WHEN DATE(p.last_active_time) = CURDATE() THEN 1 ELSE 0 END), 0),
- NOW()
- FROM lobster_user_profile p
- WHERE p.deleted = 0
- GROUP BY p.company_id
- ON DUPLICATE KEY UPDATE
- total_profiles = VALUES(total_profiles),
- enriched_count = VALUES(enriched_count),
- interaction_sum = VALUES(interaction_sum),
- avg_conversations = VALUES(avg_conversations),
- active_today = VALUES(active_today),
- refresh_time = NOW()`
- const tenantSql = `
- INSERT INTO lobster_profile_stats
- (company_id, total_profiles, enriched_count, interaction_sum, avg_conversations, active_today, refresh_time)
- SELECT
- 0,
- COUNT(1),
- IFNULL(SUM(CASE WHEN p.value_score > 10 OR p.total_purchase > 0
- OR (p.internal_tags IS NOT NULL AND p.internal_tags != '') THEN 1 ELSE 0 END), 0),
- IFNULL(SUM(p.interaction_count), 0),
- IFNULL(ROUND(AVG(p.interaction_count), 2), 0),
- IFNULL(SUM(CASE WHEN DATE(p.last_active_time) = CURDATE() THEN 1 ELSE 0 END), 0),
- NOW()
- FROM lobster_user_profile p
- WHERE p.deleted = 0
- ON DUPLICATE KEY UPDATE
- total_profiles = VALUES(total_profiles),
- enriched_count = VALUES(enriched_count),
- interaction_sum = VALUES(interaction_sum),
- avg_conversations = VALUES(avg_conversations),
- active_today = VALUES(active_today),
- refresh_time = NOW()`
- await conn.query(companySql)
- await conn.query(tenantSql)
- console.log(' stats cache backfill done')
- }
- async function ensureUserSegmentTable(conn) {
- if (!(await tableExists(conn, 'lobster_user_segment'))) {
- await conn.query(SEGMENT_DDL)
- console.log(' created lobster_user_segment')
- }
- }
- async function applySegmentSeed(conn) {
- await ensureUserSegmentTable(conn)
- await conn.query(SEGMENT_SQL)
- console.log(' segment seed applied')
- }
- async function removeSummaryMenu(conn) {
- const statements = MENU_SQL.split(';').map(s => s.trim()).filter(Boolean)
- for (const sql of statements) {
- await conn.query(sql)
- }
- console.log(' conversation-summary menu removed')
- }
- async function verifyTenant(conn) {
- return {
- lobster_profile_stats: await tableExists(conn, 'lobster_profile_stats'),
- latest_stage: await columnExists(conn, 'lobster_user_profile', 'latest_stage'),
- latest_summary_text: await columnExists(conn, 'lobster_user_profile', 'latest_summary_text'),
- interaction_sum: await columnExists(conn, 'lobster_profile_stats', 'interaction_sum'),
- builtinSegmentRows: (await tableExists(conn, 'lobster_user_segment'))
- ? (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
- : 0,
- statsCacheRows: (await tableExists(conn, 'lobster_profile_stats'))
- ? (await conn.query('SELECT COUNT(*) AS cnt FROM lobster_profile_stats'))[0][0].cnt
- : 0
- }
- }
- async function applyTenant(tenant) {
- const conn = await mysql.createConnection({
- host: tenant.host,
- port: tenant.port,
- user: tenant.user,
- password: tenant.password,
- database: tenant.database,
- multipleStatements: true,
- charset: 'utf8mb4'
- })
- try {
- console.log(`\n[RUN] tenantId=${tenant.id} db=${tenant.database}`)
- await ensureProfileSchema(conn)
- await ensureSummaryColumns(conn)
- await backfillSummaryProjection(conn)
- await backfillStatsCache(conn)
- await applySegmentSeed(conn)
- await removeSummaryMenu(conn)
- const checks = await verifyTenant(conn)
- console.log(` verify: ${JSON.stringify(checks)}`)
- } finally {
- await conn.end()
- }
- }
- async function main() {
- const filterIds = parseArgs()
- const masterConn = await mysql.createConnection({ ...MASTER, charset: 'utf8mb4' })
- const tenants = []
- try {
- const [rows] = await masterConn.query(
- 'SELECT id, db_name, db_url, db_account, db_pwd, status FROM tenant_info ORDER BY id'
- )
- for (const row of rows) {
- if (row.status !== 1) continue
- if (filterIds.length > 0 && !filterIds.includes(Number(row.id))) continue
- tenants.push({
- id: row.id,
- status: row.status,
- ...parseJdbcUrl(row.db_url, row.db_name, row.db_account, row.db_pwd)
- })
- }
- } finally {
- await masterConn.end()
- }
- console.log(`Active tenants to migrate: ${tenants.length}`)
- let ok = 0
- let fail = 0
- for (const tenant of tenants) {
- try {
- await applyTenant(tenant)
- ok++
- } catch (err) {
- fail++
- console.error(`[FAIL] tenantId=${tenant.id} db=${tenant.database}: ${err.message}`)
- }
- }
- console.log(`\nDone. ok=${ok} fail=${fail}`)
- if (fail > 0) process.exit(1)
- }
- main().catch(err => {
- console.error(err)
- process.exit(1)
- })
|