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

