Skip to content

事件总线

HMX平台内置了事件总线机制,支持进程内和跨服务的事件发布/订阅,以及命令的发送/处理。事件总线是实现服务解耦、异步通信的核心组件。

架构概览

mermaid
graph TB
    subgraph EventBus["事件总线 (IEventBusPlus)"]
        Event["Event(广播)<br/>所有订阅者收到"]
        Command["Command(命令)<br/>只有一个消费者"]
    end
    
    subgraph RabbitMQ["RabbitMQ 实现"]
        Fanout["Fanout Exchange<br/>hmx.events.fanout"]
        StreamQueue["Stream Queue<br/>hmx.events.stream.shared"]
        QuorumQueue["Quorum Queue<br/>cmd.类型名"]
    end
    
    subgraph Redis["Redis 实现"]
        PubSub["Pub/Sub"]
        RedisStream["Stream + Consumer Group"]
    end
    
    Event --> Fanout
    Event --> PubSub
    Command --> QuorumQueue
    Command --> RedisStream
    
    Fanout --> StreamQueue

消息模式

Event(广播事件)

特性说明
模式广播,所有订阅者都能收到
RabbitMQFanout Exchange + Stream Queue
RedisPub/Sub
适用场景通知推送、状态变更广播

Command(命令)

特性说明
模式竞争消费,只有一个消费者处理
RabbitMQQuorum Queue
RedisStream + Consumer Group
适用场景任务分发、后台处理

配置

连接配置

appsettings.json 中配置:

json
{
  "ConnectionStrings": {
    "redis": "10.11.5.49:6379,password=redis123!",
    "rabbitmq": "amqp://user:password@10.11.5.57:5672/dev"
  }
}

降级策略

配置情况使用的实现
配置了 rabbitmqRabbitMqEventBusPlus
未配置 rabbitmq,配置了 redisRedisEventBusPlus
都未配置抛出异常

注意:SignalR需要Redis支持,因此至少需要配置 redis

启用事件总线

Program.cs 中:

csharp
var builder = WebApplication.CreateBuilder(args);

// 添加事件总线服务
builder.Services.AddEventBusPlus(builder.Configuration);

var app = builder.Build();

使用方式

定义事件

继承 EventDataBase 定义事件类型:

csharp
using Hmx.Core.Data;
using Hmx.Core.EventBus;

/// <summary>
/// 订单创建事件
/// </summary>
public class OrderCreatedEvent : EventDataBase
{
    /// <summary>
    /// 订单ID
    /// </summary>
    public string OrderId { get; set; }

    /// <summary>
    /// 订单金额
    /// </summary>
    public decimal Amount { get; set; }
}

发布事件

csharp
public class OrderAppService : HmxAppServiceBase<OrderAppService>, IOrderAppService
{
    private readonly IEventBusPlus _eventBus;

    public OrderAppService(IEventBusPlus eventBus)
    {
        _eventBus = eventBus;
    }

    public async Task CreateOrderAsync(OrderInput input)
    {
        // 业务逻辑...
        
        // 发布事件(广播给所有订阅者)
        await _eventBus.PublishEventAsync(new OrderCreatedEvent
        {
            OrderId = order.Id,
            Amount = order.Amount
        });
    }
}

订阅事件

csharp
public class OrderEventHandler : IScopedDependency
{
    private readonly IEventBusPlus _eventBus;

    public OrderEventHandler(IEventBusPlus eventBus)
    {
        _eventBus = eventBus;
    }

    public void Initialize()
    {
        // 订阅订单创建事件
        _eventBus.SubscribeEvent<OrderCreatedEvent>(async evt =>
        {
            // 处理事件:如发送通知、更新统计等
            Console.WriteLine($"新订单: {evt.OrderId}, 金额: {evt.Amount}");
        });
    }
}

发布命令

csharp
// 发布命令(只有一个消费者处理)
await _eventBus.PublishCommandAsync(new ProcessOrderCommand
{
    OrderId = order.Id
});

订阅命令

csharp
_eventBus.SubscribeCommand<ProcessOrderCommand>(async cmd =>
{
    // 处理命令(耗时操作)
    await ProcessOrder(cmd.OrderId);
});

已有事件示例

事件类用途
AlertEventData内部消息通知
BroadcastEventData广播消息
AlertClickEventData消息点击事件
UserAlreadyLoginedEventData用户重复登录事件
SystemNotifyCreatedEventData系统通知创建事件

SignalR集成

事件发布时会同时推送给SignalR客户端:

csharp
// 内部实现
public async Task PublishEventAsync<T>(T eventData) where T : IEventData
{
    // 发布到消息队列
    channel.BasicPublish(EVENT_EXCHANGE, ...);
    
    // 同时推送给SignalR客户端
    await _signalRhubs.Clients.All.SendAsync(typeof(T).AssemblyQualifiedName!, eventData);
}

