【问题标题】:Azure eventhub library for python用于 python 的 Azure eventthub 库
【发布时间】:2020-12-09 08:35:56
【问题描述】:

我正在使用 eventthub 来摄取大量事件。我有多个消费者正在运行一个扩展组,从具有多个分区的 eventthub 读取这些事件。我正在使用 python 中的 Azure SDK 并且对使用什么感到困惑。有eventthubconsumerclient、eventprocessorHost ....

我想使用一个库,我的多个消费者可以使用消费者组连接,分区动态分配给每个消费者,并在存储帐户中进行检查点,就像我使用 kafka 的方式一样。

【问题讨论】:

  • 我使用了示例代码,这是我得到的错误“在平衡和声明 eventthub 所有权期间发生异常 (KeyError('offset'))”
  • 你考虑在python中使用“事件处理器主机”吗?它使用消费者组并且可以设置检查点。
  • 在以后的python sdk版本中,可以看到使用了eventthubconsumerclient。请参阅评论中的链接。也看看错误。我正在尝试运行链接中提供的相同示例程序
  • 这是一个预发布版本,不是稳定版本。不确定它是否有一些潜在的错误。

标签: python azure azure-eventhub


【解决方案1】:

更新:

对于生产使用,我建议您应该使用稳定版的事件中心 sdk。可以使用eph,示例代码为here


我可以使用pre-release eventhub 5.0.0b6 来使用消费者组以及设置检查点。

但奇怪的是,在 blob 存储中,我可以看到为 eventthub 创建的 2 个文件夹:checkpointownership 文件夹。在文件夹内,为分区创建了 blob,但 blob 为空。更奇怪的是,即使blob是空的,每次从eventthub读取时,它总是读取最新的数据(意味着它从不读取已经在同一个消费者组中读取的数据)。

您需要安装azure-eventhub 5.0.0b6 并使用pip install --pre azure-eventhub-checkpointstoreblob 安装azure-eventhub-checkpointstoreblob。对于 blob 存储,您应该安装最新的 version 12.1.0 of azure-storage-blob

我关注这个sample。在此示例中,它使用 事件中心级别 连接字符串(不是 事件中心命名空间级别 连接字符串)。您需要通过导航创建一个 事件中心级别 连接字符串到 azure 门户 -> 你的 eventhub 命名空间 -> 你的事件中心实例 -> 共享访问策略 -> 单击“添加” -> 然后指定一个策略名称,然后选择权限。如果只想接收数据,只能选择监听权限。截图如下:

创建策略后,您可以按照以下屏幕截图复制连接字符串:

然后你可以按照下面的代码:

import os
from azure.eventhub import EventHubConsumerClient
from azure.eventhub.extensions.checkpointstoreblob import BlobCheckpointStore

CONNECTION_STR = 'Endpoint=sb://ivanehubns.servicebus.windows.net/;SharedAccessKeyName=saspolicy;SharedAccessKey=xxx;EntityPath=myeventhub'
STORAGE_CONNECTION_STR = 'DefaultEndpointsProtocol=https;AccountName=xx;AccountKey=xxx;EndpointSuffix=core.windows.net'


def on_event(partition_context, event):
    # do something with event
    print(event)
    print('on event')
    partition_context.update_checkpoint(event)


if __name__ == '__main__':

    #the "a22" is the blob container name
    checkpoint_store = BlobCheckpointStore.from_connection_string(STORAGE_CONNECTION_STR, "a22")

    #the "$default" is the consumer group
    client = EventHubConsumerClient.from_connection_string(
        CONNECTION_STR, "$default", checkpoint_store=checkpoint_store)

    try:
        print('ok')
        client.receive(on_event)
    except KeyboardInterrupt:
        client.close()

测试结果:

【讨论】:

  • 谢谢,我也试过这个,但我一直收到这个错误:在命名空间'xxxxxxxx' eventthub'xxxxxxx'消费者组'$default'的list_ownership期间发生异常。例外是 KeyError('ownerid')
  • @Nipun,你能按照我的步骤,尝试使用新的 blob 容器吗?
  • 如果错误仍然存​​在,请发布您的代码,以及您正在使用的所有包(包括版本)。所以我可以调试它找到原因:)。
  • 在浏览了一些文档后,我看到 partitionmanager 需要拥有分区的所有权,我猜它正试图从空的存储 blob 中获取它
  • @Nipun,也许吧。但不确定为什么 blob 是空的,但可以正常工作。而且由于它是预发布版本,可能会在几天内修复。
【解决方案2】:

azure-eventhub v5 已于 2020 年 1 月 GAed,最新版本为 v5.2.0

在 pypi 上可用:https://pypi.org/project/azure-eventhub/

请关注migration guide from v1 to v5 迁移您的程序。

使用checkpoint收货请关注sample code:

import os
import logging
from azure.eventhub import EventHubConsumerClient
from azure.eventhub.extensions.checkpointstoreblob import BlobCheckpointStore

CONNECTION_STR = os.environ["EVENT_HUB_CONN_STR"]
EVENTHUB_NAME = os.environ['EVENT_HUB_NAME']
STORAGE_CONNECTION_STR = os.environ["AZURE_STORAGE_CONN_STR"]
BLOB_CONTAINER_NAME = "your-blob-container-name"  # Please make sure the blob container resource exists.

logging.basicConfig(level=logging.INFO)
log = logging.getLogger(__name__)


def on_event_batch(partition_context, event_batch):
    log.info("Partition {}, Received count: {}".format(partition_context.partition_id, len(event_batch)))
    # put your code here
    partition_context.update_checkpoint()


def receive_batch():
    checkpoint_store = BlobCheckpointStore.from_connection_string(STORAGE_CONNECTION_STR, BLOB_CONTAINER_NAME)
    client = EventHubConsumerClient.from_connection_string(
        CONNECTION_STR,
        consumer_group="$Default",
        eventhub_name=EVENTHUB_NAME,
        checkpoint_store=checkpoint_store,
    )
    with client:
        client.receive_batch(
            on_event_batch=on_event_batch,
            max_batch_size=100,
            starting_position="-1",  # "-1" is from the beginning of the partition.
        )


if __name__ == '__main__':
    receive_batch()

还有一点值得注意的是,在 V5 中,我们使用 blob 的元数据来存储检查点和所有权信息,而不是在 v1 中将它们存储为 blob 的内容。所以在使用 v5 sdk 时,预计 blob 的内容为空。

【讨论】:

  • 你的意思是azure.eventhub.extensions.checkpointstoreblobaioextensions.checkpointstoreblob 似乎不是一个东西。
  • azure.eventhub.extensions.checkpointstoreblob 是同步版本 -- pypi.org/project/azure-eventhub-checkpointstoreblobazure.eventhub.extensions.checkpointstoreblobaio 是异步版本 -- pypi.org/project/azure-eventhub-checkpointstoreblob-aio
  • 好的,谢谢。对于其他提出这个问题的人: async 和 sync 包是分开安装的。 azure-eventhub-checkpointstoreblob 和 pip 安装 azure-eventhub-checkpointstoreblob-aio
猜你喜欢
  • 2020-01-20
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-09-08
相关资源
最近更新 更多