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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 22:15:37