客户端订阅:

csharp
// WinForm客户端
var client = new HmxSignalRClient("http://localhost:5225");
await client.On<OrderCreatedEvent>(evt =>
{
    MessageBox.Show($"新订单: {evt.OrderId}");
});

案例:系统客户端实时通知

场景:后端发布系统通知,WinForm客户端实时接收并显示角标和Toast提示。

1. 定义事件(后端)

csharp
// Hmx.Service.Admin.Events/SystemNotifyCreatedEventData.cs
public class SystemNotifyCreatedEventData : EventDataBase
{
    public string Title { get; set; }
    public string TargetType { get; set; }  // 0=全员, 1=定向
    public List<string> TargetUserIds { get; set; }
}

2. 发布事件(后端)

csharp
// NotificationAppService.cs
public async Task PublishNotificationAsync(HmxNotify notify)
{
    // 保存通知到数据库...
    
    // 发布事件(同时推送到RabbitMQ和SignalR)
    await _eventBus.PublishEventAsync(new SystemNotifyCreatedEventData
    {
        Title = notify.CTitle,
        TargetType = notify.CTargetType.ToString(),
        TargetUserIds = notify.CTargetUserList?.Split(';').ToList()
    });
}

3. 订阅事件(WinForm客户端)

csharp
// NotifySignalRHandler.cs
public static class NotifySignalRHandler
{
    public static event Action<int>? OnNewNotify;

    public static void Subscribe(WinFormSignalRSubscriber subscriber)
    {
        // 订阅系统通知事件
        _ = subscriber.On<SystemNotifyCreatedEventData>(HandleNotifyEvent);
    }

    private static async Task HandleNotifyEvent(SystemNotifyCreatedEventData evt)
    {
        var userId = EngineContext.Current.User?.UserId;
        
        // 检查是否是发给自己的通知
        if (evt.TargetType != "0" && evt.TargetUserIds != null
            && !evt.TargetUserIds.Contains(userId))
        {
            return; // 不是发给我的,忽略
        }

        // 刷新未读数角标
        var count = await SF.Proxy<INotificationAppService>()
            .GetUnreadCountAsync(userId);
        OnNewNotify?.Invoke(count);

        // 右上角Toast提示
        MsgBox.ShowAlert($"系统通知:{evt.Title}");
    }
}

4. UI线程安全处理

csharp
// WinFormSignalRSubscriber.cs
public async Task On<T>(Func<T, Task> handler) where T : IEventData
{
    await SignalRClientManager.Instance.Client.On<T>(async data =>
    {
        if (_control.InvokeRequired)
        {
            // 跨线程调用UI
            _control.BeginInvoke(new Action(async () => await handler(data)));
        }
        else
        {
            await handler(data);
        }
    });
}

流程图:

后端发布通知 ──► EventBus ──► RabbitMQ (持久化)

                    └──► SignalR Hub ──► WinForm客户端

                                              ├──► 检查是否发给自己
                                              ├──► 刷新角标 (OnNewNotify)
                                              └──► Toast提示

5. 主窗体初始化

csharp
// HmxMainLayoutForm.cs
public partial class HmxMainLayoutForm : HmxBaseForm
{
    public HmxMainLayoutForm()
    {
        InitializeComponent();
        
        // 一行初始化通知订阅
        NotifyBadge.Initialize(barbtnNotification, this);
    }
}

NotifyBadge.Initialize 内部会创建 WinFormSignalRSubscriber 并调用 NotifySignalRHandler.Subscribe,同时设置定时轮询作为SignalR断连时的降级兜底。

RabbitMQ集成

更多信息请参考 RabbitMQ官方教程

平台内置队列

队列/交换机类型用途
hmx.events.fanoutFanout Exchange广播事件交换机
hmx.events.stream.sharedStream Queue广播事件队列(共享)
cmd.{完整类名}Quorum Queue命令队列(每个命令类型一个)

案例一:实验室数据采集

场景:实验室试验机通过Access数据库存储检测数据,系统实时监控文件变化,将数据同步到RabbitMQ,供MES主系统消费。

流程图:

试验机检测 → Access.mdb → FileSystemWatcher → RabbitMQPublisher → RabbitMQ → MES主系统消费

1. 队列与消息规划

规划项说明
Exchangelims.topicTopic Exchange,支持路由键匹配
Queuelims.testdata试验机检测数据队列
RoutingKeymechanical.*路由键,支持通配符匹配
消息TTL72小时超时自动删除
消息格式JSON序列化后的检测数据

消息体定义:

