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
相关产品推荐
相关产品推荐

