如何在Python Flask应用中实现流式传输海量数据?解决API崩溃问题
Flask 海量数据流式传输改造与性能优化方案
流式传输改造方案
原代码一次性将600万条数据加载到DataFrame,会导致内存耗尽崩溃。通过分批查询+Flask流式响应,可以每次仅发送10万条记录,避免内存过载。
改造后代码(推荐高效分页方式)
import clickhouse_connect import orjson from flask import Flask, Response app = Flask(__name__) # 全局复用ClickHouse客户端,避免每次请求重建连接 client = clickhouse_connect.get_client( host='XXXX', port=XXXX, username='XXXX', password='XXXXX', database='XXXXX' ) def data_generator(): batch_size = 100000 last_id = 0 is_first_batch = True # 输出JSON数组开头 yield '[' while True: # 基于有序主键分页(假设表有自增id字段),比OFFSET性能高 query = f""" SELECT * FROM Table_Name WHERE id > {last_id} ORDER BY id LIMIT {batch_size} """ query_result = client.query(query) rows = query_result.result_rows # 没有更多数据时退出循环 if not rows: break # 将查询行转为字典并序列化 records = [dict(zip(query_result.column_names, row)) for row in rows] serialized_data = orjson.dumps(records).decode('utf-8') # 批次间添加逗号分隔(首批次不需要) if not is_first_batch: yield ',' else: is_first_batch = False yield serialized_data # 更新last_id为当前批次最后一条记录的id last_id = rows[-1][0] # 若id不是第一列,需对应调整索引 # 输出JSON数组结尾 yield ']' @app.route("/") def send_data(): # 返回流式响应,指定JSON类型 return Response(data_generator(), mimetype='application/json')
备选:OFFSET分页(仅当无有序主键时使用)
如果表没有合适的有序字段,可以用LIMIT OFFSET分页,但大数据量下性能较差:
def data_generator(): batch_size = 100000 offset = 0 is_first_batch = True yield '[' while True: query = f""" SELECT * FROM Table_Name LIMIT {batch_size} OFFSET {offset} """ query_result = client.query(query) rows = query_result.result_rows if not rows: break records = [dict(zip(query_result.column_names, row)) for row in rows] serialized_data = orjson.dumps(records).decode('utf-8') if not is_first_batch: yield ',' else: is_first_batch = False yield serialized_data offset += batch_size yield ']'
性能优化建议
- 复用数据库连接:全局初始化ClickHouse客户端,避免每次请求重复建立连接,减少TCP握手和认证开销。
- 避免全量加载DataFrame:直接迭代ClickHouse查询结果,跳过DataFrame转换,大幅降低内存占用。
- 精准查询字段:不要用
SELECT *,明确列出需要返回的字段,减少数据传输量和内存消耗。 - 启用响应压缩:使用
flask-compress中间件压缩响应内容,降低网络传输压力:from flask_compress import Compress Compress(app) - 异步视图优化并发:Flask 2.0+支持异步视图,避免阻塞主线程,提升多请求处理能力:
@app.route("/") async def send_data(): return Response(data_generator(), mimetype='application/json') - 优化ClickHouse查询:给分页、过滤用到的字段添加索引,或创建物化视图预聚合数据,加快查询速度。
- 调整批次大小:根据服务器内存和网络带宽,微调批次大小(如10万),平衡内存占用和传输效率。
内容的提问来源于stack exchange,提问作者Nason Thomas
相关产品推荐
相关产品推荐

