如何在Python脚本中定时重启函数解决SSE事件流长时间运行冻结问题
SSE 流冻结定期重启解决方案
核心优化思路
原来的实现存在的核心问题是每次启动新线程前没有终止旧线程,导致实例堆积,优化方案核心保证两点:
- 同一时间仅运行一个SSE监听实例
- 到达指定运行时长后强制终止当前监听,立即重启新实例
具体实现方案
方案1:进程级隔离重启(更干净,不会残留I/O阻塞资源)
用多进程替代线程,到期直接终止进程,彻底避免资源泄漏和重复实例:
import multiprocessing import time def listen(self): print("listening to events .....") try: url = settings.EVENT_URL + "/v1/events" auth_key = settings.KEY headers = { "Authorization": "Basic " + auth_key, "Accept": "text/event-stream", } # 建议给请求加读超时,设置为比服务端心跳间隔大2倍即可,避免网络假死 response = self.with_urllib3(url, headers, timeout=300) client = sseclient.SSEClient(response) for event in client.events(): logger.info(event.data) process(event.data) finally: # 退出前主动关闭连接 if 'client' in locals(): client.close() if 'response' in locals(): response.close() def start_schedule(self, run_duration=10*60): while True: # 启动监听进程 listen_proc = multiprocessing.Process(target=self.listen, name="sse-listener") listen_proc.start() # 等待指定运行时长 listen_proc.join(run_duration) # 时长到后如果进程还在运行,强制终止 if listen_proc.is_alive(): listen_proc.terminate() listen_proc.join() # 可选:加1秒间隔避免服务端被频繁请求 time.sleep(1)
直接调用start_schedule()即可启动循环调度,不会产生重复实例。
方案2:线程版实现(资源占用更低)
如果不想用多进程,可以用事件+线程超时控制,需要改造listen函数支持终止信号:
import threading import time def listen(self, stop_event: threading.Event): print("listening to events .....") try: url = settings.EVENT_URL + "/v1/events" auth_key = settings.KEY headers = { "Authorization": "Basic " + auth_key, "Accept": "text/event-stream", } response = self.with_urllib3(url, headers, timeout=300) client = sseclient.SSEClient(response) # 把SSE迭代改成非阻塞轮询,响应停止信号 event_iter = client.events() while not stop_event.is_set(): try: event = next(event_iter) logger.info(event.data) process(event.data) except StopIteration: break finally: if 'client' in locals(): client.close() if 'response' in locals(): response.close() def start_schedule(self, run_duration=10*60): while True: stop_event = threading.Event() listen_thread = threading.Thread(target=self.listen, args=(stop_event,), name="sse-listener") listen_thread.start() # 等待指定时长 time.sleep(run_duration) # 发送停止信号 stop_event.set() listen_thread.join() time.sleep(1)
额外优化建议
- 可以在
listen函数外层加异常捕获,即使SSE运行中出现报错也会触发自动重启,不用等定时周期 - 服务端如果有心跳事件(比如间隔30秒发空注释事件),可以加空闲检测:超过2倍心跳间隔没收到事件就主动退出重启,不用等固定周期
内容的提问来源于stack exchange,提问作者Sam
相关产品推荐
相关产品推荐

