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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 23:42:50