760 lines
34 KiB
Python
760 lines
34 KiB
Python
"""
|
|
企业微信聊天记录 → 数据库导出 (CRM 绑定版)
|
|
只导出聊天记录, 并关联用户表生成 unionid / 企微 userid, 便于 CRM 客户绑定。
|
|
|
|
产出 (默认输出目录 crm_import/):
|
|
crm_import_mysql_*.sql MySQL 导入脚本 (建表 + INSERT 数据一体, 直接 mysql < 文件.sql 即可导入)
|
|
crm_wecom.db SQLite 数据库 (zyt_wx_chat_messages / zyt_wx_chat_users / zyt_wx_chat_conversations 三表, 本地备份/预览用)
|
|
聊天记录_仅消息.csv 干净的消息 CSV (只含聊天记录, 一行为一条消息, 含 unionid 列)
|
|
create_tables_mysql.sql MySQL 建表 SQL (单独建表用)
|
|
|
|
MySQL 导入方式:
|
|
mysql -u root -p --default-character-set=utf8mb4 你的库名 < crm_import_mysql_xxx.sql
|
|
或 MySQL 客户端里执行: source crm_import_mysql_xxx.sql;
|
|
|
|
用法:
|
|
python wxwork_export_db.py # 全量导出
|
|
python wxwork_export_db.py --days 1 # 只导今天
|
|
python wxwork_export_db.py --days 7 # 最近 7 天
|
|
python wxwork_export_db.py --date 2026-08-19 # 指定日期
|
|
"""
|
|
|
|
import argparse
|
|
import csv
|
|
import json
|
|
import os
|
|
import re
|
|
import sqlite3
|
|
import sys
|
|
from datetime import datetime, timedelta
|
|
|
|
BASE_DIR = os.path.dirname(os.path.abspath(__file__))
|
|
sys.path.insert(0, BASE_DIR)
|
|
|
|
# 复用主脚本的解析与解密逻辑
|
|
from wxwork_export_final import (
|
|
KEYS_FILE,
|
|
DEFAULT_DB_BASE,
|
|
decrypt_with_keys,
|
|
load_keys,
|
|
parse_content,
|
|
format_timestamp,
|
|
get_msg_type_name,
|
|
is_personal_chat,
|
|
is_blocked_conv_name,
|
|
DEFAULT_BLOCKED_CONV_NAMES,
|
|
connect_sqlite,
|
|
)
|
|
|
|
from wxwork_export_media import export_media as export_media_files
|
|
from wxwork_voice2text import load_voice2text, asr_missing
|
|
|
|
DEFAULT_OUTPUT = os.path.join(BASE_DIR, "crm_import")
|
|
|
|
# ---- 建表 SQL ----
|
|
|
|
SQLITE_SCHEMA = """
|
|
-- 企业微信聊天记录 (CRM 绑定版) - SQLite 建表
|
|
-- zyt_wx_chat_messages 消息表: sender_unionid 为外部联系人的微信 unionid, 可直接与 CRM 客户表关联
|
|
CREATE TABLE IF NOT EXISTS zyt_wx_chat_messages (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
account TEXT, -- 企微账号 (用户目录ID)
|
|
conversation_id TEXT, -- 会话ID (M:单聊 / S:群聊 / R:命名群 / Y:应用)
|
|
conversation_name TEXT, -- 会话名称
|
|
sender_id TEXT, -- 发送者企微用户ID
|
|
sender_name TEXT, -- 发送者名称
|
|
sender_unionid TEXT, -- 发送者微信 unionid (外部联系人, CRM 绑定键)
|
|
sender_userid TEXT, -- 发送者企微 userid (内部成员)
|
|
sender_mobile TEXT, -- 发送者手机号
|
|
msg_type INTEGER, -- 原始消息类型码
|
|
msg_type_name TEXT, -- 消息类型名称
|
|
content TEXT, -- 消息内容
|
|
send_time TEXT, -- 发送时间 (YYYY-MM-DD HH:MM:SS)
|
|
message_seq INTEGER, -- 消息序号
|
|
message_id INTEGER, -- 本地消息ID
|
|
server_id INTEGER, -- 服务器消息ID
|
|
client_id TEXT, -- 客户端消息ID
|
|
media_file_path TEXT, -- 媒体文件本地路径 (图片/语音/视频/文件)
|
|
media_url TEXT, -- CDN URL (本地无缓存时记录)
|
|
voice_text TEXT -- 语音转文字 (企微本地缓存/本地ASR)
|
|
);
|
|
CREATE INDEX IF NOT EXISTS idx_msg_time ON zyt_wx_chat_messages(send_time);
|
|
CREATE INDEX IF NOT EXISTS idx_msg_unionid ON zyt_wx_chat_messages(sender_unionid);
|
|
CREATE INDEX IF NOT EXISTS idx_msg_userid ON zyt_wx_chat_messages(sender_userid);
|
|
CREATE INDEX IF NOT EXISTS idx_msg_conv ON zyt_wx_chat_messages(conversation_id);
|
|
CREATE INDEX IF NOT EXISTS idx_msg_account ON zyt_wx_chat_messages(account);
|
|
|
|
-- zyt_wx_chat_users 联系人表 (CRM 绑定依据)
|
|
CREATE TABLE IF NOT EXISTS zyt_wx_chat_users (
|
|
id INTEGER PRIMARY KEY, -- 企微用户ID
|
|
name TEXT, -- 昵称
|
|
real_name TEXT, -- 真实姓名
|
|
account TEXT, -- 企微 userid (内部成员)
|
|
unionid TEXT, -- 微信 unionid (外部联系人)
|
|
mobile TEXT, -- 手机号
|
|
email TEXT, -- 邮箱
|
|
position TEXT, -- 职位
|
|
gender INTEGER, -- 性别
|
|
external_corp_name TEXT, -- 外部企业名
|
|
external_job TEXT, -- 外部职位
|
|
corp_id INTEGER -- 企业ID
|
|
);
|
|
CREATE INDEX IF NOT EXISTS idx_user_unionid ON zyt_wx_chat_users(unionid);
|
|
CREATE INDEX IF NOT EXISTS idx_user_account ON zyt_wx_chat_users(account);
|
|
|
|
-- zyt_wx_chat_conversations 会话表
|
|
CREATE TABLE IF NOT EXISTS zyt_wx_chat_conversations (
|
|
id TEXT PRIMARY KEY, -- 会话ID
|
|
account TEXT, -- 企微账号
|
|
name TEXT, -- 会话名称
|
|
roomname_remark TEXT, -- 群备注
|
|
session_id TEXT, -- 内部会话ID
|
|
create_time TEXT,
|
|
last_message_time TEXT,
|
|
is_sticked INTEGER,
|
|
is_marked INTEGER,
|
|
is_blocked INTEGER,
|
|
customer_room_type INTEGER
|
|
);
|
|
"""
|
|
|
|
MYSQL_SCHEMA = """-- ============================================================
|
|
-- 企业微信聊天记录 (CRM 绑定版) - MySQL 建表
|
|
-- 在 CRM 数据库执行本脚本后, 把导出的 CSV 导入 zyt_wx_chat_messages 即可
|
|
-- CRM 绑定: crm_customers.unionid = zyt_wx_chat_messages.sender_unionid
|
|
-- ============================================================
|
|
SET NAMES utf8mb4;
|
|
|
|
CREATE TABLE IF NOT EXISTS zyt_wx_chat_messages (
|
|
id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
|
|
account VARCHAR(40) DEFAULT NULL COMMENT '企微账号(用户目录ID)',
|
|
conversation_id VARCHAR(128) DEFAULT NULL COMMENT '会话ID(M:单聊/S:群聊/R:命名群/Y:应用)',
|
|
conversation_name VARCHAR(255) DEFAULT NULL COMMENT '会话名称',
|
|
sender_id VARCHAR(40) DEFAULT NULL COMMENT '发送者企微用户ID',
|
|
sender_name VARCHAR(128) DEFAULT NULL COMMENT '发送者名称',
|
|
sender_unionid VARCHAR(64) DEFAULT NULL COMMENT '发送者微信unionid(外部联系人,CRM绑定键)',
|
|
sender_userid VARCHAR(64) DEFAULT NULL COMMENT '发送者企微userid(内部成员)',
|
|
sender_mobile VARCHAR(32) DEFAULT NULL COMMENT '发送者手机号',
|
|
msg_type INT DEFAULT NULL COMMENT '原始消息类型码',
|
|
msg_type_name VARCHAR(32) DEFAULT NULL COMMENT '消息类型名称',
|
|
content LONGTEXT DEFAULT NULL COMMENT '消息内容',
|
|
send_time DATETIME DEFAULT NULL COMMENT '发送时间',
|
|
message_seq BIGINT DEFAULT NULL COMMENT '消息序号',
|
|
message_id BIGINT DEFAULT NULL COMMENT '本地消息ID',
|
|
server_id BIGINT DEFAULT NULL COMMENT '服务器消息ID',
|
|
client_id VARCHAR(64) DEFAULT NULL COMMENT '客户端消息ID',
|
|
media_file_path VARCHAR(512) DEFAULT NULL COMMENT '媒体文件本地路径(图片/语音/视频/文件)',
|
|
media_url VARCHAR(512) DEFAULT NULL COMMENT 'CDN URL(本地无缓存时记录)',
|
|
voice_text LONGTEXT DEFAULT NULL COMMENT '语音转文字(本地缓存/ASR)',
|
|
PRIMARY KEY (id),
|
|
KEY idx_msg_time (send_time),
|
|
KEY idx_msg_unionid (sender_unionid),
|
|
KEY idx_msg_userid (sender_userid),
|
|
KEY idx_msg_conv (conversation_id),
|
|
KEY idx_msg_account (account)
|
|
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='企业微信聊天记录';
|
|
|
|
CREATE TABLE IF NOT EXISTS zyt_wx_chat_users (
|
|
id BIGINT NOT NULL COMMENT '企微用户ID',
|
|
name VARCHAR(128) DEFAULT NULL COMMENT '昵称',
|
|
real_name VARCHAR(128) DEFAULT NULL COMMENT '真实姓名',
|
|
account VARCHAR(64) DEFAULT NULL COMMENT '企微userid(内部成员)',
|
|
unionid VARCHAR(64) DEFAULT NULL COMMENT '微信unionid(外部联系人)',
|
|
mobile VARCHAR(32) DEFAULT NULL COMMENT '手机号',
|
|
email VARCHAR(128) DEFAULT NULL COMMENT '邮箱',
|
|
position VARCHAR(128) DEFAULT NULL COMMENT '职位',
|
|
gender TINYINT DEFAULT NULL COMMENT '性别',
|
|
external_corp_name VARCHAR(255) DEFAULT NULL COMMENT '外部企业名',
|
|
external_job VARCHAR(128) DEFAULT NULL COMMENT '外部职位',
|
|
corp_id BIGINT DEFAULT NULL COMMENT '企业ID',
|
|
PRIMARY KEY (id),
|
|
KEY idx_user_unionid (unionid),
|
|
KEY idx_user_account (account)
|
|
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='企微联系人(CRM绑定依据)';
|
|
|
|
CREATE TABLE IF NOT EXISTS zyt_wx_chat_conversations (
|
|
id VARCHAR(128) NOT NULL COMMENT '会话ID',
|
|
name VARCHAR(255) DEFAULT NULL COMMENT '会话名称',
|
|
roomname_remark VARCHAR(255) DEFAULT NULL COMMENT '群备注',
|
|
session_id VARCHAR(128) DEFAULT NULL COMMENT '内部会话ID',
|
|
create_time DATETIME DEFAULT NULL,
|
|
last_message_time DATETIME DEFAULT NULL,
|
|
is_sticked TINYINT DEFAULT NULL,
|
|
is_marked TINYINT DEFAULT NULL,
|
|
is_blocked TINYINT DEFAULT NULL,
|
|
customer_room_type TINYINT DEFAULT NULL,
|
|
PRIMARY KEY (id)
|
|
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='企微会话';
|
|
"""
|
|
|
|
|
|
def collect_metadata(decrypted_dbs):
|
|
"""收集用户与会话元数据, 返回 user_cache / conv_cache"""
|
|
user_cache = {} # (user_dir, str(uid)) -> {name, unionid, userid, mobile, ...}
|
|
conv_cache = {} # (user_dir, conv_id) -> name
|
|
|
|
for db_path, db_name, user_dir in decrypted_dbs:
|
|
try:
|
|
conn = connect_sqlite(db_path)
|
|
cur = conn.cursor()
|
|
|
|
if db_name == "user.db":
|
|
try:
|
|
cur.execute("PRAGMA table_info(user_table)")
|
|
cols = [r[1] for r in cur.fetchall()]
|
|
if cols:
|
|
cur.execute("SELECT * FROM user_table")
|
|
col_idx = {c: i for i, c in enumerate(cols)}
|
|
for row in cur.fetchall():
|
|
uid = row[col_idx.get("id", 0)]
|
|
if uid is None:
|
|
continue
|
|
def gv(k):
|
|
i = col_idx.get(k)
|
|
return row[i] if i is not None and i < len(row) else None
|
|
name = gv("name") or gv("real_name") or gv("account") or str(uid)
|
|
user_cache[(user_dir, str(uid))] = {
|
|
"name": str(name),
|
|
"real_name": str(gv("real_name") or ""),
|
|
"unionid": str(gv("unionid") or ""),
|
|
"userid": str(gv("account") or ""),
|
|
"mobile": str(gv("mobile") or ""),
|
|
}
|
|
except sqlite3.Error:
|
|
pass
|
|
|
|
if db_name == "session.db":
|
|
try:
|
|
cur.execute("PRAGMA table_info(conversation_table)")
|
|
cols = [r[1] for r in cur.fetchall()]
|
|
if cols:
|
|
cur.execute("SELECT id, name, roomname_remark, session_id FROM conversation_table")
|
|
for cid, name, remark, sid in cur.fetchall():
|
|
if cid and is_personal_chat(cid):
|
|
nm = remark or name or sid
|
|
conv_cache[(user_dir, str(cid))] = str(nm) if nm else ""
|
|
except sqlite3.Error:
|
|
pass
|
|
|
|
conn.close()
|
|
except sqlite3.Error:
|
|
pass
|
|
|
|
return user_cache, conv_cache
|
|
|
|
|
|
def sql_escape(v):
|
|
"""转义 MySQL 字符串字面量; None -> NULL"""
|
|
if v is None:
|
|
return "NULL"
|
|
s = str(v)
|
|
s = s.replace("\\", "\\\\").replace("'", "''")
|
|
s = s.replace("\r", "\\r").replace("\n", "\\n")
|
|
s = s.replace("\x00", "")
|
|
return "'" + s + "'"
|
|
|
|
|
|
def sql_num(v):
|
|
"""数值字段; None/空 -> NULL"""
|
|
if v is None or v == "":
|
|
return "NULL"
|
|
try:
|
|
return str(int(v))
|
|
except (TypeError, ValueError):
|
|
return "NULL"
|
|
|
|
|
|
MYSQL_MSG_COLS = [
|
|
"account", "conversation_id", "conversation_name", "sender_id", "sender_name",
|
|
"sender_unionid", "sender_userid", "sender_mobile", "msg_type", "msg_type_name",
|
|
"content", "send_time", "message_seq", "message_id", "server_id", "client_id",
|
|
"media_file_path", "media_url", "voice_text",
|
|
]
|
|
|
|
MYSQL_USER_COLS = [
|
|
"id", "name", "real_name", "account", "unionid", "mobile", "email",
|
|
"position", "gender", "external_corp_name", "external_job", "corp_id",
|
|
]
|
|
|
|
MYSQL_CONV_COLS = [
|
|
"id", "name", "roomname_remark", "session_id",
|
|
"create_time", "last_message_time", "is_sticked", "is_marked",
|
|
"is_blocked", "customer_room_type",
|
|
]
|
|
|
|
|
|
def chunked_insert(table, cols, rows, chunk=500, ignore=False):
|
|
"""把行数据拆成多条 INSERT (每条最多 chunk 行), 减少 SQL 体积且兼容任意 MySQL 版本"""
|
|
if not rows:
|
|
return ""
|
|
insert = f"INSERT {'IGNORE ' if ignore else ''}INTO `{table}`\n (`{'`, `'.join(cols)}`)\nVALUES\n"
|
|
parts = []
|
|
for i in range(0, len(rows), chunk):
|
|
batch = rows[i:i + chunk]
|
|
vals = []
|
|
for r in batch:
|
|
vals.append(" (" + ", ".join(r) + ")")
|
|
parts.append(insert + ",\n".join(vals) + ";\n")
|
|
return "\n".join(parts)
|
|
|
|
|
|
def build_mysql_data_sql(messages, users, conversations,
|
|
date_from=None, date_to=None, total_label="全量"):
|
|
"""生成完整 MySQL 导入 SQL: 建表 + INSERT 数据"""
|
|
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
|
lines = []
|
|
lines.append("-- ============================================================")
|
|
lines.append("-- 企业微信聊天记录 (CRM 绑定版) - MySQL 数据导入脚本")
|
|
lines.append(f"-- 生成时间: {now}")
|
|
lines.append(f"-- 数据范围: {total_label}")
|
|
lines.append("-- 使用方式:")
|
|
lines.append("-- mysql -u root -p --default-character-set=utf8mb4 < 本文件")
|
|
lines.append("-- # 或 MySQL 客户端中执行: source 本文件路径;")
|
|
lines.append("-- CRM 绑定: crm_customers.unionid = zyt_wx_chat_messages.sender_unionid")
|
|
lines.append("-- ============================================================")
|
|
lines.append("")
|
|
lines.append("SET NAMES utf8mb4;")
|
|
lines.append("SET FOREIGN_KEY_CHECKS = 0;")
|
|
lines.append("")
|
|
lines.append(MYSQL_SCHEMA.strip())
|
|
lines.append("")
|
|
lines.append(f"-- {total_label}消息: {len(messages)} 条")
|
|
lines.append(chunked_insert("zyt_wx_chat_messages", MYSQL_MSG_COLS, messages))
|
|
lines.append("")
|
|
lines.append(f"-- 联系人: {len(users)} 个")
|
|
lines.append(chunked_insert("zyt_wx_chat_users", MYSQL_USER_COLS, users, ignore=True))
|
|
lines.append("")
|
|
lines.append(f"-- 会话: {len(conversations)} 个")
|
|
lines.append(chunked_insert("zyt_wx_chat_conversations", MYSQL_CONV_COLS, conversations, ignore=True))
|
|
lines.append("")
|
|
lines.append("SET FOREIGN_KEY_CHECKS = 1;")
|
|
lines.append("-- 导入完成")
|
|
return "\n".join(lines)
|
|
|
|
|
|
def export_to_db(decrypted_dbs, out_dir, date_from=None, date_to=None,
|
|
personal_only=True, blocked_conv_names=None,
|
|
voice2text_map=None, voice_asr=False):
|
|
"""只导出聊天记录到数据库 (消息表带 unionid 关联列)
|
|
|
|
personal_only: True=只导出单聊 (过滤群聊/应用消息/第三方应用, 默认); False=导出全部会话
|
|
blocked_conv_names: 会话名称关键词黑名单 (官方/系统账号, 如"企业微信团队"),
|
|
命中即过滤整个会话; None 使用 DEFAULT_BLOCKED_CONV_NAMES, 传空元组 () 关闭
|
|
voice2text_map: 语音转写文本 {(user_dir, str(server_id)): text}, None 自动读取本地缓存
|
|
voice_asr: True=对无本地缓存的语音用本地 AI 识别 (需装 pilk+faster-whisper)
|
|
"""
|
|
os.makedirs(out_dir, exist_ok=True)
|
|
user_cache, conv_cache = collect_metadata(decrypted_dbs)
|
|
if blocked_conv_names is None:
|
|
blocked_conv_names = DEFAULT_BLOCKED_CONV_NAMES
|
|
|
|
# 预解析单聊会话名称, 计算命中黑名单的会话ID (供媒体导出同步过滤)
|
|
def _resolve_name(conv_id, user_dir):
|
|
name = conv_cache.get((user_dir, str(conv_id)), "")
|
|
if name:
|
|
return name
|
|
conv_id = str(conv_id)
|
|
if conv_id.startswith("M:"):
|
|
uid = conv_id[2:]
|
|
return user_cache.get((user_dir, uid), {}).get("name", uid)
|
|
if conv_id.startswith("S:"):
|
|
parts = conv_id[2:].split("_")
|
|
peer_uid = ""
|
|
for p in parts:
|
|
if p != user_dir:
|
|
peer_uid = p
|
|
break
|
|
if not peer_uid:
|
|
peer_uid = parts[0] if parts else ""
|
|
peer_name = user_cache.get((user_dir, peer_uid), {}).get("name", "")
|
|
return peer_name if peer_name else ("单聊 " + conv_id)
|
|
return conv_id
|
|
|
|
blocked_conv_ids = set()
|
|
for (user_dir, cid) in conv_cache:
|
|
if is_personal_chat(cid) and (
|
|
is_blocked_conv_name(_resolve_name(cid, user_dir), blocked_conv_names)
|
|
or is_blocked_conv_name(str(cid), blocked_conv_names)):
|
|
blocked_conv_ids.add(str(cid))
|
|
if blocked_conv_ids:
|
|
print(f"[*] 按会话名称/ID过滤: {len(blocked_conv_ids)} 个官方/系统账号会话将被排除 "
|
|
f"({', '.join(sorted(blocked_conv_ids))[:120]}...)")
|
|
|
|
# 导出媒体文件 (图片/语音/视频/文件, 从 Cache 复制明文文件)
|
|
print("[*] 导出媒体文件 (图片/语音/视频/文件)...")
|
|
media_map, media_stats = export_media_files(
|
|
decrypted_dbs, out_dir, date_from, date_to, personal_only=personal_only,
|
|
blocked_conv_ids=blocked_conv_ids)
|
|
print(f"[+] 媒体导出: 复制 {media_stats['copied']} 个文件, "
|
|
f"UUID匹配 {media_stats['matched_uuid']}, 文件名匹配 {media_stats['matched_filename']}, "
|
|
f"仅URL {media_stats['url_only']}")
|
|
|
|
# 语音转文字 (本地缓存 + 可选本地 ASR 兜底)
|
|
if voice2text_map is None:
|
|
print("[*] 读取语音转写缓存 (msg_voice2text)...")
|
|
voice2text_map = load_voice2text(decrypted_dbs, log=print)
|
|
if voice_asr and media_map:
|
|
print("[*] 本地识别无缓存语音...")
|
|
voice2text_map = asr_missing(voice2text_map, media_map, log=print)
|
|
|
|
def _get_voice_text(user_dir, server_id):
|
|
if not voice2text_map:
|
|
return ""
|
|
sid = str(server_id)
|
|
return voice2text_map.get((user_dir, sid)) or voice2text_map.get(("", sid)) or ""
|
|
|
|
messages = []
|
|
users_seen = {}
|
|
|
|
# 遍历消息
|
|
for db_path, db_name, user_dir in decrypted_dbs:
|
|
if db_name != "message.db":
|
|
continue
|
|
try:
|
|
conn = connect_sqlite(db_path)
|
|
cur = conn.cursor()
|
|
cur.execute("SELECT name FROM sqlite_master WHERE type='table' AND name='message_table'")
|
|
if not cur.fetchone():
|
|
conn.close()
|
|
continue
|
|
cur.execute("PRAGMA table_info(message_table)")
|
|
columns = [r[1] for r in cur.fetchall()]
|
|
cur.execute("SELECT * FROM message_table ORDER BY send_time ASC")
|
|
rows = cur.fetchall()
|
|
conn.close()
|
|
except sqlite3.Error as e:
|
|
print(f" [警告] 读取失败 {db_path}: {e}")
|
|
continue
|
|
|
|
exported = 0
|
|
skipped_non_personal = 0
|
|
skipped_blocked_name = 0
|
|
for row in rows:
|
|
msg = dict(zip(columns, row))
|
|
conv_id = msg.get("conversation_id", "")
|
|
sender_id = msg.get("sender_id", "")
|
|
ts = msg.get("send_time") or 0
|
|
time_str = format_timestamp(ts)
|
|
day = time_str[:10]
|
|
if date_from and day < date_from:
|
|
continue
|
|
if date_to and day > date_to:
|
|
continue
|
|
# 只导出单聊 (过滤群聊 R:/ 应用 Y:/ 第三方应用 O:/ 系统会话)
|
|
if personal_only and not is_personal_chat(conv_id):
|
|
skipped_non_personal += 1
|
|
continue
|
|
|
|
content = parse_content(msg.get("content"))
|
|
extra = parse_content(msg.get("extra_content"))
|
|
if extra and not content.strip():
|
|
if not (extra.startswith('{') and '}' in extra) and 'http' not in extra[:8]:
|
|
content = extra
|
|
|
|
# 发送者信息 (含 unionid / userid / mobile)
|
|
uinfo = user_cache.get((user_dir, str(sender_id)), {})
|
|
sender_unionid = uinfo.get("unionid", "")
|
|
sender_userid = uinfo.get("userid", "")
|
|
sender_mobile = uinfo.get("mobile", "")
|
|
sender_name = uinfo.get("name", "")
|
|
|
|
# 会话名称
|
|
conv_name = _resolve_name(conv_id, user_dir)
|
|
|
|
# 按会话名称/ID过滤 (官方/系统账号, 如"企业微信团队"; 也可填会话ID)
|
|
if is_blocked_conv_name(conv_name, blocked_conv_names) \
|
|
or is_blocked_conv_name(str(conv_id), blocked_conv_names):
|
|
skipped_blocked_name += 1
|
|
continue
|
|
|
|
if not sender_name:
|
|
sender_name = str(sender_id) if sender_id else "系统"
|
|
|
|
# 媒体文件信息 (media_map key 为字符串 server_id)
|
|
media_info = media_map.get(str(msg.get("server_id")), {})
|
|
media_path = media_info.get("path", "")
|
|
media_url = media_info.get("url", "")
|
|
|
|
# 消息类型: 若有本地媒体文件, 按扩展名修正 (截图常被标记为 14/123 等)
|
|
msg_type_name = get_msg_type_name(msg.get("content_type"))
|
|
if media_path:
|
|
ext = os.path.splitext(media_path)[1].lower()
|
|
if ext in (".png", ".jpg", ".jpeg", ".gif", ".bmp", ".webp"):
|
|
msg_type_name = "图片"
|
|
elif ext in (".mp4", ".mov", ".avi", ".mkv"):
|
|
msg_type_name = "视频"
|
|
elif ext in (".silk", ".amr"):
|
|
msg_type_name = "语音"
|
|
elif ext:
|
|
msg_type_name = "文件"
|
|
|
|
# 语音转文字: 语音消息优先显示转写文本
|
|
voice_text = _get_voice_text(user_dir, msg.get("server_id"))
|
|
is_voice = str(media_path).lower().endswith((".silk", ".amr"))
|
|
if is_voice and voice_text:
|
|
content = voice_text
|
|
|
|
messages.append({
|
|
"account": user_dir,
|
|
"conversation_id": str(conv_id),
|
|
"conversation_name": conv_name,
|
|
"sender_id": str(sender_id) if sender_id is not None else "",
|
|
"sender_name": sender_name,
|
|
"sender_unionid": sender_unionid,
|
|
"sender_userid": sender_userid,
|
|
"sender_mobile": sender_mobile,
|
|
"msg_type": msg.get("content_type"),
|
|
"msg_type_name": msg_type_name,
|
|
"content": content,
|
|
"send_time": time_str,
|
|
"message_seq": msg.get("sequence"),
|
|
"message_id": msg.get("message_id"),
|
|
"server_id": msg.get("server_id"),
|
|
"client_id": msg.get("client_id"),
|
|
"media_file_path": media_path,
|
|
"media_url": media_url,
|
|
"voice_text": voice_text,
|
|
})
|
|
exported += 1
|
|
|
|
# 记录出现的用户 (供 zyt_wx_chat_users 表)
|
|
if sender_id is not None and uinfo:
|
|
users_seen[(user_dir, str(sender_id))] = (user_dir, str(sender_id), uinfo)
|
|
|
|
if skipped_non_personal or skipped_blocked_name:
|
|
print(f"[+] {user_dir}/message.db: 导出 {exported} 条消息 "
|
|
f"(过滤非单聊 {skipped_non_personal} 条, "
|
|
f"过滤官方/系统账号 {skipped_blocked_name} 条)")
|
|
else:
|
|
print(f"[+] {user_dir}/message.db: 导出 {exported} 条消息")
|
|
print(f" unionid 覆盖: {sum(1 for m in messages if m['account']==user_dir and m['sender_unionid'])} 条")
|
|
|
|
if not messages:
|
|
print("[-] 该时间范围内没有消息")
|
|
return None
|
|
|
|
messages.sort(key=lambda m: m["send_time"])
|
|
|
|
# 写入 SQLite (每次重建表结构, 保证与最新 schema 一致)
|
|
db_path = os.path.join(out_dir, "crm_wecom.db")
|
|
conn = sqlite3.connect(db_path)
|
|
for t in ("zyt_wx_chat_messages", "zyt_wx_chat_users", "zyt_wx_chat_conversations"):
|
|
conn.execute(f"DROP TABLE IF EXISTS {t}")
|
|
conn.executescript(SQLITE_SCHEMA)
|
|
|
|
cols = ["account", "conversation_id", "conversation_name", "sender_id", "sender_name",
|
|
"sender_unionid", "sender_userid", "sender_mobile", "msg_type", "msg_type_name",
|
|
"content", "send_time", "message_seq", "message_id", "server_id", "client_id",
|
|
"media_file_path", "media_url", "voice_text"]
|
|
conn.executemany(
|
|
f"INSERT INTO zyt_wx_chat_messages ({','.join(cols)}) VALUES ({','.join(['?']*len(cols))})",
|
|
[tuple(m[c] for c in cols) for m in messages])
|
|
|
|
# zyt_wx_chat_users: 全部联系人 (不只消息中出现过的, 保证 CRM 绑定覆盖面)
|
|
all_users = {}
|
|
for udb_path, udb_name, u_dir in decrypted_dbs:
|
|
if udb_name != "user.db":
|
|
continue
|
|
try:
|
|
conn2 = connect_sqlite(udb_path)
|
|
cur2 = conn2.cursor()
|
|
cur2.execute("PRAGMA table_info(user_table)")
|
|
cols2 = [r[1] for r in cur2.fetchall()]
|
|
if "id" in cols2:
|
|
cur2.execute("SELECT * FROM user_table")
|
|
ci = {c: i for i, c in enumerate(cols2)}
|
|
for row in cur2.fetchall():
|
|
uid = row[ci["id"]]
|
|
if uid is None:
|
|
continue
|
|
def gv(k):
|
|
i = ci.get(k)
|
|
return row[i] if i is not None and i < len(row) else None
|
|
all_users[(u_dir, str(uid))] = (
|
|
u_dir, str(uid), str(gv("name") or ""), str(gv("real_name") or ""),
|
|
str(gv("account") or ""), str(gv("unionid") or ""), str(gv("mobile") or ""),
|
|
str(gv("email") or ""), str(gv("position") or ""),
|
|
gv("gender"), str(gv("external_corp_name") or ""),
|
|
str(gv("external_job") or ""), gv("corp_id"),
|
|
)
|
|
conn2.close()
|
|
except sqlite3.Error:
|
|
pass
|
|
|
|
conn.executemany(
|
|
"""INSERT OR REPLACE INTO zyt_wx_chat_users
|
|
(id, name, real_name, account, unionid, mobile, email, position,
|
|
gender, external_corp_name, external_job, corp_id)
|
|
VALUES (?,?,?,?,?,?,?,?,?,?,?,?)""",
|
|
[(u[1], u[2], u[3], u[4], u[5], u[6], u[7], u[8], u[9], u[10], u[11], u[12])
|
|
for u in all_users.values()])
|
|
|
|
# zyt_wx_chat_conversations: 只写实际导出的单聊会话 (名称取解析后的真实用户名)
|
|
conv_rows = {}
|
|
for m in messages:
|
|
cid = m["conversation_id"]
|
|
if cid not in conv_rows:
|
|
conv_rows[cid] = (cid, m["conversation_name"], m["account"])
|
|
conn.executemany(
|
|
"INSERT OR REPLACE INTO zyt_wx_chat_conversations (id, name, account) VALUES (?,?,?)",
|
|
list(conv_rows.values()))
|
|
|
|
conn.commit()
|
|
conn.close()
|
|
print(f"\n[+] SQLite 数据库: {db_path}")
|
|
|
|
# 统计
|
|
conn = sqlite3.connect(db_path)
|
|
n_msg = conn.execute("SELECT COUNT(*) FROM zyt_wx_chat_messages").fetchone()[0]
|
|
n_uid = conn.execute("SELECT COUNT(*) FROM zyt_wx_chat_messages WHERE sender_unionid != ''").fetchone()[0]
|
|
n_user = conn.execute("SELECT COUNT(*) FROM zyt_wx_chat_users").fetchone()[0]
|
|
n_uid_user = conn.execute("SELECT COUNT(*) FROM zyt_wx_chat_users WHERE unionid != ''").fetchone()[0]
|
|
conn.close()
|
|
print(f" 消息 {n_msg} 条 | 带 unionid {n_uid} 条 ({n_uid*100//n_msg if n_msg else 0}%)")
|
|
print(f" 联系人 {n_user} 个 | 带 unionid {n_uid_user} 个")
|
|
|
|
# 导出干净 CSV (只聊天记录)
|
|
csv_path = os.path.join(out_dir, "聊天记录_仅消息.csv")
|
|
fieldnames = ["account", "conversation_id", "conversation_name", "sender_id", "sender_name",
|
|
"sender_unionid", "sender_userid", "sender_mobile", "msg_type", "msg_type_name",
|
|
"content", "send_time", "message_seq", "message_id", "server_id", "client_id",
|
|
"media_file_path", "media_url", "voice_text"]
|
|
with open(csv_path, "w", newline="", encoding="utf-8-sig") as f:
|
|
w = csv.DictWriter(f, fieldnames=fieldnames)
|
|
w.writeheader()
|
|
w.writerows(messages)
|
|
print(f"[+] 干净 CSV (仅聊天记录): {csv_path}")
|
|
|
|
# 生成 MySQL 建表 SQL
|
|
mysql_path = os.path.join(out_dir, "create_tables_mysql.sql")
|
|
with open(mysql_path, "w", encoding="utf-8") as f:
|
|
f.write(MYSQL_SCHEMA)
|
|
print(f"[+] MySQL 建表 SQL: {mysql_path}")
|
|
|
|
# 生成 MySQL 数据导入 SQL (建表 + INSERT 一体, 可直接导入)
|
|
msg_rows = []
|
|
for m in messages:
|
|
msg_rows.append([
|
|
sql_escape(m["account"]), sql_escape(m["conversation_id"]),
|
|
sql_escape(m["conversation_name"]), sql_escape(m["sender_id"]),
|
|
sql_escape(m["sender_name"]), sql_escape(m["sender_unionid"]),
|
|
sql_escape(m["sender_userid"]), sql_escape(m["sender_mobile"]),
|
|
sql_num(m["msg_type"]), sql_escape(m["msg_type_name"]),
|
|
sql_escape(m["content"]), sql_escape(m["send_time"]),
|
|
sql_num(m["message_seq"]), sql_num(m["message_id"]),
|
|
sql_num(m["server_id"]), sql_escape(m["client_id"]),
|
|
sql_escape(m["media_file_path"]), sql_escape(m["media_url"]),
|
|
sql_escape(m.get("voice_text", "")),
|
|
])
|
|
|
|
user_rows = []
|
|
for u in all_users.values():
|
|
# u = (u_dir, uid, name, real_name, account, unionid, mobile, email,
|
|
# position, gender, external_corp_name, external_job, corp_id)
|
|
user_rows.append([
|
|
sql_num(u[1]), sql_escape(u[2]), sql_escape(u[3]), sql_escape(u[4]),
|
|
sql_escape(u[5]), sql_escape(u[6]), sql_escape(u[7]), sql_escape(u[8]),
|
|
sql_num(u[9]), sql_escape(u[10]), sql_escape(u[11]), sql_num(u[12]),
|
|
])
|
|
|
|
conv_rows_sql = []
|
|
for c in conv_rows:
|
|
# c = (cid, cname, account); MySQL 表无 account 列
|
|
conv_rows_sql.append([
|
|
sql_escape(c[0]), sql_escape(c[1]),
|
|
"NULL", "NULL", "NULL", "NULL", "NULL", "NULL", "NULL", "NULL",
|
|
])
|
|
|
|
if date_from and date_to:
|
|
if date_from == date_to:
|
|
label = f"单日 {date_from}"
|
|
fname = f"crm_import_mysql_{date_from}.sql"
|
|
else:
|
|
label = f"{date_from} ~ {date_to}"
|
|
fname = f"crm_import_mysql_{date_from}_至_{date_to}.sql"
|
|
else:
|
|
label = "全量历史"
|
|
fname = "crm_import_mysql_all.sql"
|
|
|
|
mysql_data_sql = build_mysql_data_sql(
|
|
msg_rows, user_rows, conv_rows_sql,
|
|
date_from=date_from, date_to=date_to, total_label=label)
|
|
|
|
data_path = os.path.join(out_dir, fname)
|
|
with open(data_path, "w", encoding="utf-8") as f:
|
|
f.write(mysql_data_sql)
|
|
print(f"[+] MySQL 数据导入 SQL: {data_path}")
|
|
print(f" 直接导入: mysql -u root -p --default-character-set=utf8mb4 < {fname}")
|
|
|
|
# 生成 SQLite 建表 SQL (参考用)
|
|
sqlite_path = os.path.join(out_dir, "create_tables_sqlite.sql")
|
|
with open(sqlite_path, "w", encoding="utf-8") as f:
|
|
f.write("-- 企业微信聊天记录 (CRM 绑定版) - SQLite 建表参考\n-- 实际使用时脚本已自动创建, 此文件仅供查看/迁移\n")
|
|
f.write(SQLITE_SCHEMA)
|
|
print(f"[+] SQLite 建表 SQL: {sqlite_path}")
|
|
|
|
return db_path
|
|
|
|
|
|
def main():
|
|
parser = argparse.ArgumentParser(description="企业微信聊天记录 → 数据库 (CRM 绑定版)")
|
|
parser.add_argument("--output", default=DEFAULT_OUTPUT, help="输出目录 (默认 crm_import/)")
|
|
parser.add_argument("--db-dir", default=DEFAULT_DB_BASE, help="企业微信数据库根目录")
|
|
parser.add_argument("--date", default=None, help="导出指定日期 (YYYY-MM-DD)")
|
|
parser.add_argument("--days", type=int, default=None, help="导出最近 N 天 (含今天)")
|
|
parser.add_argument("--all", action="store_true", help="全量导出")
|
|
parser.add_argument("--with-groups", action="store_true",
|
|
help="同时导出群聊/应用消息 (默认只导出单聊)")
|
|
args = parser.parse_args()
|
|
|
|
print("=" * 60)
|
|
print(" 企业微信聊天记录 → 数据库导出 (CRM 绑定版, 仅单聊)")
|
|
print("=" * 60)
|
|
|
|
keys_map = load_keys()
|
|
if not keys_map:
|
|
print("[-] 未找到密钥文件")
|
|
return 1
|
|
print(f"[+] 已加载 {len(keys_map)} 个账号密钥")
|
|
if args.with_groups:
|
|
print("[*] 范围: 全部会话 (单聊 + 群聊 + 应用)")
|
|
else:
|
|
print("[*] 范围: 仅单聊 (群聊/应用消息/第三方应用已过滤)")
|
|
|
|
date_from = date_to = None
|
|
if args.all:
|
|
print("[*] 模式: 全量导出")
|
|
elif args.date:
|
|
date_from = date_to = args.date
|
|
print(f"[*] 模式: 导出指定日期 {args.date}")
|
|
else:
|
|
days = args.days if args.days and args.days > 0 else 1
|
|
date_to = datetime.now().strftime("%Y-%m-%d")
|
|
date_from = (datetime.now() - timedelta(days=days - 1)).strftime("%Y-%m-%d")
|
|
print(f"[*] 模式: 最近 {days} 天 ({date_from} ~ {date_to})")
|
|
|
|
print("\n[*] 解密数据库...")
|
|
decrypted_dbs = decrypt_with_keys(
|
|
args.db_dir, os.path.join(args.output, "decrypted"), keys_map, use_cache=True)
|
|
print(f"[+] 成功解密 {len(decrypted_dbs)} 个数据库")
|
|
|
|
if not decrypted_dbs:
|
|
print("[-] 没有可用的数据库")
|
|
return 1
|
|
|
|
print("\n[*] 导出聊天记录到数据库...")
|
|
db_path = export_to_db(decrypted_dbs, args.output,
|
|
date_from=date_from, date_to=date_to,
|
|
personal_only=not args.with_groups)
|
|
|
|
if db_path:
|
|
print(f"\n{'='*60}")
|
|
print(f" 完成! 数据库文件: {db_path}")
|
|
print(f" 输出目录: {args.output}")
|
|
print(f"{'='*60}")
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|