Files
wxsales/demo/ai/analyze.py
T

481 lines
19 KiB
Python
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.
#!/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 = ["其他异议"]
# 摘要文本:根据阶段和关键词生成结构化摘要
customer_name = messages[0]["sender_display_name"] if messages else "客户"
# 提取客户需求关键词
needs = []
for kw in SYMPTOM_WORDS + ["减肥", "减脂", "瘦身", "体脂", "产后"]:
if kw in all_text and kw not in needs:
needs.append(kw)
needs_str = "".join(needs[:3]) if needs else "未明确"
# 判断是否有朋友推荐
referred = "朋友推荐" if any(w in all_text for w in ["推荐", "小鹿", "小林", "朋友"]) else "主动咨询"
# 判断产品线
is_diet = any(kw in all_text for kw in ["减肥", "减脂", "瘦身", "B420", "纤益", "代餐", "体脂", "产后减肥"])
if stage == "成交":
# 提取成交产品
if "2盒" in all_text or "656" in all_text:
deal_detail = "2盒套餐656元" if is_diet else "3盒套餐798元"
elif "单盒" in all_text or "328" in all_text:
deal_detail = "单盒328元" if is_diet else "单盒298元"
elif "6盒" in all_text or "1499" in all_text:
deal_detail = "6盒套餐1499元"
else:
deal_detail = "套餐"
summary = f"客户{customer_name}{needs_str}),{referred},经产品介绍和案例展示后认可产品,最终成交{deal_detail}"
elif stage == "异议处理":
top_obj = objections[0] if objections else "价格"
summary = f"客户{customer_name}{needs_str}),{referred},对产品有一定兴趣但存在{top_obj}方面的顾虑,尚未成交,需持续跟进。"
elif stage == "报价":
summary = f"客户{customer_name}{needs_str}),{referred},已进入报价阶段,等待客户反馈。"
elif stage == "需求发现":
summary = f"客户{customer_name}{needs_str}),{referred},正在了解客户具体需求和产品适配度。"
else:
summary = f"客户{customer_name}{needs_str}),{referred},初步建立联系,尚未深入沟通。"
# 下一步建议
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())