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

Python连接Azure Event Hub报错:垃圾回收时无法调用on_state_changed

问题描述

使用Python脚本通过WebSocket协议对接Azure Event Hub API时,脚本可正常发送数据,但运行过程中会触发Can't call on_state_changed during garbage collection, please be sure to close or use a context manager错误后终止执行,无法定位错误触发来源。

复现代码如下:

import asyncio
import json
import websockets

#for evh
from azure.eventhub.aio import EventHubProducerClient
from azure.eventhub import EventData


async def cryptocompare():

    producer = EventHubProducerClient.from_connection_string(conn_str="CONN_STR", eventhub_name="EVH_NAME")

    
    # this is where you paste your api key
    api_key = "API_KEY"
    url = "wss://streamer.cryptocompare.com/v2?api_key=" + api_key
    async with websockets.connect(url) as websocket:
        await websocket.send(json.dumps({
            "action": "SubAdd",
            "subs": ["0~Coinbase~BTC~EUR","0~Coinbase~BTC~USD","0~Coinbase~BTC~CHF"],
        }))
        while True:
            try:
                data = await websocket.recv()
            except websockets.ConnectionClosed:
                break
            try:
                #data = json.loads(data)
                event_data_batch = await producer.create_batch()

                # Add events to the batch.
                #for i in data:
                event_data_batch.add(EventData(data))
                # Send the batch of events to the event hub.
                await producer.send_batch(event_data_batch)
                print(json.dumps(data, indent=4))
            except ValueError:
                print(data)


asyncio.get_event_loop().run_until_complete(cryptocompare())
错误触发原因
  • 核心问题是初始化的EventHubProducerClient实例没有被正确关闭。Azure Event Hub异步客户端内部持有AMQP连接、会话、状态回调等资源,必须显式释放,否则当Python垃圾回收机制回收未关闭的客户端实例时,会尝试触发状态变更回调,而垃圾回收阶段不允许执行这类IO/协程相关的回调逻辑,就会抛出该错误。
  • 现有代码仅通过from_connection_string创建了producer实例,全程没有调用close()方法释放资源,也没有使用上下文管理器(async with语法)自动管理客户端生命周期。当脚本运行中出现异常退出、或者事件循环终止时,未关闭的producer被GC回收就会触发报错。
可行解决方案

方案1:使用上下文管理器自动管理客户端生命周期(推荐)

将producer的创建逻辑放到async with块中,代码退出块作用域时会自动执行资源关闭逻辑,从根源避免资源泄漏。修改后的完整代码如下:

import asyncio
import json
import websockets

from azure.eventhub.aio import EventHubProducerClient
from azure.eventhub import EventData


async def cryptocompare():
    # 使用async with上下文管理器包裹producer,自动处理关闭逻辑
    async with EventHubProducerClient.from_connection_string(
        conn_str="CONN_STR", 
        eventhub_name="EVH_NAME"
    ) as producer:
        api_key = "API_KEY"
        url = "wss://streamer.cryptocompare.com/v2?api_key=" + api_key
        async with websockets.connect(url) as websocket:
            await websocket.send(json.dumps({
                "action": "SubAdd",
                "subs": ["0~Coinbase~BTC~EUR","0~Coinbase~BTC~USD","0~Coinbase~BTC~CHF"],
            }))
            while True:
                try:
                    data = await websocket.recv()
                except websockets.ConnectionClosed:
                    break
                try:
                    # 先解析收到的字符串为JSON对象,避免后续打印时二次转义
                    data = json.loads(data)
                    event_data_batch = await producer.create_batch()
                    event_data_batch.add(EventData(data))
                    await producer.send_batch(event_data_batch)
                    print(json.dumps(data, indent=4))
                except ValueError:
                    print(data)


asyncio.get_event_loop().run_until_complete(cryptocompare())

方案2:显式调用close方法关闭客户端

如果不使用上下文管理器,可以在代码逻辑的finally块中显式调用await producer.close(),确保无论正常运行还是抛出异常,客户端资源都能被释放。核心修改逻辑示例:

async def cryptocompare():
    producer = EventHubProducerClient.from_connection_string(conn_str="CONN_STR", eventhub_name="EVH_NAME")
    try:
        # 原有WebSocket连接、数据收发逻辑全部放在try块内
        api_key = "API_KEY"
        url = "wss://streamer.cryptocompare.com/v2?api_key=" + api_key
        async with websockets.connect(url) as websocket:
            # 省略原有业务逻辑
            pass
    finally:
        # 无论是否抛出异常,都确保producer被正确关闭
        await producer.close()
补充注意事项
  • 不要在每次发送消息时都新建EventHubProducerClient实例,客户端内部实现了连接池复用,全局复用单个实例性能更高,也能避免频繁创建销毁带来的资源泄漏问题。
  • 原代码中print(json.dumps(data, indent=4))存在逻辑问题:WebSocket的recv()方法返回的是字符串类型,直接传入json.dumps会导致字符串被二次转义,需要先执行json.loads(data)解析为字典对象后再做序列化打印。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 23:57:19