请求提供.NET程序读取Azure Data Explorer事件中心表数据的简单代码示例
.NET读取Azure Data Explorer事件中心数据的代码示例
所需NuGet包
安装这些包到你的.NET项目:
Install-Package Azure.Messaging.EventHubs Install-Package Azure.Messaging.EventHubs.Consumer Install-Package Azure.Storage.Blobs # 可选,用于持久化消费检查点
完整C#示例代码
using Azure.Messaging.EventHubs; using Azure.Messaging.EventHubs.Consumer; using Azure.Storage.Blobs; using System.Text; // 替换为你的实际配置参数 var eventHubConnString = "你的事件中心连接字符串"; var eventHubName = "你的事件中心名称"; var consumerGroup = EventHubConsumerClient.DefaultConsumerGroupName; var blobConnString = "你的Blob存储连接字符串"; // 用于检查点持久化,可选 var checkpointContainer = "检查点容器名称"; // 初始化检查点存储(可选) var blobContainerClient = new BlobContainerClient(blobConnString, checkpointContainer); await blobContainerClient.CreateIfNotExistsAsync(); // 创建事件中心消费者客户端 var consumerClient = new EventHubConsumerClient( consumerGroup, eventHubConnString, eventHubName, new BlobCheckpointStore(blobContainerClient)); try { // 持续读取事件中心的消息 await foreach (var partitionEvent in consumerClient.ReadEventsAsync()) { if (partitionEvent.Data.Body.Length == 0) continue; // 解析事件内容(ADX导出到事件中心的数据多为JSON格式) var messageContent = Encoding.UTF8.GetString(partitionEvent.Data.Body.ToArray()); Console.WriteLine($"来自分区 {partitionEvent.Partition.PartitionId} 的数据: {messageContent}"); // 更新检查点,避免重启后重复消费 await consumerClient.UpdateCheckpointAsync(partitionEvent); } } catch (TaskCanceledException) { // 停止接收事件时触发,属于正常终止流程 Console.WriteLine("事件接收已停止"); } finally { await consumerClient.CloseAsync(); }
关键说明
- 若不需要持久化检查点,可移除Blob相关代码,直接初始化
EventHubConsumerClient即可,但重启后会重新消费所有未确认的事件。 - ADX发送到事件中心的数据通常为JSON格式,你可以根据实际表结构将
messageContent反序列化为对应的实体类。 - 确保你的账号或服务主体拥有事件中心的
Azure Event Hubs Data Receiver权限,以及Blob存储的读写权限(如果使用检查点)。
内容的提问来源于stack exchange,提问作者Muhammad Usama Alam
相关产品推荐
相关产品推荐

