1. 新增客户联系方式查询客户每日总结的完整功能,包括前端页面、后端接口和数据库字段支持 2. 优化原有每日总结页面,新增绑定客户联系方式功能和无总结时的聊天记录兜底展示 3. 修复消息查询时的to_user数组匹配问题,新增客户微信号存储字段 4. 补充批量修复客户昵称的脚本和部署文档更新
135 lines
4.7 KiB
Python
135 lines
4.7 KiB
Python
#!/usr/bin/env python3
|
||
# -*- coding: utf-8 -*-
|
||
# 批量修复 customers 表中 name = customer_id 的客户昵称(企微 externalcontact/get 回填,多线程版)
|
||
# 在服务器上以 nohup 运行:nohup python3 /tmp/fix_customer_names.py > /tmp/fix_customer_names.log 2>&1 &
|
||
import json
|
||
import subprocess
|
||
import threading
|
||
import time
|
||
import urllib.request
|
||
import urllib.parse
|
||
from concurrent.futures import ThreadPoolExecutor
|
||
|
||
DB_HOST = 'sh-cdb-6fzlwnms.sql.tencentcdb.com'
|
||
DB_PORT = '63912'
|
||
DB_USER = 'assistant_prod'
|
||
DB_PASS = 'U%$4Tu_C3+4'
|
||
DB_NAME = 'ai_assistant'
|
||
CORPID = 'wwd483c2fba24ae30a'
|
||
SECRET = 'nU1xezyHyonS9rOJ9Kwo3i5OM9iuWc1TrjXf-rOYZx0'
|
||
WORKERS = 6
|
||
|
||
MYSQL = ['docker', 'exec', '-i', 'ai-assistant-mysql', 'mysql',
|
||
'-h', DB_HOST, '-P', DB_PORT, '-u', DB_USER, '-p' + DB_PASS,
|
||
'--default-character-set=utf8mb4', DB_NAME]
|
||
|
||
lock = threading.Lock()
|
||
state = {'token': '', 'fixed': 0, 'failed': 0, 'done': 0, 'buf': []}
|
||
|
||
|
||
def http_json(url):
|
||
return json.loads(urllib.request.urlopen(url, timeout=15).read().decode('utf-8'))
|
||
|
||
|
||
def get_token():
|
||
d = http_json('https://qyapi.weixin.qq.com/cgi-bin/gettoken?corpid=%s&corpsecret=%s' % (CORPID, SECRET))
|
||
return d.get('access_token', '')
|
||
|
||
|
||
def sql_escape(s):
|
||
return s.replace('\\', '\\\\').replace("'", "\\'")
|
||
|
||
|
||
def flush_buf():
|
||
with lock:
|
||
if not state['buf']:
|
||
return
|
||
sql = '\n'.join(state['buf']) + '\n'
|
||
state['buf'] = []
|
||
r = subprocess.run(MYSQL, input=sql.encode('utf-8'),
|
||
stdout=subprocess.PIPE, stderr=subprocess.PIPE)
|
||
if r.returncode != 0:
|
||
print('SQL apply error: %s' % r.stderr.decode('utf-8', 'ignore')[:200], flush=True)
|
||
|
||
|
||
def process(cid, total):
|
||
for attempt in range(3):
|
||
try:
|
||
with lock:
|
||
token = state['token']
|
||
url = ('https://qyapi.weixin.qq.com/cgi-bin/externalcontact/get'
|
||
'?access_token=%s&external_userid=%s' % (token, urllib.parse.quote(cid)))
|
||
d = http_json(url)
|
||
ec = d.get('errcode')
|
||
if ec == 0:
|
||
c = d.get('external_contact') or {}
|
||
name = (c.get('name') or '').strip()
|
||
if name:
|
||
sets = ["name='%s'" % sql_escape(name)]
|
||
av = (c.get('avatar') or '').strip()
|
||
if av:
|
||
sets.append("avatar='%s'" % sql_escape(av))
|
||
stmt = "UPDATE customers SET %s WHERE customer_id='%s';" % (', '.join(sets), sql_escape(cid))
|
||
with lock:
|
||
state['buf'].append(stmt)
|
||
state['fixed'] += 1
|
||
n = len(state['buf'])
|
||
if n >= 200:
|
||
flush_buf()
|
||
else:
|
||
with lock:
|
||
state['failed'] += 1
|
||
break
|
||
elif ec in (40014, 42001):
|
||
with lock:
|
||
state['token'] = get_token()
|
||
time.sleep(0.5)
|
||
elif ec == 45009:
|
||
time.sleep(2)
|
||
else:
|
||
with lock:
|
||
state['failed'] += 1
|
||
break
|
||
except Exception as e:
|
||
if attempt == 2:
|
||
with lock:
|
||
state['failed'] += 1
|
||
print('error %s: %s' % (cid, str(e)[:100]), flush=True)
|
||
time.sleep(1)
|
||
time.sleep(0.03)
|
||
with lock:
|
||
state['done'] += 1
|
||
done = state['done']
|
||
fixed = state['fixed']
|
||
failed = state['failed']
|
||
if done % 500 == 0:
|
||
print('progress %d/%d fixed=%d failed=%d' % (done, total, fixed, failed), flush=True)
|
||
|
||
|
||
def main():
|
||
q = subprocess.run(['docker', 'exec', '-i', 'ai-assistant-mysql', 'mysql',
|
||
'-h', DB_HOST, '-P', DB_PORT, '-u', DB_USER, '-p' + DB_PASS,
|
||
'-N', '--default-character-set=utf8mb4', DB_NAME,
|
||
'-e', "SELECT customer_id FROM customers WHERE name = customer_id OR name IS NULL OR name = ''"],
|
||
stdout=subprocess.PIPE, stderr=subprocess.PIPE)
|
||
cids = [l.strip() for l in q.stdout.decode('utf-8', 'ignore').splitlines()
|
||
if l.strip().startswith(('wm', 'wo'))]
|
||
total = len(cids)
|
||
print('candidates: %d' % total, flush=True)
|
||
|
||
state['token'] = get_token()
|
||
if not state['token']:
|
||
print('FATAL: no access token', flush=True)
|
||
return
|
||
print('token acquired, workers=%d' % WORKERS, flush=True)
|
||
|
||
with ThreadPoolExecutor(max_workers=WORKERS) as pool:
|
||
for cid in cids:
|
||
pool.submit(process, cid, total)
|
||
|
||
flush_buf()
|
||
print('DONE fixed=%d failed=%d' % (state['fixed'], state['failed']), flush=True)
|
||
|
||
|
||
main()
|