异步conn.commit阻塞事件循环致数据延迟递增,如何解决?
问题
我用两个异步库:一个接收服务器流式数据,另一个是异步版psycopg用来把数据存入TimescaleDB。但conn.commit()会阻塞主事件循环,导致接收数据和存入数据库的时间差随着迭代越来越大。想在不把数据库调用移到单独线程的前提下解决这个问题。
代码示例
async def exec_insert_trade(conn, trade, current_time): trade_values = await generate_trade_values(trade) columns = ', '.join(trade_values.keys()) placeholders = ', '.join(['%s'] * len(trade_values)) query = f"INSERT INTO trades ({columns}) VALUES ({placeholders})" time_difference = current_time - trade.time # 计算时间差 print(f"Time difference trade: {time_difference}") await conn.execute(query, list(trade_values.values())) await conn.commit() async def main(): async with AsyncRetryingClient(TINKOFF_READONLY_TOKEN, settings=retry_settings) as client, \ AsyncConnectionPool(POSTGRES_DATABASE_URL, min_size=2) as pool: async for marketdata in client.market_data_stream.market_data_stream(request_iterator()): current_time = datetime.now(timezone.utc) async with pool.connection() as conn: if marketdata.trade is not None: await exec_insert_trade(conn, marketdata.trade, current_time) if __name__ == "__main__": asyncio.run(main())
输出日志
.... Time difference trade: 0:00:00.233646 Time difference trade: 0:00:00.952377 Time difference trade: 0:00:01.187182 Time difference trade: 0:00:01.042835 Time difference trade: 0:00:03.101548 Time difference trade: 0:00:06.067422 Time difference trade: 0:00:07.025047 ...
解决方案
1. 批量提交代替单条提交
频繁单条commit是性能瓶颈核心——每次commit都会触发数据库WAL预写日志刷盘,直接阻塞事件循环。改成批量攒数据后再提交:
- 用
asyncio.Queue维护一个待插入数据的缓存队列 - 启动独立异步任务,定期(比如每1秒)或攒够指定条数(比如100条)时,批量执行插入并commit
- 主流程只负责把数据丢进队列,不等待插入完成,彻底避免阻塞流式数据接收
示例代码片段:
async def batch_insert_worker(pool, batch_size=100, interval=1): queue = asyncio.Queue(maxsize=1000) # 限制队列大小防止内存溢出 async def consumer(): async with pool.connection() as conn: # 复用一个连接减少开销 while True: batch = [] # 要么攒够batch_size条,要么到时间就执行 try: for _ in range(batch_size): batch.append(await asyncio.wait_for(queue.get(), timeout=interval)) except asyncio.TimeoutError: if not batch: continue # 批量插入逻辑 if batch: first_trade = batch[0] trade_values = await generate_trade_values(first_trade) columns = ', '.join(trade_values.keys()) placeholders = ', '.join(['%s'] * len(trade_values)) query = f"INSERT INTO trades ({columns}) VALUES ({placeholders})" # 准备所有插入参数 params = [list((await generate_trade_values(t)).values()) for t in batch] await conn.executemany(query, params) await conn.commit() # 标记队列任务完成 for _ in batch: queue.task_done() asyncio.create_task(consumer()) return queue async def main(): async with AsyncRetryingClient(TINKOFF_READONLY_TOKEN, settings=retry_settings) as client, \ AsyncConnectionPool(POSTGRES_DATABASE_URL, min_size=2) as pool: insert_queue = await batch_insert_worker(pool) async for marketdata in client.market_data_stream.market_data_stream(request_iterator()): current_time = datetime.now(timezone.utc) if marketdata.trade is not None: time_difference = current_time - marketdata.trade.time print(f"Time difference trade: {time_difference}") # 把数据放入队列,不等待插入完成 await insert_queue.put(marketdata.trade) # 程序退出前等待队列所有任务处理完 await insert_queue.join()
2. 调整PostgreSQL WAL刷盘策略(谨慎操作)
如果批量提交后仍有阻塞,可以调整PostgreSQL的wal_writer_delay参数(默认200ms),延长WAL刷盘间隔,减少commit的阻塞频率。注意:这会增加数据库崩溃时的数据丢失风险,仅适合对数据一致性要求不极高的场景。
修改postgresql.conf:
wal_writer_delay = 1000ms # 改为1秒刷盘一次
3. 利用psycopg的autocommit特性(配合批量)
把数据库连接设置为autocommit=True,这样批量插入后会自动提交,省去手动调用commit()的步骤,但必须配合批量操作才能有效降低阻塞:
在consumer的连接初始化时添加:
async with pool.connection() as conn: conn.autocommit = True # 开启自动提交 # ... 后续批量插入逻辑不变 ...
内容的提问来源于stack exchange,提问作者Alexander Zar
相关产品推荐
相关产品推荐

