事件总线
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(广播事件)
| 特性 | 说明 |
|---|---|
| 模式 | 广播,所有订阅者都能收到 |
| RabbitMQ | Fanout Exchange + Stream Queue |
| Redis | Pub/Sub |
| 适用场景 | 通知推送、状态变更广播 |
Command(命令)
| 特性 | 说明 |
|---|---|
| 模式 | 竞争消费,只有一个消费者处理 |
| RabbitMQ | Quorum Queue |
| Redis | Stream + Consumer Group |
| 适用场景 | 任务分发、后台处理 |
配置
连接配置
在 appsettings.json 中配置:
json
{
"ConnectionStrings": {
"redis": "10.11.5.49:6379,password=redis123!",
"rabbitmq": "amqp://user:password@10.11.5.57:5672/dev"
}
}降级策略
| 配置情况 | 使用的实现 |
|---|---|
配置了 rabbitmq | RabbitMqEventBusPlus |
未配置 rabbitmq,配置了 redis | RedisEventBusPlus |
| 都未配置 | 抛出异常 |
注意: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.fanout | Fanout Exchange | 广播事件交换机 |
hmx.events.stream.shared | Stream Queue | 广播事件队列(共享) |
cmd.{完整类名} | Quorum Queue | 命令队列(每个命令类型一个) |
案例一:实验室数据采集
场景:实验室试验机通过Access数据库存储检测数据,系统实时监控文件变化,将数据同步到RabbitMQ,供MES主系统消费。
流程图:
试验机检测 → Access.mdb → FileSystemWatcher → RabbitMQPublisher → RabbitMQ → MES主系统消费1. 队列与消息规划
| 规划项 | 值 | 说明 |
|---|---|---|
| Exchange | lims.topic | Topic Exchange,支持路由键匹配 |
| Queue | lims.testdata | 试验机检测数据队列 |
| RoutingKey | mechanical.* | 路由键,支持通配符匹配 |
| 消息TTL | 72小时 | 超时自动删除 |
| 消息格式 | 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. 队列与消息规划
| 规划项 | 值 | 说明 |
|---|---|---|
| Exchange | device.direct | Direct Exchange,精确路由 |
| Queue | device.{设备编号} | 每台设备一个队列 |
| RoutingKey | device.{设备编号} | 与队列名一致 |
| 消息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);
}
});常见问题
事件未收到
排查步骤:
- 确认已调用
SubscribeEvent订阅 - 确认事件类型完全匹配
- 检查RabbitMQ/Redis连接是否正常
- 查看服务日志是否有错误
消息丢失
可能原因:
- 非持久化消息 + 服务重启
- Stream Queue超过保留时间
解决方案:
- 使用持久化消息(默认已启用)
- 调整Stream Queue保留策略
性能问题
优化建议:
- 减少事件体大小
- 避免在事件处理中执行耗时操作
- 使用Command模式处理耗时任务