惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

D
DataBreaches.Net
L
LangChain Blog
博客园_首页
J
Java Code Geeks
博客园 - 【当耐特】
Microsoft Azure Blog
Microsoft Azure Blog
小众软件
小众软件
WordPress大学
WordPress大学
V
Visual Studio Blog
T
The Blog of Author Tim Ferriss
U
Unit 42
酷 壳 – CoolShell
酷 壳 – CoolShell
Recent Announcements
Recent Announcements
C
Check Point Blog
IT之家
IT之家
Engineering at Meta
Engineering at Meta
N
Netflix TechBlog - Medium
A
About on SuperTechFans
aimingoo的专栏
aimingoo的专栏
D
Docker
有赞技术团队
有赞技术团队
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
阮一峰的网络日志
阮一峰的网络日志
I
InfoQ

博客园 - lightsong

LoRA unsloth比transformer库本身的微调有什么优点? offline-llms +++ transformer + peft 微调 Train and Fine-Tune Sentence Transformers Models Symmetric vs. Asymmetric Semantic Search Hierarchical Navigable Small Worlds (HNSW) Vision Transformer + BentoML ML Serving/编排工具 Introducing Gemma 3 270M: The compact model for hyper-efficient AI Utopia -- 企业世界模型 trustgraph semantica semantica vs graphti Industrial-Strength Natural Language Processing seata reference with springboot and other valuable demo outbox pattern with springboot Saga pattern with springboot 基于 Sentence Transformers 的具体应用案例 Vault with Keycloak as workload IAM Ontology Reasoning System ADR Claude Code的hook The AI-Native SDLC playbook Introduction to Dapper Introduction to FluentValidation Introduction to AutoFixture Introduction to FluentAssertions Understanding Return Types: IEnumerable, IReadOnlyCollection, and List Introduction to Refit Introduction to Carter
Understanding Event-Driven Architecture
lightsong · 2026-08-25 · via 博客园 - lightsong

Understanding Event-Driven Architecture

https://jdaniel1987.github.io/EventDrivenArchitecture

这是一篇基于你提供的网页内容整理的博客文章。为了便于读者理解,我优化了排版结构,并保留了代码示例和关键图表描述。


🚀 深入理解事件驱动架构 (EDA):基于 .NET 与 Azure Service Bus 的实战指南

作者:Jaime Daniel Delgado Ortega
发布时间:2024年12月4日
阅读时间:约 5 分钟


事件驱动架构 (Event-Driven Architecture, EDA) 是一种软件设计模式,它使用“事件”作为核心通信手段,从而解耦生产者和消费者。

在这篇文章中,我们将探讨如何使用 Azure Service BusMassTransit 在 .NET 中实现事件驱动系统。虽然我们将重点放在 Azure Service Bus 上,但这些原则同样适用于 RabbitMQ、Kafka 或 Amazon SQS 等其他消息代理。


🔑 核心概念

在深入代码之前,我们需要理解 EDA 的四个基本组成部分:

  1. 事件 (Event):当发生感兴趣的事情(如“订单已下达”)时发出的通知。
  2. 生产者 (Producer):发布事件的实体。
  3. 消费者 (Consumer):处理事件的实体。
  4. 消息代理 (Message Broker):管理生产者和消费者之间事件传递的中间件。

架构图解
想象一个中心化的 消息代理 (Message Broker)

  • 左侧是 事件生产者 (Event Producers),它们将事件发送进代理。
  • 右侧是 事件消费者 (Event Consumers),它们从代理中接收并处理事件。
  • 生产者和消费者互不直接通信,完全通过代理进行解耦。
  • image


💡 为什么要使用事件驱动架构?

  • 解耦 (Decoupling):生产者和消费者不需要直接了解对方。
  • 可扩展性 (Scalability):消费者可以独立地、以自己的节奏处理事件。
  • 弹性 (Resilience):一个组件的故障不会直接导致其他组件崩溃。
  • 实时处理 (Real-time Processing):能够对关键业务事件做出近乎即时的响应。

🛠️ 常见应用场景

  • 电子商务:响应用户操作,更新库存、发送通知、处理支付。
  • 微服务:在分布式系统中协调独立服务。
  • 物联网 (IoT):处理来自传感器和设备的数据流。
  • 金融系统:实时处理交易和更新。

💻 在 .NET 中实现 EDA

1. 选择消息代理

在 .NET 生态中,常见的选择包括:

  • RabbitMQ:轻量级、灵活,适合消息队列。
  • Apache Kafka:高吞吐量,适合事件流。
  • Azure Service Bus:云原生企业级消息服务。

2. 环境准备

  • Azure 账户。
  • 在 Azure 门户中配置好的 Service Bus 命名空间和队列/主题。
  • .NET 6+ 项目。

安装必要的 NuGet 包:

dotnet add package MassTransit
dotnet add package MassTransit.Azure.ServiceBus.Core

3. 生产者示例:发布事件

这是发布事件的核心代码。我们配置 MassTransit 连接到 Azure Service Bus 并发布一个 OrderPlaced 事件。

using MassTransit;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;

var builder = Host.CreateDefaultBuilder(args);

