""" 企业微信聊天记录 → 数据库导出 (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())