基于.NET Core Web API对接第三方Azure Event Hub实现通知推送咨询
Azure Event Hub 监听与通知推送实现方案
方案选型建议
如果第三方仅提供Event Hub连接字符串,推荐使用Azure WebJob(或Azure Functions,托管更便捷)结合Event Hubs SDK实现持续监听——因为Webhook需要对方主动推送事件(需第三方支持配置Webhook转发),而通过SDK拉取事件的方式更灵活,无需依赖对方的推送配置。
具体实现步骤(.NET Core + Azure WebJob)
1. 安装依赖NuGet包
在项目中安装以下包:
Azure.Messaging.EventHubs(新版Event Hubs SDK)Azure.Messaging.EventHubs.Processor(事件处理器,用于批量处理+检查点管理)Microsoft.EntityFrameworkCore(数据库访问)FirebaseAdmin(Firebase推送)Microsoft.Extensions.Configuration(配置管理)
2. 配置项目参数
在appsettings.json中添加必要配置:
{ "EventHubConnectionString": "第三方提供的连接字符串", "EventHubName": "目标Event Hub名称", "StorageConnectionString": "你的存储账户连接字符串", "BlobContainerName": "eventhub-checkpoints", "DatabaseConnectionString": "你的数据库连接字符串", "FirebaseServiceAccountPath": "./firebase-service-account.json" }
3. 实现Event Hub事件监听
使用EventProcessorClient处理事件,自动管理检查点(避免重复消费):
using Azure.Messaging.EventHubs; using Azure.Messaging.EventHubs.Processor; using Azure.Storage.Blobs; using Microsoft.Extensions.Configuration; var config = new ConfigurationBuilder() .AddJsonFile("appsettings.json") .Build(); var storageConn = config["StorageConnectionString"]; var eventHubConn = config["EventHubConnectionString"]; var eventHubName = config["EventHubName"]; var containerName = config["BlobContainerName"]; // 初始化Blob容器用于存储检查点 var blobContainerClient = new BlobContainerClient(storageConn, containerName); await blobContainerClient.CreateIfNotExistsAsync(); // 创建事件处理器 var processor = new EventProcessorClient(blobContainerClient, "$Default", eventHubConn, eventHubName); // 绑定事件处理和错误处理逻辑 processor.ProcessEventAsync += ProcessEvent; processor.ProcessErrorAsync += ProcessError; // 启动处理 await processor.StartProcessingAsync(); Console.WriteLine("已启动事件监听,按回车停止..."); Console.ReadLine(); await processor.StopProcessingAsync(); async Task ProcessEvent(ProcessEventArgs args) { try { // 解析事件内容,筛选目标事件 var eventJson = args.Event.Body.ToString(); if (!IsTargetEvent(eventJson)) { await args.UpdateCheckpointAsync(); return; } // 查询订阅该事件的用户 var subscribedUsers = await FetchSubscribedUsers(eventJson); // 推送Firebase通知 await SendNotificationsToUsers(subscribedUsers); // 更新检查点,标记事件已处理 await args.UpdateCheckpointAsync(); } catch (Exception ex) { Console.WriteLine($"处理事件失败: {ex.Message}"); } } Task ProcessError(ProcessErrorEventArgs args) { Console.WriteLine($"监听发生错误: {args.Exception.Message}"); return Task.CompletedTask; }
4. 数据库查询逻辑
用EF Core实现订阅用户的查询:
private async Task<List<User>> FetchSubscribedUsers(string eventJson) { // 解析事件类型(根据实际事件结构调整) var eventType = JsonSerializer.Deserialize<EventModel>(eventJson).EventType; using var dbContext = new AppDbContext(config["DatabaseConnectionString"]); return await dbContext.Users .Where(u => u.SubscribedEventTypes.Contains(eventType)) .ToListAsync(); } // 示例实体类 public class User { public int Id { get; set; } public string FcmToken { get; set; } public List<string> SubscribedEventTypes { get; set; } } public class EventModel { public string EventType { get; set; } // 其他事件字段 }
5. Firebase通知推送
初始化Firebase Admin并发送通知:
private async Task SendNotificationsToUsers(List<User> users) { // 全局初始化一次即可 if (FirebaseApp.DefaultInstance == null) { var credential = GoogleCredential.FromFile(config["FirebaseServiceAccountPath"]); FirebaseApp.Create(new AppOptions { Credential = credential }); } var messaging = FirebaseMessaging.DefaultInstance; foreach (var user in users) { var message = new Message { Token = user.FcmToken, Notification = new Notification { Title = "事件触发通知", Body = $"你订阅的{eventType}事件已发生" } }; try { await messaging.SendAsync(message); } catch (FirebaseMessagingException ex) { Console.WriteLine($"推送通知给用户{user.Id}失败: {ex.Message}"); } } }
替代方案:Webhook(需第三方支持)
如果第三方可以将Event Hub事件转发到Webhook(比如通过Azure Event Grid),可以直接在.NET Core Web API中实现接收接口:
[ApiController] [Route("api/webhooks/event")] public class EventWebhookController : ControllerBase { private readonly AppDbContext _dbContext; public EventWebhookController(AppDbContext dbContext) { _dbContext = dbContext; } [HttpPost] public async Task<IActionResult> ReceiveEvent([FromBody] EventModel eventData) { if (!IsTargetEvent(eventData)) { return Ok(); } var users = await _dbContext.Users .Where(u => u.SubscribedEventTypes.Contains(eventData.EventType)) .ToListAsync(); await SendNotificationsToUsers(users); return Ok(); } // 注意:需处理Webhook验证(比如Azure Event Grid的GET验证请求) [HttpGet] public IActionResult ValidateWebhook([FromQuery] string validationCode) { // 返回验证代码完成验证 return Ok(validationCode); } }
实用资源参考
- Event Hubs事件处理器官方示例:了解检查点管理、批量处理的最佳实践
- Azure WebJobs部署指南:学习如何将WebJob部署到Azure App Service,配置自动启动
- Firebase Admin .NET SDK文档:掌握通知发送、错误处理的细节
- EF Core基础查询教程:快速实现数据库数据的筛选与获取
内容的提问来源于stack exchange,提问作者Harsh Joshi
相关产品推荐
相关产品推荐

