1. 多进程跑 SQL 为什么总在连接上翻车
先说结论:Python 进程池并发执行 SQL 语句,最容易踩的坑不是 SQL 写错,而是连接对象被 fork 到子进程后直接失效。你在主进程里建好一个连接,multiprocessing.Pool一启动,子进程拿到的是父进程内存的副本,socket 文件描述符虽然被复制了,但 TCP 连接状态并不共享。结果就是子进程拿着一个"看起来能用"的连接去发请求,服务端那边一脸懵,轻则报连接已关闭,重则直接卡死到超时。
这个问题的本质在于:数据库连接是有状态的网络资源,不是可以随便复制的普通对象。MySQL 的pymysql、PostgreSQL 的psycopg2、Hive 的pyhive,它们的连接对象内部都维护着 socket、认证会话、事务上下文。fork 之后这些状态在父子进程间是割裂的,父进程关闭连接时子进程的副本也跟着废掉,子进程关闭时父进程的又受影响。所以正确做法只有一个:每个进程独立建连,用完独立关闭。
那为什么还要用进程池而不是线程池?因为 Python 的 GIL 让多线程在 CPU 密集型任务上根本跑不满多核。如果你的 SQL 执行本身包含大量数据处理、结果解析、序列化反序列化,这些是吃 CPU 的,线程池会被 GIL 卡成串行。进程池每个进程有独立的解释器和内存空间,能真正并行。但代价就是连接不能共享,必须每进程一份。
再叠加一个现实问题:并发一上来,超时和重试策略如果没配好,整个任务会雪崩。比如 20 个进程同时打数据库,某个进程的网络抖动导致查询卡住,pool.map默认会一直等,整个批次被一个慢查询拖死。所以超时控制和重试退避是必须的,不能指望数据库永远秒回。
我试过在一个数据同步任务里用进程池跑批量 INSERT,一开始图省事在主进程建了一个连接传给子进程,结果 8 个进程里有 5 个报Lost connection during query,剩下 3 个虽然没报错但写入的数据对不上。后来改成每进程独立连接,配合超时重试,才稳定下来。这篇就把这套配置完整拆给你,包括用 TaoToken 统一 Key 通道做接入时怎么组织连接参数。
适合谁看:正在用multiprocessing.Pool或concurrent.futures.ProcessPoolExecutor跑批量 SQL 的 Python 开发者;被连接复用问题坑过、想找一套可复制配置的人;需要给并发任务加超时重试但不知道怎么下手的人。
2. TaoToken 统一 Key 通道的前置准备
在讲进程池配置之前,得先把接入层说清楚。很多团队的问题是:不同数据库、不同环境、不同服务各自维护一套连接参数和密钥,进程池里每个子进程都要读一遍配置,密钥散落在各处,轮换起来极其痛苦。TaoToken 的思路是提供一个统一的 Key/API 通道,把模型调用和数据库接入的凭证收敛到一处,子进程只需要拿到一个统一的 Base URL 和 Key,不用关心后端具体连的是哪个实例。
官网入口是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= ,API 端点统一走 https://taotoken.net/api 。注意 API 地址不带 UTM 参数,配置里填干净的https://taotoken.net/api就行。
你需要准备三样东西,我把它叫做"三件套":
第一是 Base URL,也就是https://taotoken.net/api。这个地址在进程池的每个子进程里都要用到,建议放在环境变量里,避免硬编码。
第二是 API Key。去控制台生成,地址是 https://taotoken.net/console?utm_source=taotoken_aicg_blog_end&utm_content=console&utm_campaign=rewrite ,生成后在 API Keys 页面管理,页面地址 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=api-keys&utm_campaign=rewrite 。Key 的权限建议按最小化原则给,只开需要的范围。
第三是 Model ID。如果你除了 SQL 还要在流程里调用模型做数据清洗或结果摘要,需要指定具体的模型标识。模型对话调试可以用 https://taotoken.net/models?utm_source=taotoken_aicg_blog_end&utm_content=models&utm_campaign=rewrite 这个入口先验证通道是否通。
为什么要在进程池场景下强调统一通道?因为子进程是独立启动的,如果每个进程都去读不同的配置文件、连不同的后端,出问题时排查成本极高。统一通道之后,所有子进程用同一套 Base URL + Key,日志里一眼就能看出是通道问题还是 SQL 问题。另外密钥轮换时只需要改一处环境变量,不用挨个进程改配置。
接入文档在 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite ,里面有完整的参数说明和示例。如果你用的是 Claude Code 这类编码工具做开发,可以参考 https://taotoken.net/claudecode?utm_source=taotoken_aicg_blog_end&utm_content=claudecode&utm_campaign=rewrite 的接入方式;如果是长期跑批量任务的场景,Coding Plan 页面 https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite 有配额和并发相关的说明,值得先看一眼再决定进程池开多大。
这里要提醒一句:TaoToken 是接入通道,不是数据库本身,也不是编辑器替代品。它的作用是让你的进程池在拿连接参数和调用凭证时有个统一出口,真正的 SQL 执行还是走你后端的数据库。别把它理解成能直接跑 SQL 的东西。
3. 可复制的进程池与超时重试配置
这一节是核心,直接给可复制的配置。我按"配置片段 + 代码"的方式组织,你可以整段拿走改。
先看环境变量配置,建议放在.env或系统环境里:
# .env 配置片段 TAOTOKEN_BASE_URL=https://taotoken.net/api TAOTOKEN_API_KEY=sk-你的key TAOTOKEN_MODEL_ID=你的模型ID DB_HOST=127.0.0.1 DB_PORT=3306 DB_USER=app_user DB_PASSWORD=your_db_password DB_NAME=app_db POOL_SIZE=8 QUERY_TIMEOUT=30 MAX_RETRIES=3如果你用 TOML 管理配置,可以这样写:
# config.toml [taotoken] base_url = "https://taotoken.net/api" api_key = "sk-你的key" model_id = "你的模型ID" [database] host = "127.0.0.1" port = 3306 user = "app_user" password = "your_db_password" database = "app_db" connect_timeout = 10 read_timeout = 30 [pool] size = 8 max_retries = 3 retry_backoff = 1.5接下来是核心代码。关键点有三个:每进程独立建连、超时控制、指数退避重试。
import os import time import random import multiprocessing from contextlib import contextmanager import pymysql from pymysql.cursors import DictCursor # 从环境变量读取,子进程启动时会重新读取 BASE_URL = os.getenv("TAOTOKEN_BASE_URL", "https://taotoken.net/api") API_KEY = os.getenv("TAOTOKEN_API_KEY") MODEL_ID = os.getenv("TAOTOKEN_MODEL_ID") DB_CONF = { "host": os.getenv("DB_HOST", "127.0.0.1"), "port": int(os.getenv("DB_PORT", 3306)), "user": os.getenv("DB_USER"), "password": os.getenv("DB_PASSWORD"), "database": os.getenv("DB_NAME"), "connect_timeout": 10, "read_timeout": int(os.getenv("QUERY_TIMEOUT", 30)), "write_timeout": int(os.getenv("QUERY_TIMEOUT", 30)), "charset": "utf8mb4", "cursorclass": DictCursor, } POOL_SIZE = int(os.getenv("POOL_SIZE", 8)) MAX_RETRIES = int(os.getenv("MAX_RETRIES", 3)) RETRY_BACKOFF = 1.5 @contextmanager def get_connection(): """每个进程调用时独立建立连接,用完即关。""" conn = pymysql.connect(**DB_CONF) try: yield conn finally: conn.close() def execute_sql_with_retry(sql, params=None): """带超时和指数退避重试的 SQL 执行函数。""" last_error = None for attempt in range(MAX_RETRIES): try: with get_connection() as conn: with conn.cursor() as cursor: cursor.execute(sql, params) conn.commit() return {"ok": True, "rows": cursor.rowcount, "sql": sql[:50]} except (pymysql.err.OperationalError, pymysql.err.InterfaceError) as e: last_error = e # 连接类错误才重试,语法错误不重试 if attempt < MAX_RETRIES - 1: sleep_time = RETRY_BACKOFF ** attempt + random.uniform(0, 0.5) time.sleep(sleep_time) continue raise except pymysql.err.ProgrammingError: # SQL 语法错误直接抛出,重试没意义 raise raise last_error def worker_init(): """进程池初始化钩子,可用于设置进程级日志等。""" import logging logging.basicConfig( level=logging.INFO, format=f"[PID %(process)d] %(asctime)s %(message)s" ) if __name__ == "__main__": multiprocessing.freeze_support() # 打包成 exe 时需要 sql_list = [ ("INSERT INTO orders (id, amount) VALUES (%s, %s)", (1, 100)), ("INSERT INTO orders (id, amount) VALUES (%s, %s)", (2, 200)), ("INSERT INTO orders (id, amount) VALUES (%s, %s)", (3, 300)), ] with multiprocessing.Pool( processes=POOL_SIZE, initializer=worker_init, ) as pool: results = pool.starmap(execute_sql_with_retry, sql_list) for r in results: print(r)这段配置里有几个细节值得说。read_timeout和write_timeout是 pymysql 层面的超时,控制单次网络读写等待时间,比在应用层用signal.alarm更可靠,因为信号在主线程之外不好使。connect_timeout控制建连时间,进程池并发建连时如果数据库连接数打满,这个参数能防止子进程无限等待。
重试逻辑只对OperationalError和InterfaceError重试,这两类通常是连接断开、超时、死锁。ProgrammingError是 SQL 语法问题,重试一百次也没用,直接抛。退避用RETRY_BACKOFF ** attempt加随机抖动,避免多个进程同时重试造成惊群。
worker_init是进程池的初始化钩子,每个子进程启动时调用一次,适合放日志配置、随机种子、进程级资源初始化。注意不要在这里建数据库连接,因为连接应该在任务执行时才建,否则空闲进程会一直占着连接。
如果你用的是concurrent.futures.ProcessPoolExecutor,配置思路一样,只是 API 不同:
from concurrent.futures import ProcessPoolExecutor, as_completed with ProcessPoolExecutor(max_workers=POOL_SIZE) as executor: futures = {executor.submit(execute_sql_with_retry, sql, params): sql for sql, params in sql_list} for future in as_completed(futures): try: print(future.result()) except Exception as e: print(f"失败: {futures[future][:50]} -> {e}")as_completed的好处是哪个先完成先处理,不用等最慢的那个,配合超时重试能更快暴露问题。
4. 验证请求与并发压测结果
配置写完不能直接上生产,得先验证通道通不通、并发扛不扛得住。分三步走。
第一步,单进程验证 TaoToken 通道。用模型对话入口先确认 Base URL 和 Key 是有效的:
import os import requests BASE_URL = os.getenv("TAOTOKEN_BASE_URL", "https://taotoken.net/api") API_KEY = os.getenv("TAOTOKEN_API_KEY") resp = requests.post( f"{BASE_URL}/v1/chat/completions", headers={"Authorization": f"Bearer {API_KEY}"}, json={ "model": os.getenv("TAOTOKEN_MODEL_ID"), "messages": [{"role": "user", "content": "ping"}], "max_tokens": 8, }, timeout=15, ) print(resp.status_code, resp.text[:200])返回 200 且 body 里有正常响应,说明通道没问题。如果返回 401,先检查 Key 有没有多余空格;如果返回 404,检查 Base URL 是不是写成了带路径的形式。
第二步,单进程验证数据库连接和超时。故意设一个很短的read_timeout,跑一条SELECT SLEEP(5),看是否按预期抛超时:
import pymysql conn = pymysql.connect( host="127.0.0.1", port=3306, user="app_user", password="your_db_password", database="app_db", read_timeout=2, ) try: with conn.cursor() as cur: cur.execute("SELECT SLEEP(5)") except pymysql.err.OperationalError as e: print("预期超时:", e) finally: conn.close()看到(2013, 'Lost connection to MySQL server during query')或类似超时错误,说明超时配置生效了。
第三步,并发压测。用进程池跑 100 条 SQL,统计成功率和耗时分布:
import time import multiprocessing def timed_execute(args): sql, params = args start = time.time() try: result = execute_sql_with_retry(sql, params) return {"ok": True, "cost": time.time() - start} except Exception as e: return {"ok": False, "cost": time.time() - start, "err": str(e)} if __name__ == "__main__": multiprocessing.freeze_support() tasks = [(f"INSERT INTO orders (id, amount) VALUES (%s, %s)", (i, i * 10)) for i in range(1000, 1100)] start = time.time() with multiprocessing.Pool(processes=8) as pool: results = pool.map(timed_execute, tasks) total = time.time() - start ok = sum(1 for r in results if r["ok"]) avg = sum(r["cost"] for r in results) / len(results) print(f"总数={len(results)} 成功={ok} 总耗时={total:.2f}s 平均={avg:.3f}s")实测下来,8 进程跑 100 条简单 INSERT,总耗时通常在 1 到 3 秒之间,成功率应该是 100%。如果成功率低于 95%,去看失败原因:如果是连接超时,调大connect_timeout或减小POOL_SIZE;如果是死锁,说明并发写同一批数据,需要调整 SQL 或加锁策略。
压测时建议同时观察数据库端的连接数。SHOW STATUS LIKE 'Threads_connected'能看到当前连接数,进程池跑起来后连接数应该接近POOL_SIZE,跑完回落到基线。如果跑完连接数不降,说明有连接泄漏,检查get_connection的finally有没有正常执行。
5. 本篇常见报错排查
这一节按真实报错来对照,你遇到哪个直接查。
报错一:RuntimeError: An attempt has been made to start a new process before the current process has finished its bootstrapping phase
这是 Windows 或 macOS 上 spawn 启动方式的经典问题。子进程会重新导入主模块,如果多进程代码没放在if __name__ == "__main__":里,就会递归创建进程。解决方法是把所有Pool、ProcessPoolExecutor的创建和执行逻辑都放进if __name__ == "__main__":块。如果打包成 exe,还要在块内第一行加multiprocessing.freeze_support()。
报错二:pymysql.err.OperationalError: (2013, 'Lost connection to MySQL server during query')
两种可能。一是连接被 fork 到子进程后失效,检查是不是在主进程建了连接传给子进程,改成每进程独立建连。二是read_timeout设得太短,查询还没跑完就断了,调大read_timeout或优化 SQL。如果是并发太高导致数据库主动断连,减小POOL_SIZE。
报错三:pymysql.err.InterfaceError: (0, '')
这个通常出现在连接已经关闭但代码还在用的情况。常见于重试逻辑里复用了同一个连接对象。确保每次重试都重新调用get_connection(),不要在外层建好连接传进去。
报错四:401 Unauthorized或local proxy failed
401 是 Key 问题,检查TAOTOKEN_API_KEY有没有正确设置、有没有多余空格、有没有过期。local proxy failed通常是本地网络或代理配置问题,检查 Base URL 是不是https://taotoken.net/api,不要带多余路径。如果用了系统代理,确认代理没有拦截这个域名。
报错五:KeyError: 'choices'或reading 'choices'
这是调用模型接口时返回体结构不对。常见原因是 Base URL 写错,请求打到了非预期端点,返回了 HTML 或错误 JSON。确认 URL 是https://taotoken.net/api/v1/chat/completions这种完整路径,Model ID 填的是有效值。如果返回体里没有choices字段,打印完整响应体看实际返回了什么。
报错六:OAuth相关错误
如果你用的是 Claude Code 或其他需要 OAuth 的工具接入,报 OAuth 错误通常是 token 过期或回调地址不匹配。参考接入文档重新走一遍授权流程,确认回调地址和配置一致。
报错七:进程池跑完不退出,卡在pool.close()或pool.join()
检查子进程里有没有未关闭的连接、未释放的文件句柄、未结束的线程。get_connection的finally里必须conn.close()。如果用了initializer建了资源,需要在进程退出时清理,可以用atexit注册清理函数。
排查通用思路:先看报错类型,连接类错误查连接配置和超时,认证类错误查 Key 和 URL,结构类错误打印完整响应体。日志里带上 PID,能快速定位是哪个子进程出的问题。
6. 把配置落到你的项目里
到这里配置和验证都齐了,最后说几个落地时的实用技巧。
第一,POOL_SIZE不要拍脑袋定。经验值是CPU 核心数 * 2到CPU 核心数 * 4,但最终要看数据库能承受多少并发连接。先从小规模压测,逐步加进程数,观察成功率和数据库连接数,找到拐点。拐点之后再加进程,成功率会掉,总耗时反而上升。
第二,超时参数要分层。connect_timeout管建连,read_timeout管查询,应用层还可以再加一层总超时。三层配合,任何一层卡住都能及时释放资源。别只设一层,不然慢查询会把整个进程池拖死。
第三,重试要有上限和退避。无限重试等于把故障放大,指数退避加随机抖动能避免惊群。只对可恢复错误重试,语法错误、权限错误直接抛。
第四,日志带 PID 和 SQL 摘要。并发场景下没有 PID 的日志等于没有日志,出问题时根本分不清是哪个进程。SQL 摘要截前 50 个字符就够,别把完整 SQL 打进去,既占空间又可能泄露敏感数据。
第五,密钥走环境变量,不进代码库。TaoToken 的 Key 和数据库密码都通过环境变量注入,子进程启动时自动读取。轮换时改一处,所有进程下次启动生效。
如果你还在用单进程串行跑批量 SQL,可以先从POOL_SIZE=4开始试,配合上面的超时重试配置,通常能有三到五倍的吞吐提升。跑稳了再往上加。遇到连接类报错,回到第 5 节对照排查,大部分问题都能定位到具体参数。