Python 进程池并发执行 SQL 语句:TaoToken 统一 Key 通道下的连接复用与超时重试配置
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_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_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_sourcetaotoken_aicg_blog_endutm_contentconsoleutm_campaignrewrite 生成后在 API Keys 页面管理页面地址 https://taotoken.net/api-keys?utm_sourcetaotoken_aicg_blog_endutm_contentapi-keysutm_campaignrewrite 。Key 的权限建议按最小化原则给只开需要的范围。第三是 Model ID。如果你除了 SQL 还要在流程里调用模型做数据清洗或结果摘要需要指定具体的模型标识。模型对话调试可以用 https://taotoken.net/models?utm_sourcetaotoken_aicg_blog_endutm_contentmodelsutm_campaignrewrite 这个入口先验证通道是否通。为什么要在进程池场景下强调统一通道因为子进程是独立启动的如果每个进程都去读不同的配置文件、连不同的后端出问题时排查成本极高。统一通道之后所有子进程用同一套 Base URL Key日志里一眼就能看出是通道问题还是 SQL 问题。另外密钥轮换时只需要改一处环境变量不用挨个进程改配置。接入文档在 https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_contentdocutm_campaignrewrite 里面有完整的参数说明和示例。如果你用的是 Claude Code 这类编码工具做开发可以参考 https://taotoken.net/claudecode?utm_sourcetaotoken_aicg_blog_endutm_contentclaudecodeutm_campaignrewrite 的接入方式如果是长期跑批量任务的场景Coding Plan 页面 https://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_contentcoding-planutm_campaignrewrite 有配额和并发相关的说明值得先看一眼再决定进程池开多大。这里要提醒一句TaoToken 是接入通道不是数据库本身也不是编辑器替代品。它的作用是让你的进程池在拿连接参数和调用凭证时有个统一出口真正的 SQL 执行还是走你后端的数据库。别把它理解成能直接跑 SQL 的东西。3. 可复制的进程池与超时重试配置这一节是核心直接给可复制的配置。我按配置片段 代码的方式组织你可以整段拿走改。先看环境变量配置建议放在.env或系统环境里# .env 配置片段 TAOTOKEN_BASE_URLhttps://taotoken.net/api TAOTOKEN_API_KEYsk-你的key TAOTOKEN_MODEL_ID你的模型ID DB_HOST127.0.0.1 DB_PORT3306 DB_USERapp_user DB_PASSWORDyour_db_password DB_NAMEapp_db POOL_SIZE8 QUERY_TIMEOUT30 MAX_RETRIES3如果你用 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, paramsNone): 带超时和指数退避重试的 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( levellogging.INFO, formatf[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( processesPOOL_SIZE, initializerworker_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_workersPOOL_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: fBearer {API_KEY}}, json{ model: os.getenv(TAOTOKEN_MODEL_ID), messages: [{role: user, content: ping}], max_tokens: 8, }, timeout15, ) print(resp.status_code, resp.text[:200])返回 200 且 body 里有正常响应说明通道没问题。如果返回 401先检查 Key 有没有多余空格如果返回 404检查 Base URL 是不是写成了带路径的形式。第二步单进程验证数据库连接和超时。故意设一个很短的read_timeout跑一条SELECT SLEEP(5)看是否按预期抛超时import pymysql conn pymysql.connect( host127.0.0.1, port3306, userapp_user, passwordyour_db_password, databaseapp_db, read_timeout2, ) 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 [(fINSERT INTO orders (id, amount) VALUES (%s, %s), (i, i * 10)) for i in range(1000, 1100)] start time.time() with multiprocessing.Pool(processes8) 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 failed401 是 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_SIZE4开始试配合上面的超时重试配置通常能有三到五倍的吞吐提升。跑稳了再往上加。遇到连接类报错回到第 5 节对照排查大部分问题都能定位到具体参数。