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

Azure Event Hub多容器Checkpoint实现及全存储账户订阅问询

能否订阅Azure存储账户下的多容器并将Event Hub消费者 checkpoint指向整个账户?

目前需要实现与多容器交互的业务场景,事件通过这些容器分发,想确认是否支持订阅整个存储账户(多容器),并将Azure Event Hub消费者的checkpoint指向多个容器或整个存储账户。查阅文档及示例后发现仅支持单容器订阅,以下是当前可正常捕获文件元数据变更事件的单容器监听实现:

sharedCredential, errAzCred := container.NewSharedKeyCredential(ae.AccountName, ae.AccountKey)
if errAzCred != nil {
    log.WithError(errAzCred).Error("Creating sharedCredential failed")
    return errAzCred
}
checkClient, errAzClient := container.NewClientWithSharedKeyCredential(ae.ContanerUrl, sharedCredential, nil)
if errAzClient != nil {
    log.WithError(errAzClient).Error("Creating container client failed")
    return errAzClient
}

checkpointStore, errCheckPt := checkpoints.NewBlobStore(checkClient, nil)
if errCheckPt != nil {
    log.WithError(errCheckPt).Error("Creating checkpoint store failed")
    return errCheckPt
}
consumerClient, errClient := azeventhubs.NewConsumerClientFromConnectionString(ae.HubNamespaceConnectionString, ae.EventHubName, azeventhubs.DefaultConsumerGroup, nil)
if errClient != nil {
    log.WithError(errClient).Error("Creating event hub consumer failed")
    return errClient
}
defer consumerClient.Close(context.TODO())
processor, errProcessor := azeventhubs.NewProcessor(consumerClient, checkpointStore, nil)
if errProcessor != nil {
    log.WithError(errProcessor).Error("Creating event hub processor failed")
    return errProcessor
}

我尝试了以下实现,代码可运行但无法捕获文件元数据变更事件:

properties, err := consumerClient.GetEventHubProperties(context.TODO(), nil)
if err != nil {
    return err
}
var wg sync.WaitGroup
for _, partitionID := range properties.PartitionIDs {
    log.Debugf("start consuming partition %s", partitionID)
    partitionClient, errPC := consumerClient.NewPartitionClient(partitionID, nil)
    if errPC != nil {
        log.WithError(errPC).Error("couldn't create a partition client")
        return errPC
    }
    wg.Add(1)
    go func() {
        defer wg.Done()
        if errProcess := ae.ProcessEvents(partitionClient); errProcess != nil {
            log.WithError(errProcess).Error("failed to process the event")
        }
    }()
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 20:42:52