builder.ConfigureServices((context, services) =>
{
    services.AddMassTransit(x =>
    {
        x.UsingAzureServiceBus((context, cfg) =>
        {
            // 替换为你的 Azure Service Bus 连接字符串
            cfg.Host("<your-azure-service-bus-connection-string>");
        });
    });

    services.AddMassTransitHostedService();
});

var host = builder.Build();
using var scope = host.Services.CreateScope();
var bus = scope.ServiceProvider.GetRequiredService<IBus>();

// 发布事件
await bus.Publish(new OrderPlaced { 
    OrderId = Guid.NewGuid(), 
    Timestamp = DateTime.UtcNow 
});

4. 消费者示例:处理事件

消费者负责监听特定的队列并处理逻辑。

public class OrderPlacedConsumer : IConsumer<OrderPlaced>
{
    public async Task Consume(ConsumeContext<OrderPlaced> context)
    {
        Console.WriteLine($"Order received: {context.Message.OrderId} at {context.Message.Timestamp}");
        // 在此处添加你的业务逻辑
    }
}

var builder = Host.CreateDefaultBuilder(args);

builder.ConfigureServices((context, services) =>
{
    services.AddMassTransit(x =>
    {
        x.AddConsumer<OrderPlacedConsumer>();

        x.UsingAzureServiceBus((context, cfg) =>
        {
            cfg.Host("<your-azure-service-bus-connection-string>");
            cfg.ReceiveEndpoint("order-queue", e =>
            {
                e.ConfigureConsumer<OrderPlacedConsumer>(context);
            });
        });
    });

    services.AddMassTransitHostedService();
});

await builder.Build().RunAsync();

5. 事件契约 (Event Contract)

生产者和消费者需要共享相同的事件定义(通常放在共享库中)。

public record OrderPlaced
{
    public Guid OrderId { get; init; }
    public DateTime Timestamp { get; init; }
}

🛡️ 处理故障与重试

为了确保系统的弹性,必须有效处理瞬时错误。

  1. 重试策略 (Retry Policies):使用如 Polly 这样的库,在判定失败前多次尝试处理消息。
  2. 死信队列 (Dead Letter Queue, DLQ):如果重试次数耗尽,消息会被路由到 DLQ。这有助于:
    • 识别问题:检查和分析有问题的消息。
    • 防止中断:隔离失败消息,防止阻塞主处理管道。
    • 增强可见性:团队可以监控 DLQ 以主动解决重复出现的问题。

使用 Polly 实现重试的代码示例:

using Azure.Messaging.ServiceBus;
using Polly;
using Polly.Retry;
using System;
using System.Threading.Tasks;

class Program
{
    private const string ConnectionString = "<Your-Service-Bus-Connection-String>";
    private const string QueueName = "example-queue";

    static async Task Main(string[] args)
    {
        var client = new ServiceBusClient(ConnectionString);
        var sender = client.CreateSender(QueueName);

        // 定义重试策略
        var retryPolicy = Policy
            .Handle<ServiceBusException>(ex => ex.IsTransient)
            .WaitAndRetryAsync(
                retryCount: 3,
                sleepDurationProvider: attempt => TimeSpan.FromSeconds(Math.Pow(2, attempt)),
                onRetry: (exception, duration, attempt, context) =>
                {
                    Console.WriteLine($"Retrying due to: {exception.Message}. Attempt: {attempt}");
                });

        var message = new ServiceBusMessage("Hello, Service Bus!");

        try
        {
            await retryPolicy.ExecuteAsync(async () =>
            {
                Console.WriteLine("Sending message...");
                await sender.SendMessageAsync(message);
                Console.WriteLine("Message sent successfully!");
            });
        }
        catch (Exception ex)
        {
            Console.WriteLine($"Failed to send message after retries. Exception: {ex.Message}");
        }
        finally
        {
            await sender.DisposeAsync();
            await client.DisposeAsync();
        }
    }
}

⚖️ 消息代理对比

不同的消息代理适用于不同的场景。以下是主要特性的对比:

特性Azure Service BusRabbitMQKafkaAmazon SQS
类型 消息队列 消息队列 事件日志 消息队列
持久性 可选
可扩展性 中等 非常高
排序 保证 可选 保证 可选
适用场景 企业级应用 轻量级系统 流分析 云原生系统

总结:Azure Service Bus 非常适合需要死信处理、会话和事务等高级功能的企业级系统。


📊 监控与可观测性

在生产环境中,你需要知道系统是否健康。可以使用以下工具来跟踪事件处理指标:

  • Prometheus & Grafana
  • Application Insights

📌 结语

事件驱动架构是一种灵活的模式,能显著增强现代应用程序的解耦性、可扩展性和弹性。通过利用 Azure Service Bus 和 MassTransit 等工具,你可以快速构建健壮的分布式系统。

根据你的项目需求,也可以考虑 RabbitMQ、Kafka 或 Amazon SQS 等其他选择。


本文基于 CC BY 4.0 许可发布。

出处:http://www.cnblogs.com/lightsong/ 本文版权归作者和博客园共有,欢迎转载,但未经作者同意必须保留此段声明,且在文章页面明显位置给出原文连接。