#!/usr/bin/env python3 """AI 分析服务 — 客户识别 + 沟通摘要(规则引擎 / LLM 模式)。 用法: python analyze.py --mode rule python analyze.py --mode llm python analyze.py --mode llm --force 环境变量: DASHSCOPE_API_KEY — 千问 API Key """ from __future__ import annotations import argparse import datetime as dt import json import os import pathlib import sys import urllib.request import urllib.error import psycopg2 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" DASHSCOPE_URL = "https://dashscope.aliyuncs.com/compatible-mode/v1/chat/completions" DASHSCOPE_MODEL = "qwen-plus" # --------------------------------------------------------------------------- # 关键词定义 # --------------------------------------------------------------------------- SYMPTOM_WORDS = [ "湿疹", "过敏", "鼻炎", "腹泻", "便秘", "免疫力", "体质", "红疹", "胀气", "拉肚子", "打喷嚏", "流鼻涕", "感冒", "食物过敏", "牛奶蛋白过敏", "减肥", "减脂", "瘦身", "产后", "体脂", "肚子", "大腿", "反弹", "代餐", ] PRODUCT_WORDS = [ "益生菌", "菌株", "进口", "配方", "成分", "丹麦", "鼠李糖", "无敏", "LGG", "临床", "调理", "肠道", "免疫", "B420", "纤益", "脂肪代谢", "减脂", "肠道代谢", "膳食纤维", ] PRICE_WORDS = [ "价格", "多少钱", "费用", "报价", "优惠", "套餐", "盒", "周期", "298", "798", "1499", "328", "656", "贵", "便宜", "划算", ] DEAL_WORDS = [ "下单", "付款", "发货", "来几盒", "定了", "来3盒", "来6盒", "已付", "转过去", "签收", "已经转", "麻烦尽快发", "来2盒", "付了", "已付款", ] OBJECTION_WORDS = [ "贵", "太贵", "考虑", "商量", "老公", "副作用", "安全吗", "有没有效", "没用过", "合生元", "对比", "担心", "再看看", "网上", "几十块", "哺乳期", "反弹", "代餐", "节食", "拉肚子", ] # 阶段关键词(按优先级排序) STAGE_KEYWORDS = [ ("成交", DEAL_WORDS), ("报价", ["多少钱", "价格", "套餐", "298", "798", "1499", "328", "656", "报价"]), ("异议处理", OBJECTION_WORDS), ("需求发现", SYMPTOM_WORDS + ["什么情况", "多大了", "几个月", "症状"]), ("建立联系", ["你好", "在吗", "请问"]), ] def contains_any(text: str, words: list[str]) -> bool: return any(w in text for w in words) def count_hits(text: str, words: list[str]) -> int: return sum(1 for w in words if w in text) # --------------------------------------------------------------------------- # 客户识别 # --------------------------------------------------------------------------- def identify_customer(messages: list[dict], salesperson_name: str | None = None) -> dict | None: """规则引擎判断联系人是否为客户。返回 None 表示不是客户。""" if len(messages) <= 5: return None # 检查销售是否参与了对话(非客户联系人不会有销售回复) if salesperson_name: has_sales_reply = any( m["sender_display_name"] == salesperson_name for m in messages ) if not has_sales_reply: return None all_text = " ".join(m["normalized_content"] or "" for m in messages) symptom_hits = count_hits(all_text, SYMPTOM_WORDS) product_hits = count_hits(all_text, PRODUCT_WORDS) price_hits = count_hits(all_text, PRICE_WORDS) deal_hits = count_hits(all_text, DEAL_WORDS) total_keyword_hits = symptom_hits + product_hits + price_hits + deal_hits if total_keyword_hits < 2: return None # 意向等级 if deal_hits >= 1 and price_hits >= 1: intent_level = "high" elif (symptom_hits >= 1 and product_hits >= 1) or price_hits >= 1: intent_level = "medium" else: intent_level = "low" # 提取关键需求 key_needs = [] for symptom in SYMPTOM_WORDS: if symptom in all_text and symptom not in key_needs: key_needs.append(symptom) if len(key_needs) >= 3: break # 行业:减肥相关标记为年轻女性,否则宝妈 diet_words = ["减肥", "减脂", "瘦身", "B420", "纤益", "代餐", "体脂", "产后减肥"] industry = "年轻女性" if any(w in all_text for w in diet_words) else "宝妈" return { "is_customer": True, "customer_name": None, # 后续用 display_name 填充 "industry": industry, "intent_level": intent_level, "key_needs": key_needs, "reason": f"消息数{len(messages)}条,命中关键词{total_keyword_hits}个(症状{symptom_hits}+产品{product_hits}+价格{price_hits}+成交{deal_hits})", } # --------------------------------------------------------------------------- # 沟通摘要 # --------------------------------------------------------------------------- def generate_summary(messages: list[dict], mode: str = "rule") -> dict: """生成沟通摘要。mode='rule' 用规则引擎,mode='llm' 用千问 LLM。""" if mode == "llm": llm_result = generate_summary_llm(messages) if llm_result: return llm_result print("[WARN] LLM 调用失败,回退到规则模式", file=sys.stderr) return generate_summary_rule(messages) def generate_summary_llm(messages: list[dict]) -> dict | None: """调用千问 LLM 生成沟通摘要。""" api_key = os.environ.get("DASHSCOPE_API_KEY") if not api_key: print("[WARN] 未设置 DASHSCOPE_API_KEY,回退到规则模式", file=sys.stderr) return None # 构建对话文本 dialog_lines = [] for m in messages: sender = m["sender_display_name"] or "未知" content = m["normalized_content"] or f"[{m['message_type']}]" dialog_lines.append(f"{sender}: {content}") dialog_text = "\n".join(dialog_lines) # 截断过长的对话 if len(dialog_text) > 8000: dialog_text = dialog_text[:4000] + "\n...(中间部分省略)...\n" + dialog_text[-4000:] prompt = f"""你是一个微信销售对话分析助手。以下是销售与客户的微信对话记录,请分析并输出 JSON 格式的沟通摘要。 ## 产品背景 1. 益童宝儿童益生菌粉:丹麦进口菌株,主打小儿抗过敏(湿疹、鼻炎、食物过敏),298元/盒,3盒套餐798元,6盒套餐1499元。目标客户是宝妈。 2. 纤益减肥益生菌:含专利B420菌株,主打健康减脂不反弹,不节食不腹泻,通过调节肠道菌群促进脂肪代谢。328元/盒(30天用量),2盒套餐656元。目标客户是18-35岁有体重管理需求的年轻女性。 ## 对话记录 {dialog_text} ## 输出要求 请输出严格的 JSON,字段如下: - summary: 沟通核心要点总结(2-4句话,概括客户需求、销售策略、关键转折和结果) - stage: 销售阶段,取值之一:建立联系/需求发现/报价/异议处理/成交 - key_points: 如果已成交,列出关键成交点(促成成交的关键因素,如信任建立、案例背书、价格拆解等);如果未成交,列出需要突破的问题点(阻碍成交的核心原因,如价格顾虑、安全担忧、决策链阻碍等)。返回数组,每条一句话。 - objections: 客户提出的异议或顾虑列表(如价格贵、安全顾虑、效果质疑等,空数组表示无) - next_action: 下一步行动建议(具体可执行) 只输出 JSON,不要其他文字。""" body = json.dumps({ "model": DASHSCOPE_MODEL, "messages": [ {"role": "system", "content": "你是微信销售对话分析助手,擅长从对话中提取核心信息。"}, {"role": "user", "content": prompt}, ], "response_format": {"type": "json_object"}, "temperature": 0.3, }).encode("utf-8") req = urllib.request.Request( DASHSCOPE_URL, data=body, headers={ "Authorization": f"Bearer {api_key}", "Content-Type": "application/json", }, method="POST", ) try: print(f" [LLM] 调用千问 ({len(messages)} 条消息)...", file=sys.stderr, flush=True) with urllib.request.urlopen(req, timeout=60) as resp: result = json.loads(resp.read().decode("utf-8")) content = result["choices"][0]["message"]["content"] parsed = json.loads(content) return { "summary": parsed.get("summary", ""), "stage": parsed.get("stage", "建立联系"), "key_points": parsed.get("key_points", []), "objections": parsed.get("objections", []), "next_action": parsed.get("next_action", "继续跟进"), } except Exception as exc: print(f"[WARN] LLM 调用失败: {exc}", file=sys.stderr) return None def generate_summary_rule(messages: list[dict]) -> dict: """规则引擎生成沟通摘要。""" all_text = " ".join(m["normalized_content"] or "" for m in messages) # 阶段判断(按优先级) stage = "建立联系" for stage_name, keywords in STAGE_KEYWORDS: if contains_any(all_text, keywords): stage = stage_name break # 异议提取 objections = [] for word in OBJECTION_WORDS: if word in all_text and word not in objections: objections.append(word) if not objections and stage == "异议处理": objections = ["其他异议"] # 摘要文本:取前 5 条和后 3 条消息拼接 summary_parts = [] for m in messages[:5]: sender = m["sender_display_name"] or "未知" content = m["normalized_content"] or f"[{m['message_type']}]" summary_parts.append(f"{sender}: {content}") if len(messages) > 8: summary_parts.append("...") for m in messages[-3:]: sender = m["sender_display_name"] or "未知" content = m["normalized_content"] or f"[{m['message_type']}]" summary_parts.append(f"{sender}: {content}") summary = " | ".join(summary_parts) # 下一步建议 next_action = "继续跟进" if stage == "成交": next_action = "安排发货并跟进使用效果,适时推荐复购" elif stage == "异议处理": next_action = "发送案例和科普文章,3天后跟进" elif stage == "报价": next_action = "等待客户反馈,适时推荐套餐" elif stage == "需求发现": next_action = "深入了解客户需求,介绍产品优势" elif stage == "建立联系": next_action = "保持互动,寻找需求切入点" # 关键点(成交点或突破点) key_points = [] if stage == "成交": if "临床" in all_text or "验证" in all_text: key_points.append("临床验证数据增强信任") if "反馈" in all_text or "案例" in all_text: key_points.append("真实案例背书打消顾虑") if "798" in all_text or "套餐" in all_text: key_points.append("套餐价格拆解降低价格敏感度") if "售后" in all_text or "指导" in all_text: key_points.append("售后承诺降低试错风险") if not key_points: key_points.append("客户需求明确,销售及时跟进促成成交") else: if "贵" in all_text or "太贵" in all_text: key_points.append("价格顾虑:需进一步拆解日均成本或强调产品差异") if "安全" in all_text or "副作用" in all_text: key_points.append("安全顾虑:需提供更多无敏配方和临床数据") if "考虑" in all_text or "商量" in all_text or "老公" in all_text: key_points.append("决策链阻碍:需提供资料支持客户与家人沟通") if "合生元" in all_text or "对比" in all_text: key_points.append("竞品对比:需强化菌株差异和临床优势") if "效果" in all_text and ("不好" in all_text or "没用" in all_text): key_points.append("效果质疑:需提供更多真实反馈和售后保障") if not key_points: key_points.append("客户意向尚不明确,需持续跟进挖掘需求") return { "summary": summary[:500], "stage": stage, "key_points": key_points, "objections": objections, "next_action": next_action, } # --------------------------------------------------------------------------- # 主流程 # --------------------------------------------------------------------------- def main(): ap = argparse.ArgumentParser(description="AI 分析服务") ap.add_argument("--dsn", default=DEFAULT_DSN) ap.add_argument("--mode", choices=["rule", "llm"], default="rule") ap.add_argument("--force", action="store_true", help="强制重新分析") args = ap.parse_args() if args.mode == "llm" and not os.environ.get("DASHSCOPE_API_KEY"): print("[WARN] 未设置 DASHSCOPE_API_KEY 环境变量,将使用规则模式", file=sys.stderr) pg = psycopg2.connect(args.dsn) pg.autocommit = False try: with pg.cursor() as cur: # 获取需要分析的联系人 if args.force: cur.execute(""" SELECT c.id, c.salesperson_id, c.display_name, c.is_group, c.wx_username FROM contact c WHERE c.is_group = FALSE ORDER BY c.id """) else: cur.execute(""" SELECT c.id, c.salesperson_id, c.display_name, c.is_group, c.wx_username FROM contact c LEFT JOIN customer cu ON cu.contact_id = c.id AND cu.salesperson_id = c.salesperson_id WHERE c.is_group = FALSE AND (cu.id IS NULL OR cu.last_analysis IS NULL) ORDER BY c.id """) contacts = cur.fetchall() customers_created = 0 summaries_generated = 0 for contact_id, sp_id, display_name, is_group, wx_username in contacts: # 获取该联系人的所有消息 with pg.cursor() as cur: cur.execute(""" SELECT m.normalized_content, m.sender_display_name, m.message_type, m.created_at, conv.id FROM message m JOIN conversation conv ON conv.id = m.conversation_id WHERE conv.contact_id = %s ORDER BY m.created_at """, (contact_id,)) rows = cur.fetchall() if not rows: continue messages = [ { "normalized_content": r[0], "sender_display_name": r[1], "message_type": r[2], "created_at": r[3], } for r in rows ] # 获取销售名用于过滤 with pg.cursor() as cur2: cur2.execute("SELECT name FROM salesperson WHERE id=%s", (sp_id,)) sp_name = cur2.fetchone()[0] # 客户识别 result = identify_customer(messages, salesperson_name=sp_name) if result is None: # 不是客户,跳过 continue conv_id = rows[0][4] # 沟通摘要 summary_result = generate_summary(messages, mode=args.mode) # 写入 customer 表 with pg.cursor() as cur: cur.execute(""" INSERT INTO customer (salesperson_id, contact_id, customer_name, industry, intent_level, key_needs, reason, summary, stage, key_points, objections, next_action, last_analysis) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, now()) ON CONFLICT (salesperson_id, contact_id) DO UPDATE SET customer_name=excluded.customer_name, industry=excluded.industry, intent_level=excluded.intent_level, key_needs=excluded.key_needs, reason=excluded.reason, summary=excluded.summary, stage=excluded.stage, key_points=excluded.key_points, objections=excluded.objections, next_action=excluded.next_action, last_analysis=now() """, ( sp_id, contact_id, display_name, result["industry"], result["intent_level"], result["key_needs"], result["reason"], summary_result["summary"], summary_result["stage"], summary_result["key_points"], summary_result["objections"], summary_result["next_action"], )) customers_created += 1 summaries_generated += 1 pg.commit() # 统计 with pg.cursor() as cur: cur.execute("SELECT count(*) FROM customer") total_customers = cur.fetchone()[0] cur.execute("SELECT intent_level, count(*) FROM customer GROUP BY intent_level ORDER BY intent_level") intent_dist = {r[0]: r[1] for r in cur.fetchall()} cur.execute("SELECT stage, count(*) FROM customer GROUP BY stage ORDER BY count(*) DESC") stage_dist = {r[0]: r[1] for r in cur.fetchall()} print(json.dumps({ "mode": args.mode, "contacts_analyzed": len(contacts), "customers_identified": customers_created, "summaries_generated": summaries_generated, "total_customers": total_customers, "intent_distribution": intent_dist, "stage_distribution": stage_dist, }, 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())