csharp
public class TestDataEvent
{
    public string MachineId { get; set; }    // 试验机编号
    public string SampleId { get; set; }     // 试样编号
    public decimal TestValue { get; set; }   // 检测值
    public string Unit { get; set; }         // 单位
    public DateTime Timestamp { get; set; }  // 检测时间
}

2. 发送方示例

csharp
public class MechanicalDbSyncHostedService : IHostedService
{
    private RabbitMQPublisher? _rabbit;

    public async Task StartAsync(CancellationToken ct)
    {
        // 初始化RabbitMQ发布者
        _rabbit = new RabbitMQPublisher(
            connectionString: "amqp://mes:password@10.11.5.57:5672/dev",
            queueName: "lims.testdata",
            routingKey: "mechanical.*",
            messageTtl: TimeSpan.FromHours(72));

        // 启动文件监控
        SetupFileWatcher(opt.AccessFile);
    }

    private async Task ProcessFileAsync()
    {
        // 从Access数据库读取数据
        var rows = await _db.QueryListAsync("SELECT * FROM TestData WHERE Synced = 0");
        
        if (rows.Count > 0)
        {
            // 批量发布到RabbitMQ
            _rabbit.BatchPublish(rows.Select(r => new TestDataEvent
            {
                MachineId = r["MachineId"].ToString(),
                SampleId = r["SampleId"].ToString(),
                TestValue = Convert.ToDecimal(r["TestValue"]),
                Timestamp = DateTime.Now
            }));
        }
    }
}

RabbitMQPublisher核心特性:

特性配置说明
消息持久化Persistent = true服务重启后消息不丢失
自动重试3次,间隔2秒发布失败自动重试
延迟连接构造时不连接保证服务启动安全
批量发布BatchPublish<T>()支持批量发送,提高吞吐

3. 处理方示例

csharp
public class TestDataConsumer : BackgroundService
{
    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        var factory = new ConnectionFactory { Uri = new Uri("amqp://mes:password@10.11.5.57:5672/dev") };
        using var connection = factory.CreateConnection();
        using var channel = connection.CreateModel();

        // 声明队列
        channel.QueueDeclare("lims.testdata", durable: true, exclusive: false, autoDelete: false);

        var consumer = new EventingBasicConsumer(channel);
        consumer.Received += async (model, ea) =>
        {
            var json = Encoding.UTF8.GetString(ea.Body.ToArray());
            var eventData = JsonSerializer.Deserialize<TestDataEvent>(json);

            try
            {
                // 处理检测数据:写入MES数据库
                await SaveToMES(eventData);
                
                // 确认消息
                channel.BasicAck(ea.DeliveryTag, false);
            }
            catch (Exception ex)
            {
                // 处理失败,重新入队
                channel.BasicNack(ea.DeliveryTag, false, requeue: true);
            }
        };

        channel.BasicConsume("lims.testdata", autoAck: false, consumer: consumer);
    }
}

4. 场景总结

要点说明
数据源Access数据库文件(试验机本地存储)
触发方式FileSystemWatcher监控文件变化
Exchange类型Topic Exchange,支持路由键通配符
消息特点批量、异步、持久化
容错机制自动重试 + 消息TTL
适用场景异步数据同步、批量数据导入

案例二:自动化设备对接

场景:BXCOM作为设备通讯网关,通过TCP/IP与PLC、仪表等设备通讯,采集的实时数据转发到RabbitMQ,供MES/SCADA系统消费。

流程图:

PLC/仪表 → TCP Socket采集 → BXCOM网关解析 → RabbitMQ → MES/SCADA系统

1. 队列与消息规划

规划项说明
Exchangedevice.directDirect Exchange,精确路由
Queuedevice.{设备编号}每台设备一个队列
RoutingKeydevice.{设备编号}与队列名一致
消息TTL无(实时消费)数据实时处理
消息格式JSON统一格式便于解析

消息体定义:

csharp
public class DeviceDataEvent
{
    public string DeviceId { get; set; }      // 设备编号
    public string DeviceType { get; set; }    // 设备类型(PLC/仪表)
    public string PointAddress { get; set; }  // 数据点地址
    public object Value { get; set; }         // 数据值
    public DateTime CollectTime { get; set; } // 采集时间
    public string Quality { get; set; }       // 数据质量(Good/Bad)
}

2. 发送方示例

csharp
public class BxcomDataForwarder
{
    private IConnection? _connection;
    private IModel? _channel;
    private const string EXCHANGE_NAME = "device.direct";

    public void Initialize(string rabbitMQConnectionString)
    {
        var factory = new ConnectionFactory 
        { 
            Uri = new Uri(rabbitMQConnectionString),
            AutomaticRecoveryEnabled = true
        };
        _connection = factory.CreateConnection();
        _channel = _connection.CreateModel();
        _channel.ExchangeDeclare(EXCHANGE_NAME, ExchangeType.Direct, durable: true);
    }

