Files
wzdj/scripts/import-lytiao-from-source-mongo.ts

475 lines
20 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

/**
* 从 KR/腾讯云 源库【直接复制】用户到玩值电竞库(非绑定,整条写入)
* - 以手机号为唯一键,从多表合并后复制到玩值 users在后台用户列表直接登记
* - 补足用户数据用户估值、详细材料、主播与公会关联guildMembers + streamerBindings
*
* 环境变量(.env.local
* - MONGODB_URI玩值电竞库写入目标
* - SOURCE_MONGODB_URI源库腾讯云/KR必填才执行复制
* - SOURCE_DB_NAME源库库名KR 可设为 KR
* - SOURCE_COLLECTIONS_KR逗号分隔如 用户估值,用户资产整合,存客宝用户资产,老坑爹论坛 www.lkdie.com_已扫描,老坑爹商店 shop.lkdie.com_已扫描
*
* 执行pnpm run import:lytiao:mongo
*/
import { config } from "dotenv"
import { resolve } from "path"
config({ path: resolve(process.cwd(), ".env.local") })
import { MongoClient, ObjectId } from "mongodb"
import { getDb, closeMongo } from "../lib/mongodb/client"
import { COLLECTIONS } from "../lib/db/mongo/collections"
import { LYTIAO_STREAMERS, LYTIAO_SOURCE } from "../lib/data/lytiao-streamers"
const now = () => new Date().toISOString()
const SOURCE_DB_NAME = process.env.SOURCE_DB_NAME || "asset_db"
const SOURCE_COLLECTION_USERS = process.env.SOURCE_COLLECTION_USERS || ""
const SOURCE_COLLECTION_ORDERS = process.env.SOURCE_COLLECTION_ORDERS || ""
const SOURCE_COLLECTIONS_KR = (process.env.SOURCE_COLLECTIONS_KR || "").split(",").map((s) => s.trim()).filter(Boolean)
const GUILD_NAME_WANZHI = "玩值电竞"
/** 确保存在公会「玩值电竞」并将所有主播加入该公会 */
async function ensureWanzhiGuildAndAssignStreamers() {
const db = await getDb()
const guildsColl = db.collection(COLLECTIONS.guilds)
const streamersColl = db.collection(COLLECTIONS.streamers)
const usersColl = db.collection(COLLECTIONS.users)
const firstUser = await usersColl.findOne({})
const ownerId = firstUser ? String((firstUser as { _id: ObjectId })._id) : ""
let guild = await guildsColl.findOne({ name: GUILD_NAME_WANZHI })
if (!guild) {
const res = await guildsColl.insertOne({
name: GUILD_NAME_WANZHI,
logo: "",
coverImage: "",
description: "玩值电竞官方公会",
ownerId,
memberCount: 0,
maxMembers: 5000,
level: 10,
totalIncome: 0,
commissionRate: 0.1,
tags: ["官方", "电竞"],
games: ["魔兽世界", "英雄联盟", "王者荣耀"],
requirements: "",
benefits: ["流量扶持", "课程分成"],
status: "active",
createdAt: now(),
updatedAt: now(),
})
guild = await guildsColl.findOne({ _id: res.insertedId })
console.log("[import-lytiao] 公会「玩值电竞」已创建")
}
const guildId = (guild as { _id: ObjectId })._id.toString()
const r = await streamersColl.updateMany({}, { $set: { guildId, updatedAt: now() } })
if (r.modifiedCount > 0) console.log("[import-lytiao] 主播已归属公会「玩值电竞」:", r.modifiedCount, "条")
}
async function ensureLytiaoStreamersInTarget() {
const db = await getDb()
const streamersColl = db.collection(COLLECTIONS.streamers)
const usersColl = db.collection(COLLECTIONS.users)
let defaultUserId: ObjectId
const firstUser = await usersColl.findOne({})
if (firstUser && (firstUser as { _id?: ObjectId })._id) {
defaultUserId = (firstUser as { _id: ObjectId })._id
} else {
const ins = await usersColl.insertOne({
phone: "13800138000",
name: "系统",
registerSource: "system",
status: "active",
createdAt: now(),
updatedAt: now(),
})
defaultUserId = ins.insertedId
}
let inserted = 0
for (const s of LYTIAO_STREAMERS) {
const exists = await streamersColl.findOne({ source: LYTIAO_SOURCE, name: s.name })
if (exists) continue
const avatar = s.avatar || `/streamers/lytiao/${s.name.replace(/[\u3000]/g, "").trim()}.jpg`
await streamersColl.insertOne({
userId: defaultUserId.toString(),
name: s.name,
title: s.title,
intro: s.intro,
game: s.game,
gameId: s.gameId,
tags: s.tags,
avatar,
coverImage: "",
level: 20,
fans: Math.max(0, s.douyuFans ?? s.douyinFans ?? 0),
hotValue: 0,
isLive: false,
totalGifts: 0,
totalIncome: 0,
commissionRate: 0.7,
status: "approved",
verifiedAt: now(),
source: LYTIAO_SOURCE,
city: s.city ?? "深圳",
createdAt: now(),
updatedAt: now(),
})
inserted++
}
console.log("[import-lytiao] 老游条主播 已存在或本次新增:", inserted)
}
/** 从单集合中取手机号(支持多种字段名) */
function normalizePhone(raw: Record<string, unknown>): string {
const p = String(raw.phone ?? raw.mobile ?? raw.tel ?? raw. ?? raw.Phone ?? raw. ?? "").trim()
return p.replace(/\s/g, "").replace(/^\+86/, "") || ""
}
/** 合并两条记录:优先保留非空、字符串取更长者 */
function mergeRecord(base: Record<string, unknown>, from: Record<string, unknown>): void {
for (const [k, v] of Object.entries(from)) {
if (k === "_id" || v === undefined || v === null) continue
const cur = base[k]
if (cur === undefined || cur === null) {
base[k] = v
continue
}
if (typeof v === "string" && typeof cur === "string" && v.length > (cur as string).length) base[k] = v
else if (typeof v === "number" && (typeof cur !== "number" || cur === 0)) base[k] = v
else if (typeof v === "object" && !Array.isArray(v) && typeof cur === "object" && !Array.isArray(cur))
mergeRecord(cur as Record<string, unknown>, v as Record<string, unknown>)
}
}
const BATCH_SIZE = 6000
/** 全量流式:逐集合、逐批读取源库,与目标库已有记录合并后 upsert控制内存 */
async function importUsersFromSourceStreaming(
sourceDb: { collection: (n: string) => { find: (q: object) => AsyncIterable<Record<string, unknown>> } },
targetUsersColl: { find: (q: object) => Promise<Record<string, unknown>[]>; bulkWrite: (ops: object[]) => Promise<{ upsertedCount: number; modifiedCount: number }> },
collectionNames: string[]
) {
let totalUpserted = 0
let totalModified = 0
for (const collName of collectionNames) {
if (!collName.trim()) continue
const sourceColl = sourceDb.collection(collName.trim()) as unknown as { find: (q: object) => AsyncIterable<Record<string, unknown>> }
const cursor = sourceColl.find({})
let batch: Record<string, unknown>[] = []
let count = 0
try {
for await (const doc of cursor) {
const raw = doc as Record<string, unknown>
const phone = normalizePhone(raw)
if (!phone || phone.length < 6) continue
batch.push({ ...raw, phone })
count++
if (batch.length >= BATCH_SIZE) {
const byPhoneInBatch = new Map<string, Record<string, unknown>>()
for (const raw of batch) {
const phone = String(raw.phone).trim()
const existing = byPhoneInBatch.get(phone)
if (existing) mergeRecord(existing, raw)
else byPhoneInBatch.set(phone, { ...raw, sourceCollections: [collName] })
}
const phones = [...byPhoneInBatch.keys()]
const existingList = (await (targetUsersColl as { find: (q: object) => { toArray: () => Promise<Record<string, unknown>[]> } }).find({ phone: { $in: phones } }).toArray()) as { phone: string; [k: string]: unknown }[]
const byPhone = new Map<string, Record<string, unknown>>()
for (const e of existingList) byPhone.set(String(e.phone).trim(), e)
const ops: { updateOne: { filter: { phone: string }; update: { $set: Record<string, unknown> }; upsert: boolean } }[] = []
for (const [phone, merged] of byPhoneInBatch) {
const base = byPhone.get(phone)
const mergedFull = base ? { ...base, ...merged } : merged
if (!Array.isArray(mergedFull.sourceCollections)) mergedFull.sourceCollections = [collName]
else if (!mergedFull.sourceCollections.includes(collName)) mergedFull.sourceCollections.push(collName)
const doc = buildWanzhiUser(mergedFull)
ops.push({
updateOne: {
filter: { phone },
update: { $set: doc },
upsert: true,
},
})
}
const result = await targetUsersColl.bulkWrite(ops as never)
totalUpserted += result.upsertedCount ?? 0
totalModified += result.modifiedCount ?? 0
batch = []
}
}
if (batch.length > 0) {
const byPhoneInBatch = new Map<string, Record<string, unknown>>()
for (const raw of batch) {
const phone = String(raw.phone).trim()
const existing = byPhoneInBatch.get(phone)
if (existing) mergeRecord(existing, raw)
else byPhoneInBatch.set(phone, { ...raw, sourceCollections: [collName] })
}
const phones = [...byPhoneInBatch.keys()]
const existingList = (await (targetUsersColl as { find: (q: object) => { toArray: () => Promise<Record<string, unknown>[]> } }).find({ phone: { $in: phones } }).toArray()) as { phone: string; [k: string]: unknown }[]
const byPhone = new Map<string, Record<string, unknown>>()
for (const e of existingList) byPhone.set(String(e.phone).trim(), e)
const ops: { updateOne: { filter: { phone: string }; update: { $set: Record<string, unknown> }; upsert: boolean } }[] = []
for (const [phone, merged] of byPhoneInBatch) {
const base = byPhone.get(phone)
const mergedFull = base ? { ...base, ...merged } : merged
if (!Array.isArray(mergedFull.sourceCollections)) mergedFull.sourceCollections = [collName]
else if (!mergedFull.sourceCollections.includes(collName)) mergedFull.sourceCollections.push(collName)
const doc = buildWanzhiUser(mergedFull)
ops.push({ updateOne: { filter: { phone }, update: { $set: doc }, upsert: true } })
}
const result = await targetUsersColl.bulkWrite(ops as never)
totalUpserted += result.upsertedCount ?? 0
totalModified += result.modifiedCount ?? 0
}
console.log("[import-lytiao] 集合", collName, "全量处理条数:", count)
} catch (e) {
console.warn("[import-lytiao] 集合", collName, "失败:", (e as Error).message)
}
}
return { totalUpserted, totalModified }
}
/** 从源记录中提取分类(多字段兼容),用于后台用户管理展示 */
function extractCategory(merged: Record<string, unknown>): string[] {
const out: string[] = []
const push = (v: unknown) => {
if (typeof v === "string" && v.trim()) out.push(v.trim())
else if (Array.isArray(v)) v.forEach((x) => (typeof x === "string" && x.trim() ? out.push(x.trim()) : null))
}
push(merged.)
push(merged.)
push(merged.)
push(merged.)
push(merged.category)
push(merged.classification)
push(merged.)
push(merged.tags)
push(merged.)
push(merged.source)
const seen = new Set<string>()
return out.filter((s) => !seen.has(s) && seen.add(s))
}
/** 从合并记录中取最佳姓名(优先真实姓名、姓名) */
function pickName(merged: Record<string, unknown>): string {
const s = String(
merged. ?? merged. ?? merged.name ?? merged. ?? merged.nickname ?? merged. ?? merged.username ?? ""
).trim()
return s || "导入用户"
}
/** 从合并记录中取最佳昵称(优先昵称、用户名) */
function pickNickname(merged: Record<string, unknown>): string {
const s = String(
merged. ?? merged.nickname ?? merged.username ?? merged. ?? merged.name ?? merged. ?? ""
).trim()
return s || pickName(merged)
}
/** 将合并后的用户直接复制到玩值 users以手机号登记补全姓名/昵称/估值/材料等) */
function buildWanzhiUser(merged: Record<string, unknown>): Record<string, unknown> {
const phone = String(merged.phone ?? "").trim()
const name = pickName(merged)
const nickname = pickNickname(merged)
const registerTime = (merged.createdAt ?? merged.registerTime ?? merged. ?? merged. ?? now()) as string
const valuation =
typeof merged.valuation === "number"
? merged.valuation
: typeof merged. === "number"
? merged.估值
: typeof merged. === "number"
? merged.用户估值
: undefined
const materials = merged.materials ?? merged. ?? merged. ?? merged. ?? merged.content ?? merged.
const odlId = merged._id ? String((merged._id as ObjectId).toString()) : undefined
const category = extractCategory(merged)
const doc: Record<string, unknown> = {
phone,
name,
nickname: nickname || name,
registerSource: merged.registerSource ?? "lytiao",
registerTime: typeof registerTime === "string" ? registerTime : now(),
status: "active",
createdAt: merged.createdAt ?? now(),
updatedAt: now(),
}
if (odlId) doc.odlId = odlId
if (valuation !== undefined) doc.valuation = valuation
if (materials !== undefined && materials !== null) doc.materials = materials
if (Array.isArray(merged.sourceCollections)) doc.sourceCollections = merged.sourceCollections
if (category.length > 0) doc.category = category
if (merged.avatar !== undefined && merged.avatar !== null) doc.avatar = merged.avatar
if (merged.level !== undefined && typeof merged.level === "number") doc.level = merged.level
if (merged.balance !== undefined && typeof merged.balance === "number") doc.balance = merged.balance
if (merged.tags !== undefined && Array.isArray(merged.tags)) doc.tags = merged.tags
if (merged.gender !== undefined && merged.gender !== null) doc.gender = merged.gender
if (merged.address !== undefined && merged.address !== null) doc.address = merged.address
if (merged.city !== undefined && merged.city !== null) doc.city = merged.city
return doc
}
/** 从源库多表全量流式复制到玩值 users补全姓名/昵称,以手机号登记 */
async function importUsersFromSource() {
const uri = process.env.SOURCE_MONGODB_URI
if (!uri) {
console.log("[import-lytiao] 未设置 SOURCE_MONGODB_URI跳过用户复制")
return
}
const collectionNames =
SOURCE_COLLECTIONS_KR.length > 0 ? SOURCE_COLLECTIONS_KR : SOURCE_COLLECTION_USERS ? [SOURCE_COLLECTION_USERS] : []
if (collectionNames.length === 0) {
console.log("[import-lytiao] 未配置 SOURCE_COLLECTIONS_KR 或 SOURCE_COLLECTION_USERS跳过用户复制")
return
}
const client = new MongoClient(uri, { serverSelectionTimeoutMS: 15000 })
await client.connect()
const sourceDb = client.db(SOURCE_DB_NAME)
const db = await getDb()
const usersColl = db.collection(COLLECTIONS.users)
const { totalUpserted, totalModified } = await importUsersFromSourceStreaming(sourceDb, usersColl as never, collectionNames)
await client.close()
if (totalUpserted > 0 || totalModified > 0) {
console.log("[import-lytiao] 全量复制到玩值电竞 新增:", totalUpserted, "更新:", totalModified, "(姓名/昵称已补全,以手机号登记)")
}
}
/** 为已导入用户补足公会成员与主播绑定(玩值电竞公会 + 老游条主播) */
async function ensureImportedUsersGuildAndStreamerBindings() {
const db = await getDb()
const guildsColl = db.collection(COLLECTIONS.guilds)
const guildMembersColl = db.collection(COLLECTIONS.guildMembers)
const streamersColl = db.collection(COLLECTIONS.streamers)
const bindingsColl = db.collection(COLLECTIONS.streamerBindings)
const usersColl = db.collection(COLLECTIONS.users)
const guild = await guildsColl.findOne({ name: GUILD_NAME_WANZHI })
if (!guild) return
const guildId = (guild as { _id: ObjectId })._id.toString()
const streamers = await streamersColl.find({}).limit(50).toArray()
if (streamers.length === 0) return
const users = await usersColl.find({ phone: { $exists: true, $ne: "" } }).toArray()
const nowStr = now()
let guildAdded = 0
let bindAdded = 0
for (const u of users) {
const uid = (u as { _id: ObjectId })._id.toString()
const existingMember = await guildMembersColl.findOne({ guildId, userId: uid })
if (!existingMember) {
await guildMembersColl.insertOne({
guildId,
userId: uid,
role: "member",
status: "active",
createdAt: nowStr,
updatedAt: nowStr,
})
guildAdded++
}
const firstStreamer = streamers[0] as { _id: ObjectId }
const existingBinding = await bindingsColl.findOne({ userId: uid, streamerId: firstStreamer._id.toString() })
if (!existingBinding) {
await bindingsColl.insertOne({
userId: uid,
streamerId: firstStreamer._id.toString(),
status: "active",
createdAt: nowStr,
updatedAt: nowStr,
})
bindAdded++
}
}
if (guildAdded > 0 || bindAdded > 0) {
console.log("[import-lytiao] 用户补足公会/主播数据 公会成员:", guildAdded, "主播绑定:", bindAdded)
}
}
/** 用玩值现有用户与老游条主播建立绑定 + 示例课程购买与分成流水 */
async function linkUsersToStreamersAndCreateDemoIncome() {
const db = await getDb()
const usersColl = db.collection(COLLECTIONS.users)
const streamersColl = db.collection(COLLECTIONS.streamers)
const bindingsColl = db.collection(COLLECTIONS.streamerBindings)
const coursePurchasesColl = db.collection(COLLECTIONS.coursePurchases)
const transactionsColl = db.collection(COLLECTIONS.transactions)
const streamers = await streamersColl.find({ source: LYTIAO_SOURCE }).toArray()
const users = await usersColl.find({}).limit(200).toArray()
if (streamers.length === 0 || users.length === 0) {
console.log("[import-lytiao] 无老游条主播或用户,跳过绑定与收益示例")
return
}
let bindingsAdded = 0
let purchasesAdded = 0
let txAdded = 0
const nowStr = now()
for (let i = 0; i < users.length; i++) {
const u = users[i] as { _id: ObjectId }
const streamer = streamers[i % streamers.length] as { _id: ObjectId; name: string }
const existingBinding = await bindingsColl.findOne({ userId: u._id.toString(), streamerId: streamer._id.toString() })
if (!existingBinding) {
await bindingsColl.insertOne({
userId: u._id.toString(),
streamerId: streamer._id.toString(),
status: "active",
createdAt: nowStr,
updatedAt: nowStr,
})
bindingsAdded++
}
if (i < 30 && Math.random() > 0.5) {
const amount = [259, 398, 998][Math.floor(Math.random() * 3)]
const orderNo = `LYT-${Date.now()}-${i}`
await coursePurchasesColl.insertOne({
userId: u._id.toString(),
streamerId: streamer._id.toString(),
courseId: `course-${streamer.name}-1`,
courseName: `${streamer.name} 精品课`,
amount,
orderNo,
status: "paid",
createdAt: nowStr,
updatedAt: nowStr,
})
purchasesAdded++
const commissionRate = 0.7
const streamerIncome = Math.round(amount * commissionRate)
await transactionsColl.insertOne({
userId: u._id.toString(),
orderNo,
type: "course_income",
amount: streamerIncome,
remark: `课程分成-${streamer.name}`,
createdAt: nowStr,
})
txAdded++
}
}
console.log("[import-lytiao] 主播-用户绑定 新增:", bindingsAdded, "课程购买示例:", purchasesAdded, "分成流水:", txAdded)
}
async function main() {
try {
await ensureWanzhiGuildAndAssignStreamers()
await ensureLytiaoStreamersInTarget()
await importUsersFromSource()
await ensureImportedUsersGuildAndStreamerBindings()
await linkUsersToStreamersAndCreateDemoIncome()
} finally {
await closeMongo()
}
}
main().catch((e) => {
console.error(e)
process.exit(1)
})