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

Python 3.6中如何用KeyboardInterrupt终止aiohttp ClientSession事件流

问题原因分析
  1. 重复输出问题:复用同一个事件循环但未彻底清理残留任务。首次中断后,SSE监听的异步任务仍在后台运行,再次调用函数时旧任务的数据流会与新任务叠加,导致重复输出。
  2. KeyboardInterrupt捕获失败:该异常在主线程触发,异步函数内部的try/except无法捕获,必须在启动事件循环的外层代码中处理。
  3. 关闭循环报错:事件循环一旦调用loop.close()就彻底失效,无法复用。3.6中直接关闭循环会导致后续调用报错。
修复方案与代码示例

以下是适配Python3.6的完整修复代码,核心是正确处理任务取消、循环清理与异常捕获:

import asyncio
import aiohttp  # 替换为你实际使用的SSE客户端库

async def _sse_listener(url, stop_trigger, result_buffer):
    """内部SSE监听逻辑:处理数据流,直到停止信号触发或满足自定义终止条件"""
    async with aiohttp.ClientSession() as session:
        async with session.get(url) as resp:
            async for raw_line in resp.content:
                if stop_trigger.is_set():
                    break
                # 替换为你的数据处理逻辑
                processed = raw_line.decode('utf-8').strip()
                if processed:
                    result_buffer.append(processed)
                    # 可添加自定义终止条件,比如:
                    # if processed == "TERMINATE_SIGNAL":
                    #     stop_trigger.set()
                    #     break

def listen_sse_until_interrupt(url):
    result = []
    loop = asyncio.get_event_loop()
    stop_event = asyncio.Event()

    # 创建监听任务
    listen_task = loop.create_task(_sse_listener(url, stop_event, result))

    try:
        loop.run_until_complete(listen_task)
    except KeyboardInterrupt:
        # 触发停止信号,取消任务并等待清理完成
        stop_event.set()
        listen_task.cancel()
        # 忽略取消异常,确保任务彻底结束
        loop.run_until_complete(asyncio.gather(listen_task, return_exceptions=True))
        return []
    else:
        # 正常终止(满足自定义条件),返回处理后的数据
        return result
    finally:
        # 清理循环中所有残留的未完成任务
        pending_tasks = asyncio.Task.all_tasks(loop=loop)
        if pending_tasks:
            loop.run_until_complete(asyncio.gather(*pending_tasks, return_exceptions=True))

# 测试调用
if __name__ == "__main__":
    first_run = listen_sse_until_interrupt("http://your-sse-endpoint.com/stream")
    print(f"首次调用结果: {first_run}")
    
    # 再次调用不会出现重复输出
    second_run = listen_sse_until_interrupt("http://your-sse-endpoint.com/stream")
    print(f"二次调用结果: {second_run}")
关键修复细节
  • 任务生命周期管理:收到中断时,先通过stop_event通知监听函数主动关闭SSE连接,再取消任务并等待其彻底结束,避免资源泄漏。
  • 循环复用与清理:在finally块中清理所有未完成的任务,确保循环状态干净,可安全复用。
  • 异常捕获位置:将KeyboardInterrupt的捕获放在loop.run_until_complete()外层,确保能捕获主线程的中断信号。
  • 数据隔离:每次调用函数时重新创建结果列表,避免多次调用间的数据污染。

如果你的SSE客户端有特定的关闭方法,需确保在stop_event触发时调用对应方法,彻底关闭连接。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 07:20:01