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
相关产品推荐
相关产品推荐

