使用Azure AD读取Event Hub事件失败,返回0事件求助
问题:Azure AD认证后无法读取Event Hub事件,但写入正常
我尝试使用Azure Active Directory(Azure AD)读取Event Hub中的事件,认证后可成功向Event Hub写入事件,但读取时程序返回0事件。我能看到Event Hub中有等待的消息,且使用共享访问策略URL可以正常读取这些消息。
程序代码
using System; using System.Configuration; using System.Linq; using System.Text; using System.Threading.Tasks; using Microsoft.Azure.EventHubs; using Microsoft.Identity.Client; namespace EventHubsSenderReceiverRbac { class claimInfo { public string providerID { get; set; } public string data1 { get; set; } public string patientId { get; set; } } class Program { static readonly string TenantId = ConfigurationManager.AppSettings["tenantId"]; static readonly string ClientId = ConfigurationManager.AppSettings["clientId"]; static readonly string EventHubNamespace = ConfigurationManager.AppSettings["eventHubNamespaceFQDN"]; static readonly string EventHubName = ConfigurationManager.AppSettings["eventHubName"]; private const int numOfEvents = 5; static async Task Main() { var ehClient = generateClient(); Console.WriteLine("Press s to send {0} events\n", numOfEvents); Console.WriteLine("Press r to receive events\n"); Console.WriteLine("Press any other key to exit\n"); var userInput = Console.ReadKey(); Console.WriteLine("\n"); switch (userInput.Key) { case ConsoleKey.S: await sendEvents(ehClient); break; case ConsoleKey.R: await receiveEvents(ehClient); break; default: Console.WriteLine("Return"); return; } } static EventHubClient generateClient() { TokenProvider tp = TokenProvider.CreateAzureActiveDirectoryTokenProvider( async (audience, authority, state) => { Console.WriteLine("Client ID: {0}", ClientId); Console.WriteLine("Aud: {0}", audience); IConfidentialClientApplication app = ConfidentialClientApplicationBuilder.Create(ClientId) .WithAuthority(authority) .WithClientSecret(ConfigurationManager.AppSettings["clientSecret"]) .Build(); Console.WriteLine("App: {0} ", app); var authResult = await app.AcquireTokenForClient(new string[] { $"{audience}.default" }).ExecuteAsync(); Console.WriteLine("Access Token: {0}", authResult.AccessToken); return authResult.AccessToken; }, $"https://login.microsoftonline.com/{TenantId}"); var ehClient = EventHubClient.CreateWithTokenProvider(new Uri($"sb://{EventHubNamespace}/"), EventHubName, tp); Console.WriteLine("EH Client: {0}", ehClient.ClientId); return ehClient; } static async Task sendEvents(EventHubClient ehClient) { var clientStub = new claimInfo(); clientStub.providerID = "provider1"; clientStub.data1 = "some data"; clientStub.patientId = "xxxyyy001"; Console.WriteLine("Sending event"); for (int i = 1; i <= numOfEvents; i++) { await ehClient.SendAsync(new EventData(Encoding.UTF8.GetBytes($"{clientStub}"))); } Console.WriteLine("Send done"); await receiveEvents(ehClient); } static async Task receiveEvents(EventHubClient ehClient) { Console.WriteLine("Fetching eventhub description to discover partitions"); var ehDesc = await ehClient.GetRuntimeInformationAsync(); Console.WriteLine("ehDesc {0}", ehDesc); Console.WriteLine($"Discovered partitions as {string.Join(", ", ehDesc.PartitionIds)}"); var receiveTasks = ehDesc.PartitionIds.Select(async partitionId => { Console.WriteLine($"Initiating receiver on partition {partitionId}"); var receiver = ehClient.CreateReceiver(PartitionReceiver.DefaultConsumerGroupName, partitionId, EventPosition.FromEnd()); while (true) { Console.WriteLine("Starting Loop"); Console.WriteLine("{0}", receiver.ClientId); Console.WriteLine("{0}", receiver.ConsumerGroupName); Console.WriteLine("{0}", receiver.EventHubClient); Console.WriteLine("Part ID {0}", receiver.PartitionId); var events = await receiver.ReceiveAsync(10, TimeSpan.FromSeconds(120)); Console.WriteLine("Event: {0}", events); if (events == null) { Console.WriteLine("Ending Loop"); break; } var eventData = events.FirstOrDefault(); Console.WriteLine("Events: {0}", events); foreach (EventData c in events) { Console.WriteLine("Data: {0}", Encoding.UTF8.GetString(c.Body.Array)); } Console.WriteLine($"Received from partition {partitionId} with message content '" + Encoding.UTF8.GetString(eventData.Body.Array) + "'"); } await receiver.CloseAsync(); }).ToList<Task>(); Console.WriteLine("Waiting for receivers to complete"); await Task.WhenAll(receiveTasks); Console.WriteLine("All receivers completed"); await ehClient.CloseAsync(); } } }
问题排查与解决方案
1. 直接原因:接收位置设置错误
代码中创建接收器时使用了EventPosition.FromEnd(),这会让接收器仅接收创建时刻之后产生的新事件,完全跳过已存在的待处理消息。这是你看不到已有事件的核心原因。
2. 权限确认(必查)
虽然你能用SAS策略读取,但仍需确保Azure AD应用拥有正确权限:
- 给ClientId对应的Azure AD服务主体分配
Azure Event Hubs Data Receiver或Azure Event Hubs Data Owner角色,作用域需覆盖目标Event Hub或其命名空间。 - 角色分配后需等待5-10分钟生效,若刚配置完请等待权限同步。
3. 代码修复
修改接收器的起始位置:
- 若要读取所有未处理的历史消息,使用
EventPosition.FromStart() - 若要实现断点续读,建议结合Azure Blob存储持久化检查点(保存偏移量)
修复后的关键代码:
// 替换原CreateReceiver行,改为从事件流起始位置读取 var receiver = ehClient.CreateReceiver(PartitionReceiver.DefaultConsumerGroupName, partitionId, EventPosition.FromStart());
4. 额外优化点
- 发送逻辑中,
$"{clientStub}"会输出对象类型名而非实际数据,建议用JSON序列化(需引入Newtonsoft.Json包):await ehClient.SendAsync(new EventData(Encoding.UTF8.GetBytes(Newtonsoft.Json.JsonConvert.SerializeObject(clientStub)))); - 接收循环中,当前
events为空时直接退出,若需要持续监听新事件,可移除break语句保持循环。
内容的提问来源于stack exchange,提问作者CJW
相关产品推荐
相关产品推荐

