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

Python3 长时数据库循环进程嵌入HTTP服务器的进度统计问题

这个场景我之前也碰到过,十亿次迭代的循环确实容易把Tornado的IOLoop堵死,导致HTTP请求完全得不到响应。给你几个实用的解决方案,你可以根据现有代码的情况来选:

方案一:把数据库循环放到独立线程(最省心的改造方式)

如果你的数据库处理逻辑是同步的,不想改太多代码,直接把循环扔到单独线程里跑,主线程专心处理Tornado的HTTP请求就行。记得用锁保护进度变量,避免线程安全问题,另外一定要用服务器端游标,不然十亿行数据直接加载到内存会直接炸。

示例代码:

import threading
import tornado.web
import tornado.ioloop
import psycopg2  # 换成你用的数据库驱动

# 进度统计变量,加锁保证线程安全
progress_lock = threading.Lock()
total_rows = 0
processed_rows = 0

class ProgressHandler(tornado.web.RequestHandler):
    def get(self):
        with progress_lock:
            self.write({
                "total": total_rows,
                "processed": processed_rows,
                "progress": round((processed_rows / total_rows) * 100, 2) if total_rows > 0 else 0
            })

def run_db_query():
    global total_rows, processed_rows
    # 连接数据库,用服务器端游标(name参数指定)
    conn = psycopg2.connect("dbname=your_db user=your_user password=your_pwd host=your_host")
    cur = conn.cursor(name='large_result_cursor')
    cur.execute("SELECT * FROM your_huge_table")
    
    # 先获取总行数(不同数据库获取方式可能不同,比如MySQL需要单独查COUNT)
    with progress_lock:
        total_rows = cur.rowcount
    
    # 迭代处理每一行
    for row in cur:
        # 这里替换成你的行处理逻辑
        process_single_row(row)
        
        # 更新进度,记得加锁
        with progress_lock:
            processed_rows += 1
    
    cur.close()
    conn.close()

def process_single_row(row):
    # 示例:简单打印或者复杂业务逻辑
    pass

if __name__ == "__main__":
    app = tornado.web.Application([(r"/progress", ProgressHandler)])
    app.listen(8888)
    
    # 启动数据库处理线程
    db_thread = threading.Thread(target=run_db_query, daemon=True)
    db_thread.start()
    
    # 启动Tornado的IOLoop
    tornado.ioloop.IOLoop.current().start()

这个方案的好处是几乎不用改原有的数据库处理代码,只要把循环包到线程里就行,HTTP请求能实时响应。

方案二:用异步数据库驱动(最符合Tornado设计的方式)

如果你的数据库有异步驱动(比如PostgreSQL的asyncpg、MySQL的aiomysql),直接改成异步流程,这样整个程序都在单线程的IOLoop里跑,数据库查询和HTTP请求能同时处理,完全不用手动让出控制权。

示例代码:

import tornado.web
import tornado.ioloop
import asyncpg

# 异步场景下不用锁,因为单线程执行,没有竞态问题
total_rows = 0
processed_rows = 0

class ProgressHandler(tornado.web.RequestHandler):
    async def get(self):
        self.write({
            "total": total_rows,
            "processed": processed_rows,
            "progress": round((processed_rows / total_rows) * 100, 2) if total_rows > 0 else 0
        })

async def run_db_query():
    global total_rows, processed_rows
    # 异步连接数据库
    conn = await asyncpg.connect(user='your_user', database='your_db', password='your_pwd', host='your_host')
    
    # 先获取总行数
    total_rows = await conn.fetchval("SELECT COUNT(*) FROM your_huge_table")
    
    # 异步迭代结果集
    async with conn.transaction():
        async for record in conn.cursor("SELECT * FROM your_huge_table"):
            # 如果你的处理逻辑是CPU密集型,记得扔到线程池里,别阻塞IOLoop
            await tornado.ioloop.IOLoop.current().run_in_executor(None, process_single_row, record)
            processed_rows += 1
    
    await conn.close()

def process_single_row(record):
    # 你的行处理逻辑,同步/异步都可以,CPU密集型用线程池
    pass

if __name__ == "__main__":
    app = tornado.web.Application([(r"/progress", ProgressHandler)])
    app.listen(8888)
    
    # 把异步数据库任务加到IOLoop
    tornado.ioloop.IOLoop.current().add_callback(run_db_query)
    
    tornado.ioloop.IOLoop.current().start()

这个方案性能最优,完全贴合Tornado的异步模型,不过需要你替换成异步数据库驱动,可能要改一些原有的查询代码。

方案三:定期主动让出IOLoop(妥协方案)

如果你既不想开线程,也不想换异步驱动,那只能在循环里定期给IOLoop留时间处理HTTP请求。比如每处理N行,就让IOLoop跑一次pending的回调。

示例代码:

import tornado.web
import tornado.ioloop
import psycopg2

processed_rows = 0
total_rows = 0
# 每处理1000行就让出一次控制权,根据你的处理速度调整这个值
BATCH_SIZE = 1000

class ProgressHandler(tornado.web.RequestHandler):
    def get(self):
        self.write({
            "total": total_rows,
            "processed": processed_rows,
            "progress": round((processed_rows / total_rows) * 100, 2) if total_rows > 0 else 0
        })

def run_db_query():
    global processed_rows, total_rows
    conn = psycopg2.connect("dbname=your_db user=your_user")
    cur = conn.cursor(name='large_result_cursor')
    cur.execute("SELECT * FROM your_huge_table")
    total_rows = cur.rowcount
    
    count = 0
    for row in cur:
        process_single_row(row)
        processed_rows += 1
        count += 1
        
        # 每处理BATCH_SIZE行,让IOLoop处理一次请求
        if count % BATCH_SIZE == 0:
            tornado.ioloop.IOLoop.current().run_sync(lambda: None)
    
    cur.close()
    conn.close()

def process_single_row(row):
    pass

if __name__ == "__main__":
    app = tornado.web.Application([(r"/progress", ProgressHandler)])
    app.listen(8888)
    
    tornado.ioloop.IOLoop.current().add_callback(run_db_query)
    tornado.ioloop.IOLoop.current().start()

这个方案的缺点是HTTP请求的响应延迟取决于BATCH_SIZE的大小:太大的话请求会卡顿,太小的话会拖慢循环速度,需要你根据实际情况调试。


内容的提问来源于stack exchange,提问作者Martin Witkowski

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 09:29:29