#!/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()