如何使用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
相关产品推荐
相关产品推荐

