You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

异步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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 04:36:10