如何在Python Azure Event Hub的onEvent回调内停止接收器
该需求完全可以实现,以下是具体实现方案:
实现思路
同步版本的EventHubConsumerClient的receive()方法是阻塞调用,只需要在事件回调中匹配到目标条件后,主动调用消费客户端的close()方法,即可立即终止阻塞的接收逻辑。
修改后完整代码示例
import json from azure.eventhub import EventHubConsumerClient # 声明全局变量方便回调中访问消费客户端 consumer_client = None def on_event(partition_context, event): print(f"Received event from partition: {partition_context.partition_id}.") # 推荐使用body_as_str()获取事件载荷,避免携带额外元数据导致解析失败 event_json = json.loads(event.body_as_str()) # 匹配目标条件(修正原代码拼写错误 successull -> successful) if event_json.get("command_status") == "successful": print("匹配到符合条件的事件,停止接收器") # 调用close方法终止接收 consumer_client.close() return def on_partition_initialize(partition_context): print(f"Partition: {partition_context.partition_id} has been initialized.") def on_partition_close(partition_context, reason): print(f"Partition: {partition_context.partition_id} has been closed, reason for closing: {reason}.") def on_error(partition_context, error): if partition_context: print(f"An exception: {error} occurred during receiving from Partition: {partition_context.partition_id}.") else: print(f"An exception: {error} occurred during the load balance process.") # 初始化消费客户端 consumer_client = EventHubConsumerClient.from_connection_string( conn_str=self.conn_string_event_hub, consumer_group=self.consumer_group, eventhub_name=self.event_hub_name, ) with consumer_client: consumer_client.receive( on_event=on_event, on_partition_initialize=on_partition_initialize, on_partition_close=on_partition_close, on_error=on_error, # 如果需求是收到任意第一个事件就停止,可添加max_batch_size=1参数,回调中无需判断直接关即可 # max_batch_size=1 )
其他可选方案
- 如果不需要销毁整个消费客户端,仅需要停止单分区接收,可在回调中调用
partition_context.close()关闭对应分区的消费链路。 - 如果使用异步版本的
AsyncEventHubConsumerClient,逻辑完全一致,仅需要将close()改为await调用即可。
内容的提问来源于stack exchange,提问作者manuel mourato
相关产品推荐
相关产品推荐

