#!/usr/bin/env python3 """模拟销售设备采集代理 — 生成模拟微信数据并同步到中央 PostgreSQL。 用法: python mock_sync.py --salesperson-id 1 python mock_sync.py --salesperson-id 1 --dsn "host=127.0.0.1 port=5434 dbname=wxchat_sales" """ from __future__ import annotations import argparse import datetime as dt import hashlib import json import pathlib import random import sys import psycopg2 from psycopg2.extras import execute_values TZ = dt.timezone(dt.timedelta(hours=8)) # --------------------------------------------------------------------------- # 数据库连接 # --------------------------------------------------------------------------- DEFAULT_SOCKET = str(pathlib.Path(__file__).resolve().parent.parent / "db" / "socket") DEFAULT_DSN = f"host={DEFAULT_SOCKET} port=5434 dbname=wxchat_sales" # --------------------------------------------------------------------------- # 销售人员信息 # --------------------------------------------------------------------------- SALESPERSONS = { 1: {"name": "张伟", "wx_account": "wxid_zhangwei", "device_id": "macbook-zw-001"}, 2: {"name": "李娜", "wx_account": "wxid_lina", "device_id": "macbook-ln-002"}, 3: {"name": "王强", "wx_account": "wxid_wangqiang", "device_id": "macbook-wq-003"}, } # --------------------------------------------------------------------------- # 场景模板 — 真实客户对话 # --------------------------------------------------------------------------- SCENARIO_TEMPLATES = [ { "name": "湿疹宝宝咨询后成交", "stage": "成交", "contact": { "nickname": "辰辰妈妈", "remark": "辰辰妈-湿疹-2岁", "industry": "宝妈", "intent_level": "high", }, "messages": [ ("customer", "text", "你好,我家宝宝2岁,湿疹反复好几个月了,朋友推荐你这边益生菌"), ("sales", "text", "辰辰妈您好!宝宝湿疹确实让人心疼,请问现在湿疹主要在哪些部位?有用过什么药吗?"), ("customer", "text", "脸上和手臂都有,医生开了激素药膏,但停了就复发"), ("sales", "text", "理解,激素药膏只能暂时压制。益生菌是从肠道调节免疫,从根本上降低过敏反应。我们用的是丹麦进口的鼠李糖乳杆菌,有专门针对儿童湿疹的临床验证"), ("sales", "image", "[图片:临床验证报告截图]"), ("customer", "text", "这个是进口的?安全吗?2岁能吃吗?"), ("sales", "text", "是的,丹麦进口菌株,0岁以上就能用,无敏配方,不含牛奶蛋白和麸质。很多宝妈反馈坚持吃2-3个月湿疹明显好转"), ("customer", "text", "多少钱?怎么卖的?"), ("sales", "text", "单盒298元30袋,建议先吃3盒一个周期,3盒套餐798元算下来每天不到9块钱"), ("sales", "image", "[图片:产品包装图]"), ("customer", "text", "3盒798是吧,效果不好怎么办?"), ("sales", "text", "我们有售后指导,期间有任何问题随时找我。另外我发您几个同情况宝妈的反馈看看"), ("sales", "link", "[链接:宝妈真实反馈合集]"), ("customer", "text", "好的,那先来3盒试试"), ("sales", "text", "好的辰辰妈!3盒套餐798元,您方便现在付款吗?我这边给您安排发货"), ("customer", "text", "已经转过去了,麻烦尽快发货"), ("sales", "text", "收到!今天就给您发顺丰,预计后天到"), ("customer", "emoji", "[表情:谢谢]"), ], }, { "name": "过敏性鼻炎咨询后犹豫", "stage": "异议处理", "contact": { "nickname": "乐乐妈", "remark": "乐乐妈-鼻炎-5岁", "industry": "宝妈", "intent_level": "medium", }, "messages": [ ("customer", "text", "你好,听说益生菌对过敏性鼻炎有帮助?"), ("sales", "text", "乐乐妈您好!是的,益生菌可以调节肠道免疫,对过敏性鼻炎有改善作用。请问宝宝现在鼻炎什么症状?"), ("customer", "text", "早上起来一直打喷嚏流鼻涕,医生说是过敏性鼻炎"), ("sales", "text", "益生菌从肠道调节免疫系统,坚持吃可以逐渐降低过敏反应。我们的菌株是丹麦进口的,有临床验证"), ("customer", "text", "有没有副作用?安全吗?"), ("sales", "text", "很安全的,无敏配方,不含牛奶蛋白和麸质,5岁完全没问题"), ("customer", "text", "多少钱?"), ("sales", "text", "单盒298,3盒套餐798"), ("customer", "text", "有点贵啊..."), ("sales", "text", "算下来每天不到9块钱,比吃药划算多了"), ("customer", "text", "我再想想,和老公商量一下"), ("sales", "text", "好的,不着急。您可以先看看这个科普文章"), ("sales", "link", "[链接:益生菌与儿童过敏科普]"), ("customer", "text", "好的,我看看"), ], }, { "name": "朋友推荐来咨询", "stage": "需求发现", "contact": { "nickname": "果果妈妈", "remark": "果果妈-推荐来的", "industry": "宝妈", "intent_level": "medium", }, "messages": [ ("customer", "text", "你好,是小林推荐我加你的,她说你家益生菌不错"), ("sales", "text", "果果妈您好!谢谢小林推荐,请问宝宝多大了?有什么情况吗?"), ("customer", "text", "1岁半,最近总是拉肚子,小林说吃益生菌调理一下"), ("sales", "text", "拉肚子多久了?有去看过医生吗?"), ("customer", "text", "断断续续一个多月了,医生说没什么大问题,让注意饮食"), ("sales", "text", "那益生菌确实可以帮到宝宝。我们的菌株是丹麦进口的,专门针对儿童肠道和免疫调节。1岁半完全可以吃"), ("customer", "text", "成分安全吗?"), ("sales", "text", "无敏配方,不含牛奶蛋白、麸质和大豆,很多敏宝妈妈都在用"), ("customer", "text", "好的,我先了解一下"), ("sales", "text", "没问题,我给您发个产品详情,您先看看"), ("sales", "link", "[链接:产品详情页]"), ], }, { "name": "长期跟进后复购", "stage": "成交", "contact": { "nickname": "糖糖妈妈", "remark": "糖糖妈-复购-湿疹已好转", "industry": "宝妈", "intent_level": "high", }, "messages": [ ("customer", "text", "你好,我想了解一下益生菌"), ("sales", "text", "糖糖妈您好!请问宝宝什么情况?"), ("customer", "text", "3岁,湿疹,反反复复的"), ("sales", "text", "湿疹确实需要从内调理。我们的益生菌丹麦进口菌株,专门针对儿童过敏"), ("customer", "text", "多少钱?"), ("sales", "text", "单盒298,3盒套餐798"), ("customer", "text", "先了解一下,回头再说"), ("sales", "text", "好的,不着急"), ("sales", "text", "糖糖妈,宝宝最近湿疹怎么样了?"), ("customer", "text", "还是老样子,断不了根"), ("sales", "text", "湿疹确实需要坚持调理,益生菌一般吃2-3个月能看到明显改善"), ("customer", "text", "我再考虑考虑"), ("sales", "text", "糖糖妈,今天有个宝妈反馈说吃了2个月湿疹好多了,发您看看"), ("sales", "image", "[图片:宝妈反馈截图]"), ("customer", "text", "真的吗?那我也试试吧"), ("sales", "text", "好的!建议先拿3盒一个周期,798元"), ("customer", "text", "好的,付款吧"), ("sales", "text", "好的!我马上安排发货"), ("sales", "text", "糖糖妈,宝宝吃了这段时间效果怎么样?"), ("customer", "text", "比之前好多了,湿疹没那么频繁了"), ("sales", "text", "太好了!建议继续吃巩固一下,现在有6盒套餐1499更划算"), ("customer", "text", "好,那来6盒"), ("sales", "text", "好的!马上安排"), ], }, { "name": "对比竞品后选择", "stage": "方案匹配", "contact": { "nickname": "壮壮妈妈", "remark": "壮壮妈-对比合生元", "industry": "宝妈", "intent_level": "medium", }, "messages": [ ("customer", "text", "你好,我想给孩子买益生菌,但在对比你们和合生元"), ("sales", "text", "壮壮妈您好!对比是正常的。请问宝宝主要什么情况?"), ("customer", "text", "经常感冒,免疫力不太好"), ("sales", "text", "我们的菌株是丹麦进口的鼠李糖乳杆菌GG,有专门针对儿童免疫力的临床验证。合生元用的是国产菌株"), ("customer", "text", "但你们比合生元贵啊"), ("sales", "text", "菌株不一样,进口菌株的临床数据更充分,效果更有保障。算下来每天不到9块钱"), ("customer", "text", "有什么区别吗?"), ("sales", "image", "[图片:菌株对比表]"), ("sales", "text", "主要是菌株来源和临床验证数量不同。我们的LGG是全球研究最多的益生菌菌株之一"), ("customer", "text", "我再看看吧"), ("sales", "text", "好的,有问题随时问我"), ], }, ] # --------------------------------------------------------------------------- # 场景变体 — 给不同销售使用 # --------------------------------------------------------------------------- SCENARIO_VARIANTS = [ { "name": "食物过敏咨询后成交", "stage": "成交", "contact": { "nickname": "依依妈", "remark": "依依妈-食物过敏-3岁", "industry": "宝妈", "intent_level": "high", }, "messages": [ ("customer", "text", "你好,我家宝宝对牛奶蛋白过敏,听说益生菌能帮助脱敏?"), ("sales", "text", "依依妈您好!是的,益生菌可以帮助调节肠道免疫,对食物过敏有辅助改善作用。请问宝宝现在多大了?"), ("customer", "text", "3岁,查出来对牛奶蛋白和鸡蛋过敏"), ("sales", "text", "理解,食物过敏的宝宝确实需要从肠道调理。我们的菌株是丹麦进口LGG,有临床数据支持对食物过敏的改善"), ("customer", "text", "无敏配方吗?我家对牛奶蛋白过敏"), ("sales", "text", "是的,完全无敏配方,不含牛奶蛋白、麸质和大豆,过敏宝宝可以放心吃"), ("sales", "image", "[图片:成分表]"), ("customer", "text", "看着成分确实干净,多少钱?"), ("sales", "text", "单盒298,3盒套餐798,建议先吃一个周期"), ("customer", "text", "好,来3盒吧"), ("sales", "text", "好的依依妈!3盒798元,给您安排发货"), ("customer", "text", "付款了,麻烦尽快发"), ("sales", "text", "收到!今天就发顺丰"), ], }, { "name": "腹泻咨询后成交", "stage": "成交", "contact": { "nickname": "阳阳妈", "remark": "阳阳妈-腹泻-1岁", "industry": "宝妈", "intent_level": "high", }, "messages": [ ("customer", "text", "你好,宝宝1岁,最近老拉肚子,朋友说益生菌有用"), ("sales", "text", "阳阳妈您好!益生菌对宝宝腹泻确实有帮助。请问拉肚子多久了?"), ("customer", "text", "一个多星期了,吃了药好点,停了又拉"), ("sales", "text", "这种情况益生菌可以帮到宝宝调节肠道菌群。我们的丹麦进口菌株专门针对儿童肠道健康"), ("customer", "text", "1岁能吃吗?"), ("sales", "text", "可以的,0岁以上就能用,无敏配方很安全"), ("customer", "text", "多少钱?"), ("sales", "text", "单盒298,3盒套餐798"), ("customer", "text", "先来3盒试试"), ("sales", "text", "好的!3盒798元,马上安排发货"), ("customer", "text", "好的,已付款"), ], }, { "name": "价格异议后放弃", "stage": "异议处理", "contact": { "nickname": "可可妈妈", "remark": "可可妈-价格异议", "industry": "宝妈", "intent_level": "low", }, "messages": [ ("customer", "text", "益生菌多少钱?"), ("sales", "text", "可可妈您好!单盒298,3盒套餐798"), ("customer", "text", "这么贵?网上益生菌才几十块"), ("sales", "text", "我们的是丹麦进口菌株,有临床验证,和网上几十块的不一样"), ("customer", "text", "但是也太贵了"), ("sales", "text", "算下来每天不到9块钱,效果有保障"), ("customer", "text", "我再看看吧"), ("sales", "text", "好的,有需要随时找我"), ], }, { "name": "安全顾虑后犹豫", "stage": "异议处理", "contact": { "nickname": "贝贝妈妈", "remark": "贝贝妈-安全顾虑", "industry": "宝妈", "intent_level": "medium", }, "messages": [ ("customer", "text", "你好,宝宝6个月,能吃益生菌吗?"), ("sales", "text", "贝贝妈您好!6个月以上的宝宝可以吃的。请问宝宝什么情况?"), ("customer", "text", "老是胀气,睡眠也不好"), ("sales", "text", "益生菌可以帮宝宝调节肠道,改善胀气。我们的丹麦进口菌株0岁以上可用"), ("customer", "text", "这么小的宝宝吃安全吗?有没有副作用?"), ("sales", "text", "很安全的,无敏配方,不含牛奶蛋白和麸质"), ("customer", "text", "我还是有点担心,毕竟宝宝才6个月"), ("sales", "text", "理解您的顾虑,很多宝妈也有同样的担心。其实益生菌是很成熟的品类,我们的菌株有大量临床数据支持"), ("customer", "text", "我再想想吧"), ("sales", "text", "好的,不着急。有问题随时问我"), ], }, { "name": "过早报价后被拒", "stage": "异议处理", "contact": { "nickname": "欢欢妈妈", "remark": "欢欢妈-过早报价", "industry": "宝妈", "intent_level": "low", }, "messages": [ ("customer", "text", "益生菌多少钱?"), ("sales", "text", "您好!单盒298,3盒798,6盒1499"), ("customer", "text", "哦,这么贵"), ("sales", "text", "丹麦进口菌株,效果很好的"), ("customer", "text", "我先看看吧"), ("sales", "text", "好的,需要随时联系"), ("customer", "text", "嗯"), ], }, ] # --------------------------------------------------------------------------- # 非客户联系人模板 # --------------------------------------------------------------------------- NON_CUSTOMER_TEMPLATES = [ { "nickname": "韩国代购-小美", "remark": "", "is_group": False, "messages": [ ("other", "text", "【韩国直邮】儿童维生素团购开始啦,DHA+益生菌组合装只要199!"), ("other", "text", "需要的姐妹私我~"), ], }, { "nickname": "母婴用品批发", "remark": "", "is_group": False, "messages": [ ("other", "text", "姐,我们新款益生菌做活动,比你现在用的便宜一半,要不要看看?"), ("other", "text", "保证正品,支持验货"), ("other", "emoji", "[表情:微笑]"), ], }, { "nickname": "朵朵阿姨", "remark": "", "is_group": False, "messages": [ ("other", "text", "在吗?"), ("other", "text", "你朋友圈那个是什么产品?"), ("other", "text", "好的知道了"), ], }, { "nickname": "老婆 ❤️", "remark": "", "is_group": False, "messages": [ ("other", "text", "今晚回来吃饭吗?"), ("other", "text", "孩子放学了"), ("other", "text", "周末去哪玩?"), ("other", "text", "好的"), ], }, { "nickname": "快递小哥", "remark": "", "is_group": False, "messages": [ ("other", "text", "你的快递到了,放门口了"), ("other", "text", "签收一下"), ], }, { "nickname": "大学同学", "remark": "", "is_group": False, "messages": [ ("other", "text", "好久不见啊"), ("other", "text", "最近怎么样"), ("other", "text", "改天聚聚"), ], }, { "nickname": "同行-小张", "remark": "", "is_group": False, "messages": [ ("other", "text", "张哥,最近怎么样?"), ("other", "text", "你们那边转化率怎么样?"), ("other", "text", "交流一下经验"), ], }, { "nickname": "健身教练", "remark": "", "is_group": False, "messages": [ ("other", "text", "明天还来训练吗?"), ("other", "text", "记得带毛巾"), ], }, { "nickname": "房产中介小王", "remark": "", "is_group": False, "messages": [ ("other", "text", "哥,最近有考虑换房吗?"), ("other", "text", "新盘有优惠"), ], }, { "nickname": "表妹", "remark": "", "is_group": False, "messages": [ ("other", "text", "哥,你那边益生菌我朋友想了解一下"), ("other", "text", "我推她微信给你"), ("other", "text", "好的"), ], }, ] # --------------------------------------------------------------------------- # 群聊模板 # --------------------------------------------------------------------------- GROUP_TEMPLATES = [ { "nickname": "宝妈育儿交流群", "is_group": True, "messages": [ ("other", "text", "[群公告] 本群禁发广告,违者移出"), ("other", "text", "有没有宝妈推荐好的儿童面霜?"), ("other", "text", "我家用的那个丝塔芙还不错"), ("other", "text", "谢谢推荐"), ("other", "text", "有没有宝妈推荐好的儿童霜?"), ("other", "text", "大家宝宝都几岁上幼儿园的?"), ("other", "text", "我家3岁送的"), ("other", "text", "这么早吗?会不会太小了"), ("other", "text", "还好,适应期过了就好了"), ], }, { "nickname": "过敏宝宝互助群", "is_group": True, "messages": [ ("other", "text", "我家宝宝湿疹终于好了"), ("other", "text", "怎么好的?用的什么?"), ("other", "text", "吃了益生菌加上注意饮食"), ("other", "text", "你们都用的什么益生菌?"), ("other", "text", "益生菌真的有用吗?"), ("other", "text", "我觉得有用的,坚持吃了2个月"), ("other", "text", "我也想试试"), ("other", "text", "可以试试,从肠道调理确实有道理"), ("other", "text", "有没有副作用?"), ("other", "text", "我们吃的没什么副作用"), ("other", "text", "好的,谢谢分享"), ], }, { "nickname": "育儿干货分享群", "is_group": True, "messages": [ ("other", "text", "今天分享一篇关于儿童免疫力的文章"), ("other", "link", "[链接:儿童免疫力科普]"), ("other", "text", "写得好,收藏了"), ("other", "text", "谢谢分享"), ("other", "text", "有没有关于过敏的科普?"), ("other", "text", "有的,我找找"), ], }, ] # --------------------------------------------------------------------------- # 每个销售的联系人分配方案 # --------------------------------------------------------------------------- # 格式: (scenario_index_or_variant_name, contact_nickname, contact_remark) # scenario_index 指向 SCENARIO_TEMPLATES, variant_name 指向 SCENARIO_VARIANTS SALESPERSON_ASSIGNMENTS = { 1: [ # 张伟 — 资深 ("scenario_0", None, None), # 湿疹成交 ("scenario_2", None, None), # 朋友推荐 ("variant_食物过敏咨询后成交", None, None), ("variant_腹泻咨询后成交", None, None), ("scenario_3", None, None), # 长期跟进复购 # 潜在客户 ("scenario_2", "小宇妈", "小宇妈-咨询"), ("scenario_2", "安安妈妈", "安安妈-咨询"), # 非客户 ("non_customer", "韩国代购-小美", ""), ("non_customer", "母婴用品批发", ""), ("non_customer", "朵朵阿姨", ""), ("non_customer", "老婆 ❤️", ""), # 群聊 ("group", "宝妈育儿交流群", ""), ("group", "过敏宝宝互助群", ""), ], 2: [ # 李娜 — 新人 ("scenario_1", None, None), # 鼻炎犹豫 ("variant_价格异议后放弃", None, None), ("scenario_4", None, None), # 竞品对比 ("variant_过早报价后被拒", None, None), ("scenario_2", "甜甜妈", "甜甜妈-推荐咨询"), # 潜在客户 ("scenario_2", "球球妈妈", "球球妈-咨询"), ("scenario_2", "米米妈", "米米妈-咨询"), # 非客户 ("non_customer", "同行-小张", ""), ("non_customer", "快递小哥", ""), ("non_customer", "大学同学", ""), # 群聊 ("group", "宝妈育儿交流群", ""), ("group", "育儿干货分享群", ""), ], 3: [ # 王强 — 中等 ("scenario_4", None, None), # 竞品对比 ("scenario_3", None, None), # 长期跟进复购 ("variant_腹泻咨询后成交", None, None), ("variant_安全顾虑后犹豫", None, None), ("scenario_2", "星星妈", "星星妈-推荐咨询"), # 潜在客户 ("scenario_2", "晨晨妈妈", "晨晨妈-咨询"), ("scenario_2", "露露妈", "露露妈-咨询"), # 非客户 ("non_customer", "健身教练", ""), ("non_customer", "房产中介小王", ""), ("non_customer", "表妹", ""), # 群聊 ("group", "过敏宝宝互助群", ""), ("group", "宝妈育儿交流群", ""), ], } # --------------------------------------------------------------------------- # 生成与同步逻辑 # --------------------------------------------------------------------------- def get_variant(name: str): for v in SCENARIO_VARIANTS: if v["name"] == name: return v raise KeyError(f"variant not found: {name}") def generate_contacts_for_salesperson(sp_id: int): """返回 [(nickname, remark, is_group, wx_username, scenario_or_template)]""" assignments = SALESPERSON_ASSIGNMENTS.get(sp_id, []) result = [] for i, (kind, nick_override, remark_override) in enumerate(assignments): if kind.startswith("scenario_"): idx = int(kind.split("_")[1]) tpl = SCENARIO_TEMPLATES[idx] nick = nick_override or tpl["contact"]["nickname"] remark = remark_override or tpl["contact"].get("remark", "") is_group = False wx_username = f"{SALESPERSONS[sp_id]['wx_account']}_contact_{i+1}" result.append((nick, remark, is_group, wx_username, tpl)) elif kind.startswith("variant_"): tpl = get_variant(kind[8:]) nick = nick_override or tpl["contact"]["nickname"] remark = remark_override or tpl["contact"].get("remark", "") is_group = False wx_username = f"{SALESPERSONS[sp_id]['wx_account']}_contact_{i+1}" result.append((nick, remark, is_group, wx_username, tpl)) elif kind == "non_customer": nick = nick_override remark = remark_override or "" # 找到对应模板 tpl = None for t in NON_CUSTOMER_TEMPLATES: if t["nickname"] == nick: tpl = t break if tpl is None: tpl = NON_CUSTOMER_TEMPLATES[i % len(NON_CUSTOMER_TEMPLATES)] nick = tpl["nickname"] is_group = False wx_username = f"{SALESPERSONS[sp_id]['wx_account']}_contact_{i+1}" result.append((nick, remark, is_group, wx_username, tpl)) elif kind == "group": nick = nick_override remark = remark_override or "" tpl = None for t in GROUP_TEMPLATES: if t["nickname"] == nick: tpl = t break if tpl is None: tpl = GROUP_TEMPLATES[i % len(GROUP_TEMPLATES)] nick = tpl["nickname"] is_group = True wx_username = f"{SALESPERSONS[sp_id]['wx_account']}_group_{i+1}@chatroom" result.append((nick, remark, is_group, wx_username, tpl)) return result def sync_to_central(pg, sp_id: int, contacts_data, run_counter: int = 1): """将联系人、会话、消息写入中央数据库。""" sp_info = SALESPERSONS[sp_id] source_shard = f"mock_{sp_id}" base_time = dt.datetime(2026, 6, 15, 10, 0, 0, tzinfo=TZ) # 每次运行追加的时间偏移 time_offset = dt.timedelta(days=run_counter * 7) contact_ids = {} conversation_ids = {} with pg.cursor() as cur: # 写入联系人 for nick, remark, is_group, wx_username, tpl in contacts_data: display_name = remark or nick cur.execute(""" INSERT INTO contact (salesperson_id, wx_username, nickname, remark, display_name, is_group) VALUES (%s, %s, %s, %s, %s, %s) ON CONFLICT (salesperson_id, wx_username) DO UPDATE SET nickname=excluded.nickname, remark=excluded.remark, display_name=excluded.display_name, is_group=excluded.is_group RETURNING id """, (sp_id, wx_username, nick, remark, display_name, is_group)) contact_ids[wx_username] = cur.fetchone()[0] # 写入会话 for nick, remark, is_group, wx_username, tpl in contacts_data: conv_type = "group" if is_group else "single" cur.execute(""" INSERT INTO conversation (salesperson_id, contact_id, wx_identifier, conv_type, last_synced_at) VALUES (%s, %s, %s, %s, now()) ON CONFLICT (salesperson_id, wx_identifier) DO UPDATE SET contact_id=excluded.contact_id, last_synced_at=now() RETURNING id """, (sp_id, contact_ids[wx_username], wx_username, conv_type)) conversation_ids[wx_username] = cur.fetchone()[0] # 写入消息 local_id_counter = (run_counter - 1) * 10000 # 每次运行从不同区间开始 msg_count = 0 batch = [] for nick, remark, is_group, wx_username, tpl in contacts_data: conv_id = conversation_ids[wx_username] contact_id = contact_ids[wx_username] messages = tpl["messages"] # 确定发送者 sp_wx = sp_info["wx_account"] sp_name = sp_info["name"] for j, (sender_role, msg_type, content) in enumerate(messages): local_id_counter += 1 # 时间递增:每条消息间隔 2~30 分钟 msg_time = base_time + time_offset + dt.timedelta(minutes=j * random.randint(2, 30)) if sender_role == "customer": sender_wx = wx_username if not is_group else f"wxid_customer_{j}" sender_name = nick elif sender_role == "sales": sender_wx = sp_wx sender_name = sp_name else: # other sender_wx = f"wxid_other_{j}" if is_group else wx_username sender_name = nick source_table = "Msg_" + hashlib.md5(wx_username.encode()).hexdigest()[:16] batch.append(( sp_id, conv_id, sender_wx, sender_name, msg_type, content, content, msg_time, source_shard, source_table, local_id_counter, )) msg_count += 1 if len(batch) >= 2000: execute_values(cur, """ INSERT INTO message (salesperson_id, conversation_id, sender_wx_username, sender_display_name, message_type, raw_content, normalized_content, created_at, source_shard, source_table, source_local_id) VALUES %s ON CONFLICT (source_shard, source_table, source_local_id) DO NOTHING """, batch, page_size=2000) pg.commit() batch = [] if batch: execute_values(cur, """ INSERT INTO message (salesperson_id, conversation_id, sender_wx_username, sender_display_name, message_type, raw_content, normalized_content, created_at, source_shard, source_table, source_local_id) VALUES %s ON CONFLICT (source_shard, source_table, source_local_id) DO NOTHING """, batch, page_size=2000) pg.commit() return len(contact_ids), msg_count def get_run_counter(pg, sp_id: int) -> int: """根据已有消息数推断运行次数。""" with pg.cursor() as cur: cur.execute(""" SELECT count(DISTINCT source_local_id / 10000) FROM message WHERE salesperson_id = %s AND source_shard = %s """, (sp_id, f"mock_{sp_id}")) result = cur.fetchone()[0] return int(result or 0) + 1 def main(): ap = argparse.ArgumentParser(description="模拟采集代理") ap.add_argument("--salesperson-id", type=int, required=True) ap.add_argument("--dsn", default=DEFAULT_DSN) args = ap.parse_args() sp_id = args.salesperson_id if sp_id not in SALESPERSONS: print(f"[ERROR] 未知销售 ID: {sp_id}", file=sys.stderr) return 1 sp_info = SALESPERSONS[sp_id] print(f"[INFO] 模拟采集代理启动: {sp_info['name']} (id={sp_id})") pg = psycopg2.connect(args.dsn) pg.autocommit = False try: run_counter = get_run_counter(pg, sp_id) print(f"[INFO] 第 {run_counter} 次同步") contacts_data = generate_contacts_for_salesperson(sp_id) print(f"[INFO] 生成 {len(contacts_data)} 个联系人") contact_count, msg_count = sync_to_central(pg, sp_id, contacts_data, run_counter) with pg.cursor() as cur: cur.execute("SELECT count(*) FROM message WHERE salesperson_id=%s", (sp_id,)) total_msgs = cur.fetchone()[0] print(json.dumps({ "salesperson_id": sp_id, "salesperson_name": sp_info["name"], "run": run_counter, "contacts": contact_count, "new_messages": msg_count, "total_messages": total_msgs, }, ensure_ascii=False, indent=2)) return 0 except Exception as exc: pg.rollback() print(f"[ERROR] {exc}", file=sys.stderr) return 1 finally: pg.close() if __name__ == "__main__": sys.exit(main())