    // 设备数据采集后转发
    public void ForwardDeviceData(string deviceId, Dictionary<string, object> dataPoints)
    {
        var queueName = $"device.{deviceId}";
        
        // 声明队列(幂等)
        _channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false);
        _channel.QueueBind(queueName, EXCHANGE_NAME, queueName);

        foreach (var point in dataPoints)
        {
            var eventData = new DeviceDataEvent
            {
                DeviceId = deviceId,
                PointAddress = point.Key,
                Value = point.Value,
                CollectTime = DateTime.Now,
                Quality = "Good"
            };

            var json = JsonSerializer.Serialize(eventData);
            var body = Encoding.UTF8.GetBytes(json);

            var props = _channel.CreateBasicProperties();
            props.Persistent = true;
            props.Headers = new Dictionary<string, object>
            {
                ["device_id"] = deviceId,
                ["device_type"] = "PLC"
            };

            _channel.BasicPublish(EXCHANGE_NAME, queueName, false, props, body);
        }
    }
}

3. 处理方示例

csharp
public class DeviceDataConsumer : BackgroundService
{
    private readonly string _deviceId;
    private readonly string _connectionString;

    public DeviceDataConsumer(string deviceId, string connectionString)
    {
        _deviceId = deviceId;
        _connectionString = connectionString;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        var factory = new ConnectionFactory { Uri = new Uri(_connectionString) };
        using var connection = factory.CreateConnection();
        using var channel = connection.CreateModel();

        var queueName = $"device.{_deviceId}";
        channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false);

        // 限制预取数量,避免堆积
        channel.BasicQos(0, 10, false);

        var consumer = new EventingBasicConsumer(channel);
        consumer.Received += async (model, ea) =>
        {
            var json = Encoding.UTF8.GetString(ea.Body.ToArray());
            var eventData = JsonSerializer.Deserialize<DeviceDataEvent>(json);

            try
            {
                // 实时处理:更新SCADA画面、写入实时数据库、触发报警判断
                await ProcessRealtimeData(eventData);
                
                channel.BasicAck(ea.DeliveryTag, false);
            }
            catch (Exception ex)
            {
                channel.BasicNack(ea.DeliveryTag, false, requeue: true);
            }
        };

        channel.BasicConsume(queueName, autoAck: false, consumer: consumer);
    }
}

4. 场景总结

要点说明
数据源PLC/仪表实时采集
触发方式TCP Socket实时接收
Exchange类型Direct Exchange,精确路由到设备队列
消息特点实时、高频、点对点
队列设计每台设备独立队列,便于横向扩展
适用场景实时数据采集、设备监控、SCADA集成

最佳实践

1. 事件命名

csharp
// ✅ 推荐:使用过去式命名
public class OrderCreatedEvent : EventDataBase { }
public class UserDeletedEvent : EventDataBase { }

// ❌ 避免:使用动词
public class CreateOrderEvent : EventDataBase { }

2. 事件体设计

csharp
// ✅ 推荐:事件体包含必要信息
public class OrderCreatedEvent : EventDataBase
{
    public string OrderId { get; set; }
    public decimal Amount { get; set; }
    public string CustomerId { get; set; }
}

// ❌ 避免:事件体过大或包含敏感信息
public class OrderCreatedEvent : EventDataBase
{
    public string FullOrderJson { get; set; }  // 过大
    public string CustomerPassword { get; set; }  // 敏感
}

3. 处理幂等性

csharp
_eventBus.SubscribeEvent<OrderCreatedEvent>(async evt =>
{
    // 处理幂等:相同事件多次处理结果一致
    var exists = await db.Orders.AnyAsync(x => x.Id == evt.OrderId);
    if (!exists)
    {
        await ProcessOrder(evt.OrderId);
    }
});

4. 错误处理

csharp
_eventBus.SubscribeEvent<OrderCreatedEvent>(async evt =>
{
    try
    {
        await ProcessOrder(evt.OrderId);
    }
    catch (Exception ex)
    {
        // 记录日志,不影响其他消息处理
        logger.LogError(ex, "处理订单事件失败: {OrderId}", evt.OrderId);
    }
});

常见问题

事件未收到

排查步骤:

  1. 确认已调用 SubscribeEvent 订阅
  2. 确认事件类型完全匹配
  3. 检查RabbitMQ/Redis连接是否正常
  4. 查看服务日志是否有错误

消息丢失

可能原因:

  • 非持久化消息 + 服务重启
  • Stream Queue超过保留时间

解决方案:

  • 使用持久化消息(默认已启用)
  • 调整Stream Queue保留策略

性能问题

优化建议:

  • 减少事件体大小
  • 避免在事件处理中执行耗时操作
  • 使用Command模式处理耗时任务

HiMind 工业互联网平台 技术文档