
1. 为什么 pymysql 连 Doris 总在批量写入时断线Doris 前端 FE 默认走 MySQL 协议端口 9030所以 Python 侧最省事的做法就是pymysql。但真把它放到生产脚本里跑批量导入很多人会撞上同一个报错pymysql.err.OperationalError: (2013, Lost connection to MySQL server during query)这个 2013 不是 SQL 写错了而是连接在查询过程中被对端断掉了。Doris 的 FE 有wait_timeout、interactive_timeout这类空闲回收机制默认值通常在 8 小时量级但实际生产里更常见的是脚本跑一个长事务、中间隔了几分钟没发请求或者 FE 做了滚动重启、负载均衡把连接切走下一次cursor.execute就直接炸。我试过最原始的写法就是 excerpt 里那种import pymysql.cursors conn pymysql.connect(host127.0.0.1, port9030, userroot, password, databasessb, charsetutf8) mycursor conn.cursor() mycursor.execute(select * from lineorder_rt limit 1) result mycursor.fetchall() for data in result: print(data) mycursor.close() conn.close()单次查询没问题但它有三个致命伤没有事务边界、没有重连、游标和连接的生命周期全靠手写。批量写入场景下一旦中途断线前面execute过的数据可能已经进了 Doris重跑就重复不重跑就丢数据。所以这篇要解决的核心就三件事用上下文管理器管住 cursor 和事务、用ping(reconnectTrue)做断线探测、用重试骨架把「断线后重试」变成可验证的动作。适合正在写数据同步脚本、批量导入任务、定时调度 ETL 的同学代码可以直接套。顺带说一句如果你后面还要接大模型做数据清洗或 SQL 生成统一走一个 Key 通道会省很多事TaoToken 的 API 入口是 https://taotoken.net/api 模型对话、Coding Plan、API Keys 都在同一套体系里后面 §6 会给具体分流。2. 前置准备pymysql 连接 Doris 的参数与 TaoToken 通道先把环境铺好。Doris 侧你需要确认三样东西FE 的 MySQL 协议端口默认 9030、一个有写入权限的账号、目标库名。Python 侧只需要一个包pip install pymysql版本上pymysql1.0对ping(reconnectTrue)的支持更稳建议锁 1.0 以上。连接参数里几个容易踩的点参数建议值说明host / portFE 地址 / 9030走 MySQL 协议不是 8030 HTTPcharsetutf8mb4避免中文和 emoji 写入乱码autocommitFalse手动控制事务批量写入必须关connect_timeout10建连超时别用默认read_timeout60单次读超时长查询要调大write_timeout60单次写超时cursorclassDictCursor可选返回字典更易读这里有个关键认知Doris 的事务能力和 MySQL InnoDB 不是一回事。Doris 的写入是批量导入模型commit更多是提交这一批的可见性回滚能力有限。所以封装事务时重点是「保证一批要么整体提交、要么整体重试」而不是指望rollback能撤销所有已写入数据。这一点想清楚后面的重试策略才不会写歪。再说 TaoToken 的位置。它不替代 Doris也不替代 pymysql它是你脚本里调用大模型能力时的统一入口。比如你想让模型根据表结构自动生成 INSERT 语句、或者对导入失败的行做语义纠错就需要一个稳定的 API 通道。TaoToken 提供统一的 Key 和 API 地址省去你到处配不同厂商 endpoint 的麻烦。接入前先在控制台建 Keyhttps://taotoken.net/console?utm_sourcetaotoken_aicg_blog_endutm_contentconsoleutm_campaignrewrite建完 Key 后模型对话调试入口https://taotoken.net/models?utm_sourcetaotoken_aicg_blog_endutm_contentmodelsutm_campaignrewrite如果你是要长期跑编码类 Agent比如让模型帮你写 Doris 导入脚本、做 SQL 审查那更适合用 Coding Planhttps://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_contentcoding-planutm_campaignrewriteAPI Key 管理页https://taotoken.net/api-keys?utm_sourcetaotoken_aicg_blog_endutm_contentapi-keysutm_campaignrewrite文档https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_contentdocutm_campaignrewriteClaude Code 接入 Anthropic 的说明页https://taotoken.net/claude-code-anthropic?utm_sourcetaotoken_aicg_blog_endutm_contentClaudeCodeAnthropicutm_campaignrewrite这些链接先记着§6 会讲什么场景该点哪个。现在回到 Doris 连接本身。3. 可复制的连接封装事务骨架 自动重连配置这一节是全文核心直接给能跑的代码。先看连接工厂把参数集中管理import pymysql import pymysql.cursors import time import logging logging.basicConfig(levellogging.INFO, format%(asctime)s %(levelname)s %(message)s) logger logging.getLogger(doris) DORIS_CONF { host: 127.0.0.1, port: 9030, user: root, password: , database: ads, charset: utf8mb4, autocommit: False, connect_timeout: 10, read_timeout: 60, write_timeout: 60, cursorclass: pymysql.cursors.DictCursor, } def create_conn(): return pymysql.connect(**DORIS_CONF)注意autocommitFalse这是事务控制的前提。接下来是带重连的执行器。核心思路每次执行前先ping(reconnectTrue)探活断了就自动重建连接执行失败按次数退避重试。class DorisClient: def __init__(self, conf, max_retry3, backoff1.5): self.conf conf self.max_retry max_retry self.backoff backoff self.conn create_conn() def _ensure_alive(self): try: self.conn.ping(reconnectTrue) except Exception as e: logger.warning(ping failed, rebuild conn: %s, e) self.conn create_conn() def execute(self, sql, paramsNone, fetchNone): attempt 0 while attempt self.max_retry: attempt 1 try: self._ensure_alive() with self.conn.cursor() as cursor: cursor.execute(sql, params) self.conn.commit() if fetch one: return cursor.fetchone() if fetch all: return cursor.fetchall() if fetch many: return cursor.fetchmany(100) return cursor.rowcount except pymysql.err.OperationalError as e: code e.args[0] if e.args else None logger.error(op error code%s attempt%s msg%s, code, attempt, e) if code in (2006, 2013): self.conn create_conn() time.sleep(self.backoff ** attempt) except Exception as e: logger.error(non-retryable error: %s, e) raise raise RuntimeError(execute failed after retries: %s % sql)几个设计点解释一下。ping(reconnectTrue)是 pymysql 自带的探活它会在连接失效时尝试重连比你自己判断conn.open靠谱。with self.conn.cursor()保证游标一定被关闭不会泄漏。commit放在execute之后、fetch之前是因为 Doris 的写入需要显式提交才可见。批量写入场景建议用executemany包一层def execute_many(self, sql, rows): attempt 0 while attempt self.max_retry: attempt 1 try: self._ensure_alive() with self.conn.cursor() as cursor: cursor.executemany(sql, rows) self.conn.commit() return cursor.rowcount except pymysql.err.OperationalError as e: code e.args[0] if e.args else None logger.error(batch op error code%s attempt%s, code, attempt) if code in (2006, 2013): self.conn create_conn() time.sleep(self.backoff ** attempt) raise RuntimeError(batch insert failed after retries)如果你用配置文件管理连接可以放一份 TOML路径按你项目来比如conf/doris.toml[doris] host 127.0.0.1 port 9030 user root password database ads charset utf8mb4 autocommit false connect_timeout 10 read_timeout 60 write_timeout 60 max_retry 3 backoff 1.5读取时用tomllibPython 3.11或tomliimport tomllib with open(conf/doris.toml, rb) as f: cfg tomllib.load(f)[doris]这样连接参数和重试策略就和代码解耦了改配置不用动逻辑。到这里事务骨架和自动重连就齐了下一节验证它到底管不管用。4. 验证请求断线重试与成功结果实测代码写完不验证等于没写。这一节做两个动作正常查询验证、模拟断线验证重连。先建一张测试表并写入client DorisClient(DORIS_CONF) client.execute( CREATE TABLE IF NOT EXISTS ads.demo_orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10,2), dt DATE ) ENGINEOLAP DUPLICATE KEY(order_id) DISTRIBUTED BY HASH(order_id) BUCKETS 3 PROPERTIES (replication_num 1) ) rows [ (1001, 2001, 99.50, 2024-06-01), (1002, 2002, 150.00, 2024-06-01), (1003, 2003, 23.80, 2024-06-02), ] affected client.execute_many( INSERT INTO ads.demo_orders VALUES (%s, %s, %s, %s), rows) print(inserted rows:, affected)正常情况会打印inserted rows: 3。接着查回来one client.execute( SELECT * FROM ads.demo_orders WHERE order_id%s, (1001,), fetchone) print(fetchone:, one) all_rows client.execute( SELECT * FROM ads.demo_orders ORDER BY order_id, fetchall) print(fetchall count:, len(all_rows)) many client.execute( SELECT * FROM ads.demo_orders, fetchmany) print(fetchmany:, many)fetchone返回单条字典fetchall返回全部fetchmany(100)按批返回这三个在 Doris 上行为和 MySQL 一致。实测下来DictCursor配合fetchall在结果集不大时最省心。现在验证重连。手动把连接关掉模拟断线client.conn.close() print(conn closed, open flag:, client.conn.open) # 此时再执行_ensure_alive 会 ping 失败并重建 result client.execute( SELECT COUNT(*) AS cnt FROM ads.demo_orders, fetchone) print(after reconnect:, result)预期输出类似conn closed, open flag: False 2024-06-10 10:00:00 WARNING ping failed, rebuild conn: (0, ) after reconnect: {cnt: 3}ping(reconnectTrue)在连接已关闭时会抛异常被_ensure_alive捕获后重建连接查询照常成功。这就是自动重连的验证动作。再模拟一次「查询中途断线」的重试。这个不好直接造但可以用一个不存在的表触发非重试错误确认它不会无限重试try: client.execute(SELECT * FROM ads.not_exist_table, fetchone) except Exception as e: print(expected error:, type(e).__name__)非 2006/2013 的错误会直接抛出不会进重试循环避免把语法错误当成网络问题反复重试。这个边界很重要很多人的重连逻辑写成了「所有异常都重试」结果 SQL 写错也重试三次白白拖慢任务。如果你在脚本里还要调模型做辅助比如把失败行交给模型分析可以这样接 TaoTokenimport requests resp requests.post( https://taotoken.net/api/v1/chat/completions, headers{Authorization: Bearer YOUR_TAOTOKEN_KEY}, json{ model: MODEL_ID, messages: [{role: user, content: 分析这行导入失败原因}] }, timeout30, ) print(resp.json())注意 Base URL 用https://taotoken.net/apiKey 从 API Keys 页拿Model ID 按你选的模型填。这三件套Base URL Key Model ID缺一不可后面 §5 会专门讲配错会报什么。5. 常见报错排查2013、401、local proxy failed 与 OAuth这一节按真实报错来对。你跑上面的代码最可能撞到下面几类。第一类2013 Lost connection during query。这是本文主线。触发原因通常是空闲超时、FE 重启、网络抖动。排查顺序先看ping(reconnectTrue)有没有生效再看max_retry和backoff是否合理。如果重试后仍失败检查 FE 的wait_timeoutSHOW VARIABLES LIKE %timeout%;如果wait_timeout很小可以在连接参数里加init_command或定期发心跳。但更稳的做法就是本文的_ensure_alive每次执行前探活。第二类2006 MySQL server has gone away。和 2013 是兄弟都是连接层断了。代码里已经把它和 2013 一起处理触发重建。区别是 2006 多发生在发送请求时2013 多发生在读取结果时。第三类401 Unauthorized。这个不是 Doris 的是你调 TaoToken 或其它 API 时 Key 错了。检查三件套Base URL 是不是https://taotoken.net/apiKey 有没有多余空格Model ID 是否拼错。401 基本就是 Key 无效或没带上Authorization头。第四类local proxy failed。这个报错通常出现在你本地配了某种转发但目标不可达时。排查方向是确认你的请求地址写对了、网络能通、没有多余的本地转发配置干扰。把 Base URL 换成官方地址重试往往就好了。第五类OAuth 相关报错。如果你用 Claude Code 接 Anthropic 通道可能会遇到 OAuth 认证失败。这类问题看接入文档最直接https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_contentdocutm_campaignrewriteClaude Code 的专门说明在https://taotoken.net/claude-code-anthropic?utm_sourcetaotoken_aicg_blog_endutm_contentClaudeCodeAnthropicutm_campaignrewrite第六类reading choices 报错。这是解析模型返回时字段对不上通常是 Model ID 和返回结构不匹配。确认你用的 Model ID 和请求体里的model字段一致别把对话模型的 ID 填到别的接口上。再补一个 Doris 特有的坑cursor.execute执行 DDL 后立刻commit有时会报「table already exists」之类的并发问题。批量建表建议加IF NOT EXISTS或者串行执行。排查时记住一个原则先分清是连接层、SQL 层还是 API 层。2013/2006 是连接层语法错误是 SQL 层401/OAuth 是 API 层。分层之后定位快很多。6. 按场景选入口排障、验证模型与长期编码最后说清楚什么场景点哪个入口别只记首页。如果你现在卡在连接报错、401、重连不生效这类排障和接入问题上先去 API Keys 页确认 Key再看接入文档https://taotoken.net/api-keys?utm_sourcetaotoken_aicg_blog_endutm_contentapi-keysutm_campaignrewrite https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_contentdocutm_campaignrewrite如果你是想验证某个模型能不能胜任你的 SQL 生成、数据清洗任务用模型对话页直接试https://taotoken.net/models?utm_sourcetaotoken_aicg_blog_endutm_contentmodelsutm_campaignrewrite如果你是要长期跑编码类 Agent比如让模型持续帮你维护 Doris 导入脚本、做代码审查那 Coding Plan 更合适https://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_contentcoding-planutm_campaignrewrite控制台统一管理https://taotoken.net/console?utm_sourcetaotoken_aicg_blog_endutm_contentconsoleutm_campaignrewrite回到 Doris 本身给你一个收尾的实用技巧把DorisClient做成单例或连接池别每次查询都create_conn。批量任务里一个进程复用一个 client配合_ensure_alive探活既省建连开销又能扛住偶发断线。重试次数别设太大3 次足够退避用 1.5 的幂次避免雪崩。最后所有写入操作都包在execute_many里让事务边界清晰重跑时用业务主键去重这样即使断线重试也不会写脏数据。