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

如何使用Python3 requests库iter_lines处理流式数据时应对服务器停滞

处理流式API服务器停滞的健壮方案

嘿,这种流式连接遇到服务器停滞的问题我太熟了!咱们先拆解下核心问题:你当前的iter_lines()在服务器长时间没输出数据时会一直阻塞,线程直接僵住,完全没法处理异常情况。下面给你几个实用的改进方案,让你的代码能平稳应对这种状况。

1. 给请求加上超时限制

首先,给requests.get加上连接超时和读取超时,这样服务器停滞时不会无限期挂起。流式场景下,读取超时可以设得短一些(比如10秒),一旦超过这个时间没收到数据,就触发异常处理。

修改后的基础版代码:

import requests
import threading
import time

def do_something_with_the_line(line):
    # 替换成你的业务逻辑
    print(f"处理数据: {line.decode('utf-8')}")

def listener():
    someurl = "你的流式API地址"
    reconnect_delay = 5  # 重连等待时间(秒)
    
    while True:
        try:
            # timeout=(连接超时, 读取超时),单位秒
            resp = requests.get(someurl, stream=True, timeout=(5, 10))
            if resp.status_code == 200:
                print("成功建立流式连接")
                for line in resp.iter_lines():
                    if line:
                        do_something_with_the_line(line)
            else:
                print(f"请求失败,状态码: {resp.status_code}")
                
        except requests.exceptions.Timeout:
            print("服务器超时未响应,准备重连")
        except requests.exceptions.ConnectionError:
            print("连接被中断,准备重连")
        except Exception as e:
            print(f"发生未知错误: {str(e)}")
        
        print(f"等待 {reconnect_delay} 秒后尝试重连...")
        time.sleep(reconnect_delay)

# 注意:你原来的代码里写的是trade_thread.start(),应该是price_thread哦
price_thread = threading.Thread(target=listener, name="StreamingThread")
price_thread.start()

2. 用底层套接字检测数据(更精细的控制)

如果需要更精准的超时控制,可以直接操作请求的底层套接字,用select模块监听数据是否到达,超时就触发处理逻辑:

import requests
import threading
import select
import time

def do_something_with_the_line(line):
    print(f"处理数据: {line.decode('utf-8')}")

def listener():
    someurl = "你的流式API地址"
    reconnect_delay = 5
    timeout = 10  # 数据接收超时时间
    
    while True:
        try:
            resp = requests.get(someurl, stream=True, timeout=(5, None))
            if resp.status_code != 200:
                print(f"请求失败,状态码: {resp.status_code}")
                continue
                
            sock = resp.raw._fp.fp.raw  # 获取底层套接字
            print("成功建立流式连接")
            
            while True:
                # 监听套接字是否有可读数据
                ready, _, _ = select.select([sock], [], [], timeout)
                if not ready:
                    print(f"{timeout}秒未收到数据,判定服务器停滞")
                    break  # 跳出循环准备重连
                
                # 读取一行数据
                line = resp.raw.readline()
                if not line:
                    print("连接被服务器关闭")
                    break
                if line.strip():
                    do_something_with_the_line(line)
                    
        except Exception as e:
            print(f"异常发生: {str(e)}")
        
        print(f"等待 {reconnect_delay} 秒后尝试重连...")
        time.sleep(reconnect_delay)

price_thread = threading.Thread(target=listener, name="StreamingThread")
price_thread.start()

这个方案的优势是能精准控制“无数据”的判定时间,而且不会依赖requests的超时机制,灵活性更高。

3. 给线程加监控心跳

如果担心线程僵死,可以额外启动一个监控线程,定期检查流式线程的状态(比如记录最后一次处理数据的时间),如果超过阈值就重启流式线程:

import requests
import threading
import time

last_process_time = time.time()
lock = threading.Lock()

def do_something_with_the_line(line):
    global last_process_time
    with lock:
        last_process_time = time.time()
    print(f"处理数据: {line.decode('utf-8')}")

def listener():
    someurl = "你的流式API地址"
    while True:
        try:
            resp = requests.get(someurl, stream=True, timeout=(5, 10))
            if resp.status_code == 200:
                for line in resp.iter_lines():
                    if line:
                        do_something_with_the_line(line)
            else:
                print(f"请求失败,状态码: {resp.status_code}")
                time.sleep(5)
        except Exception as e:
            print(f"异常发生: {str(e)}")
            time.sleep(5)

def monitor_thread():
    global last_process_time
    timeout_threshold = 15  # 15秒未处理数据判定为停滞
    while True:
        time.sleep(5)
        with lock:
            elapsed = time.time() - last_process_time
        if elapsed > timeout_threshold:
            print("检测到流式线程停滞,准备重启...")
            # 停止旧线程(这里可以根据实际情况优化,比如设置退出标志)
            global price_thread
            price_thread.join(timeout=2)
            # 启动新线程
            price_thread = threading.Thread(target=listener, name="StreamingThread")
            price_thread.start()
            with lock:
                last_process_time = time.time()

# 启动流式线程和监控线程
price_thread = threading.Thread(target=listener, name="StreamingThread")
price_thread.start()
monitor = threading.Thread(target=monitor_thread, name="MonitorThread")
monitor.start()

额外小提示

  • 你原来的代码里有个笔误:trade_thread.start()应该是price_thread.start(),记得修正哦。
  • 如果API提供商明确说了停滞的触发场景,可以针对性地做处理(比如特定时间段容易停滞,就提前做好重连准备)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:32:51