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

基于.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 06:14:54