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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 08:45:00