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

Python读取多个REST API流端点的最佳实现方法咨询

多流式REST API端点的读取方案

针对你的需求,这里提供三种主流实现方案,对应不同场景,你可以根据实际情况选择:

1. 异步单进程(推荐)

这是IO密集型流式任务的最优方案,用aiohttp实现异步请求,在等待数据返回时切换到其他任务,资源占用极低,并发效率最高。

import aiohttp
import asyncio

async def stream_url(session, url, headers):
    async with session.get(url, headers=headers) as resp:
        # 异步迭代流式内容
        async for line in resp.content:
            if line:
                print(line.decode('utf-8'))  # 转字符串输出

async def main(urls, headers):
    # 创建全局会话复用连接
    async with aiohttp.ClientSession() as session:
        # 生成所有异步任务
        tasks = [stream_url(session, url, headers) for url in urls]
        # 并发执行所有任务
        await asyncio.gather(*tasks)

if __name__ == "__main__":
    # 替换为你的目标URL列表
    target_urls = ["https://api.example.com/stream1", "https://api.example.com/stream2"]
    request_headers = {"Authorization": "Bearer your_token"}
    asyncio.run(main(target_urls, request_headers))

2. 多线程(最小代码改动)

如果不想重构原有代码,多线程是快速扩展的选择。利用ThreadPoolExecutor,每个线程处理一个URL的流式数据,IO等待时GIL会释放,能达到不错的并发效果。

import requests
from concurrent.futures import ThreadPoolExecutor

def stream_single_url(url, headers):
    # 复用原有单URL逻辑
    s = requests.Session()
    resp = s.get(url, headers=headers, stream=True)
    for line in resp.iter_lines():
        if line:
            print(line.decode('utf-8'))

def main(urls, headers):
    # 根据URL数量设置线程数(也可固定值,比如5)
    with ThreadPoolExecutor(max_workers=len(urls)) as executor:
        # 映射URL和headers到处理函数
        executor.map(stream_single_url, urls, [headers]*len(urls))

if __name__ == "__main__":
    target_urls = ["https://api.example.com/stream1", "https://api.example.com/stream2"]
    request_headers = {"Authorization": "Bearer your_token"}
    main(target_urls, request_headers)

3. 多进程(仅适用于带CPU密集处理的场景)

如果每个流式数据需要大量CPU计算(比如复杂解析、数据转换),可以用多进程。但纯IO场景下不推荐,因为进程开销远大于线程/异步。

import requests
from concurrent.futures import ProcessPoolExecutor

def stream_single_url(url, headers):
    # 复用原有单URL逻辑
    s = requests.Session()
    resp = s.get(url, headers=headers, stream=True)
    for line in resp.iter_lines():
        if line:
            # 这里可以添加CPU密集型处理逻辑
            print(line.decode('utf-8'))

def main(urls, headers):
    with ProcessPoolExecutor(max_workers=len(urls)) as executor:
        executor.map(stream_single_url, urls, [headers]*len(urls))

if __name__ == "__main__":
    target_urls = ["https://api.example.com/stream1", "https://api.example.com/stream2"]
    request_headers = {"Authorization": "Bearer your_token"}
    main(target_urls, request_headers)

方案选择建议

  • 纯读取流式数据:优先选异步单进程,资源占用低、效率最高。
  • 不想改原有代码:选多线程,快速实现扩展。
  • 流式数据需大量CPU处理:选多进程,绕过GIL限制。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 04:36:29