ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

workflow-core 外抛事件:Channel 管道 + 多发布器实现

workflow-core 外抛事件:Channel 管道 + 多发布器实现 在订单履约系统里接 workflow-core 做流程编排已经快两年了。最近业务方提了个很实际的需求工作流每跑完一个步骤、或者整个流程开始/结束时外部系统BI、移动端、ERP都要能实时感知而不是靠定时任务轮询数据库。这就牵出一个很经典的问题怎么给 workflow-core 扩展“外抛事件”。这里说的“外抛事件”不是 workflow-core 里那个用来接收外部信号的 PublishEvent/WaitFor 机制而是指把工作流内部的生命周期变化以事件的形式主动推给外部系统。这几乎是所有用到 workflow-core 的生产项目都会遇到的诉求流程内部的状态不能只停留在引擎里得让外面的世界“看见”它。这篇文章会拆清楚 workflow-core 的事件体系对比几种扩展方案然后给出一套可直接抄作业的实现用 Channel 解耦、用发布器对接 RabbitMQ 和 Webhook、再配合 WaitFor 做外部回调恢复流程。无论你是刚接触 workflow-core还是已经在生产环境踩过坑这篇都应该能帮你少绕几个弯。1. 先搞清楚 workflow-core 的事件体系再动手扩展之前最忌讳的一件事就是没分清 workflow-core 里那两套“事件”到底谁是谁。这俩名字相近、方向相反很多人一上来就把它们混在一起后面代码越写越乱。1.1 两类“事件”别混淆入站事件与出站事件workflow-core 本身内置了两套事件机制一套负责“进来”一套负责“出去”。“进来”的是 PublishEvent 和 WaitFor。工作流里可以放一个 WaitFor 节点执行到那里就挂起等着外部通过host.PublishEvent(eventName, eventKey, eventData)把数据送进来然后工作流继续往后跑。典型场景就是人工审批流程走到审批节点挂起来等审批结果审批系统调用 PublishEvent 把结果送回工作流。这是事件驱动思路里标准的“入站”通道。“出去”的是它内部暴露的生命周期事件。WorkflowHost 上挂着一个WorkflowEvents事件引擎里每发生一次状态变化都会触发一次比如流程启动、流程完成、某个 Step 开始、某个 Step 出错。这些事件包含了 WorkflowInstanceId、StepId、时间戳等信息是观察工作流内部运行状态的原始信号。我见过不少同事把 PublishEvent 当成“外抛”来用方向完全反了。PublishEvent 是外部系统把事件推进工作流不是工作流把事件推出去。真正要外抛给外部系统的是 WorkflowEvents 这条出站通道。1.2 官方 WorkflowEvents 能做什么、做不到什么先看看官方的出站事件长什么样var host serviceProvider.GetRequiredServiceIWorkflowHost(); host.WorkflowEvents (sender, evt) { Console.WriteLine(${DateTime.UtcNow:O} {evt.GetType().Name}); };这段代码看起来没毛病但放到生产环境里很快就会发现三个尴尬的问题。第一这个事件是进程内委托只在当前进程里生效。进程重启、应用发布、容器漂移订阅关系就没了期间产生的事件一个都不会补给你。第二事件触发是同步的它直接运行在工作流引擎的执行线程上。如果你在 handler 里做 RabbitMQ 发送、HTTP 调用这种耗时操作整个工作流会被卡住。第三事件本身没有持久化也没有重试机制发送失败就是失败不会有任何补偿。所以要做的不是“订阅一下打个日志”而是把它扩展成一条完整的管道拿到引擎抛出的原始事件包装成统一格式放进缓冲队列由后台服务异步派发给各种外部通道并且要能处理重试和失败。1.3 扩展目标一条不阻塞引擎的外抛管道基于上面的痛点我在设计时给这个扩展定了四个硬性目标。核心原则是低侵入。不修改 workflow-core 的源码不重写它的执行器只在外面包一层订阅者。官方的 WorkflowEvents 就是给扩展留的口子顺着它走最稳。然后是异步化。事件 handler 只做一件事把事件包装好塞进 Channel立刻返回。真正发消息的活交给后台消费者去干这样无论外部通道多慢都不会拖累工作流引擎。第三是通道可插拔。同一个事件信封可以同时发到 RabbitMQ、HTTP Webhook、本地日志甚至后续加 Kafka 也只是多实现一个发布器而已不能把通道写死在代码里。最后是要可控。事件类型可能很多步骤级事件尤其高频必须能通过配置选择哪些事件发出去、哪些不发不然一个几万步的大流程跑起来光外抛事件就能把下游系统打到崩溃。这四个目标定完之后扩展的方向就非常清楚了做一条异步事件管道而不是在引擎里塞一堆业务代码。2. 方案选型与整体设计目标明确之后我花了一点时间对比了三条实现路线。每条路线都有它的适用场景但最终选定的方案差别很大。2.1 三条路线对比为什么要选订阅转发路线 A直接订阅 WorkflowHost.WorkflowEvents在事件 handler 中做转发。这是最接近官方设计意图的做法几乎零侵入不需要动引擎内部的任何东西。劣势是拿不到执行器内部那些更底层的调用栈细节但对我们外抛业务事件这个需求来说WorkflowEvents 提供的信息已经足够了。路线 B替换或包装 WorkflowCore 的持久化接口、执行器接口。这个路线能拿到更细的粒度你可以在一次 Step 执行时捕获到方法进出、数据变化全过程。但代价是侵入性太高workflow-core 的内部实现版本迭代很快升级时很容易被破坏。除非你要做字节码级调试或者深度性能分析否则我不建议走这条路。路线 C在每个自定义 Step 里手动写发事件的代码。这是最直接的做法也是我最不推荐的做法。它把事件发送逻辑散落在业务步骤里业务代码和工作流引擎耦合在一起换一个事件通道要改一大堆地方。我见过一个项目就是因为前期这么干后期把 RabbitMQ 换 Kafka 时前前后后改了二十几个 Step还漏了几个导致线上事件丢失。三条路线放一起对比就很直观路线侵入性可控性维护成本推荐度订阅 WorkflowEvents 转发低中低强烈推荐替换执行器/持久化层高高高不推荐自定义 Step 内手动发中低高不推荐2.2 事件管道设计从引擎到外部通道的四个环节选订阅转发这条路之后整个管道我设计成了四个环节。最前端是事件源也就是 WorkflowHost 的 WorkflowEvents。第二环是一个订阅桥它在应用启动时挂上订阅、停止时摘掉订阅并且负责把原始事件转换成统一信封写入内存队列。第三环是后台转发服务它从队列里取出信封经过事件类型过滤再调用发布器。最后一环是发布器负责把信封真正送到外部系统。用文字描述这条链路就是WorkflowCore 引擎 → WorkflowEvents 原始事件 → WorkflowEventBridge 订阅并入队 → ChannelEventEnvelope 缓冲队列 → EventForwardingService 后台消费 → IExternalEventPublisher 发布器 → RabbitMQ / Webhook / 控制台日志这里用 Channel 而不是直接用 BlockingCollection主要是看中它的异步读写能力。工作流引擎的线程把事件写入 Channel 时用的是 TryWrite几乎零等待不会阻塞引擎。后台可以同时启动多个转发消费者处理不了积压时还能横向加消费者比直接同步调用优雅得多。2.3 统一事件信封EventEnvelope 的字段设计外抛事件要想被下游系统稳定消费统一的信封格式比什么都重要。我定义的 EventEnvelope 是这样的public sealed class EventEnvelope { public string EventId { get; set; } string.Empty; public DateTime Timestamp { get; set; } public string Source { get; set; } string.Empty; public string WorkflowDefinitionId { get; set; } string.Empty; public string WorkflowInstanceId { get; set; } string.Empty; public string? StepId { get; set; } public string? StepName { get; set; } public string EventType { get; set; } string.Empty; public object? Data { get; set; } }字段不多但每个都有它的用途。EventId 是全局唯一的事件编号用于下游幂等去重必须生成一次就定下来不能重新生成。Timestamp 用 UTC 时间避免不同时区的服务解读出偏差。Source 标记事件来源系统比如订单服务、履约服务在多系统联动时能快速定位事件是从哪里来的。WorkflowDefinitionId 和 WorkflowInstanceId 是定位一次流程运行的关键维度下游可以用它们做关联查询。StepId 和 StepName 在步骤级事件里才有值流程级事件里为空。EventType 是事件类型标识我用的是点分字符串例如workflow.started、step.completed方便在消息队列的路由键里直接使用。Data 字段放的是事件附带的数据比如出错时的异常信息、审批请求的上下文由具体场景决定。2.4 配置项设计事件不是发得越多越好事件管道建好后最容易被忽略的是配置设计。我见过有人把步骤级事件统统外抛结果一个流程跑完下游系统收到了五百多个事件把消息队列都打满了。所以我在扩展里加了一个 ExternalEventOptions用配置控制事件外抛的边界。public sealed class ExternalEventOptions { public bool Enabled { get; set; } true; public string Source { get; set; } unknown; public HashSetstring IncludeEventTypes { get; set; } new(); public HashSetstring ExcludeEventTypes { get; set; } new(); public RabbitMqOptions RabbitMq { get; set; } new(); public WebhookOptions Webhook { get; set; } new(); }IncludeEventTypes 和 ExcludeEventTypes 相当于白名单和黑名单。当 IncludeEventTypes 有值时只有列表内的事件会被外抛ExcludeEventTypes 则相反列表内的事件会被屏蔽。实际项目里我通常只抛流程级事件workflow.started、workflow.completed、workflow.errored和特定业务的步骤事件步骤级的事件默认全关按需打开这样下游的压力会小很多。还有一个常被忽略的配置是 Source这个字段需要每个微服务设置成自己的名字否则多个服务共用一个 RabbitMQ 时下游根本分辨不出事件到底是从哪个系统发出来的。3. 实操实现一个可上线的外抛事件扩展设计定了代码就好写了。这一部分我会把核心实现拆成几块每一块都给出可以直接用的代码。整个扩展我建议放在独立的类库项目里这样 workflow-core 的业务项目干净不掺杂这些基础设施代码。3.1 先定义发布器接口把通道差异隔离掉外抛事件的通道可能是 RabbitMQ、Kafka、Webhook也可能是测试时用的控制台。为了不把这些通道细节散落在业务代码里我设计了一个很小的接口public interface IExternalEventPublisher { Task PublishAsync(EventEnvelope envelope, CancellationToken cancellationToken default); }发布器只干一件事把信封发出去。至于是发到队列还是 HTTP那是实现类自己的事。实际项目中可能会同时配置多个发布器比如既发 RabbitMQ 又发 Webhook。这种情况下可以用一个聚合发布器把事件转发给所有注册的发布器public sealed class CompositeEventPublisher : IExternalEventPublisher { private readonly IEnumerableIExternalEventPublisher _publishers; public CompositeEventPublisher(IEnumerableIExternalEventPublisher publishers) { _publishers publishers; } public async Task PublishAsync(EventEnvelope envelope, CancellationToken cancellationToken default) { foreach (var publisher in _publishers) { await publisher.PublishAsync(envelope, cancellationToken); } } }这个聚合器让主流程不需要关心到底挂着几个发布器新增通道时只要往容器里注册新实现就行。3.2 订阅与入队WorkflowEventBridge 的正确姿势接下来是管道的入口。我把它写成一个 IHostedService目的是让订阅生命周期和应用生命周期绑定应用启动时订阅应用停止时退订避免因宿主重启导致事件订阅丢失。public sealed class WorkflowEventBridge : IHostedService { private readonly IWorkflowHost _host; private readonly ChannelEventEnvelope _channel; private readonly EventEnvelopeFactory _factory; private readonly ExternalEventOptions _options; public WorkflowEventBridge( IWorkflowHost host, ChannelEventEnvelope channel, EventEnvelopeFactory factory, IOptionsExternalEventOptions options) { _host host; _channel channel; _factory factory; _options options.Value; } private void OnWorkflowEvent(object? sender, WorkflowEvent evt) { if (!_options.Enabled) { return; } var envelope _factory.Create(evt); if (envelope is null) { return; } _channel.Writer.TryWrite(envelope); } public Task StartAsync(CancellationToken cancellationToken) { _host.WorkflowEvents OnWorkflowEvent; return Task.CompletedTask; } public Task StopAsync(CancellationToken cancellationToken) { _host.WorkflowEvents - OnWorkflowEvent; return Task.CompletedTask; } }这段代码看起来简单但有一个关键点事件 handler 里只做 TryWrite不做任何重活。我刚开始写的时候第一版直接在这个 handler 里调了 RabbitMQ结果一次流程跑下来原本几十毫秒的步骤硬生生膨胀到几秒。WorkflowEvents 是同步触发的它跑在引擎工作线程里任何阻塞都会拖慢整个流程。改成 Channel 异步消费之后这个问题立刻消失。TryWrite 而不是 WriteAsync 也是同样的原因。handler 是同步上下文不应该在这里等队列容量。Unbounded Channel 的 TryWrite 永远不会因为队列满而阻塞这正是我要的行为。3.3 事件工厂把 WorkflowEvent 转换成 EventEnvelopeWorkflowEventBridge 把转换逻辑交给了 EventEnvelopeFactory。这个工厂负责解读各种生命周期事件填出一个完整的信封。这里通常会用到模式匹配把不同类型的原始事件映射到标准 EventType。public sealed class EventEnvelopeFactory { private readonly ExternalEventOptions _options; public EventEnvelopeFactory(IOptionsExternalEventOptions options) { _options options.Value; } public EventEnvelope? Create(WorkflowEvent evt) { if (!IsAllowed(evt.GetType().Name)) { return null; } var envelope new EventEnvelope { EventId Guid.NewGuid().ToString(N), Timestamp evt.EventTimeUtc, Source _options.Source }; switch (evt) { case WorkflowStarted w: envelope.EventType workflow.started; envelope.WorkflowDefinitionId w.WorkflowDefinitionId; envelope.WorkflowInstanceId w.WorkflowInstanceId; break; case WorkflowCompleted c: envelope.EventType workflow.completed; envelope.WorkflowDefinitionId c.WorkflowDefinitionId; envelope.WorkflowInstanceId c.WorkflowInstanceId; break; case WorkflowError e: envelope.EventType workflow.errored; envelope.WorkflowDefinitionId e.WorkflowDefinitionId; envelope.WorkflowInstanceId e.WorkflowInstanceId; envelope.Data new { e.Exception.Message, e.Exception.StackTrace }; break; case StepStarted s: envelope.EventType step.started; envelope.WorkflowInstanceId s.WorkflowInstanceId; envelope.StepId s.StepId; envelope.StepName s.StepName; break; case StepCompleted s: envelope.EventType step.completed; envelope.WorkflowInstanceId s.WorkflowInstanceId; envelope.StepId s.StepId; envelope.StepName s.StepName; break; case StepError s: envelope.EventType step.errored; envelope.WorkflowInstanceId s.WorkflowInstanceId; envelope.StepId s.StepId; envelope.StepName s.StepName; envelope.Data new { s.Exception.Message }; break; default: return null; } return envelope; } private bool IsAllowed(string eventType) { if (_options.IncludeEventTypes.Count 0 !_options.IncludeEventTypes.Contains(eventType)) { return false; } if (_options.ExcludeEventTypes.Contains(eventType)) { return false; } return true; } }这里补充一句WorkflowEvent 及各个子类的属性名在不同小版本里可能略有差异写代码时以你实际引用的 NuGet 包源码为准。我用的 WorkflowInstanceId、StepId、StepName 这几个属性在目前的主流版本里都存在编译不过的话去 WorkflowCore 源码里确认一下即可。另外注意 Data 字段我给的是一个匿名对象而不是整个异常对象。整段 StackTrace 塞进事件里会让信封变得很大而且下游不一定需要把 Message 和关键信息挑出来是更克制也更安全的做法。3.4 后台转发服务消费队列并执行发布器入队搞定后需要一个后台服务消费 Channel 里的信封。这个服务是 ASP.NET Core 的 BackgroundService应用启动时会自动运行。public sealed class EventForwardingService : BackgroundService { private readonly ChannelEventEnvelope _channel; private readonly IExternalEventPublisher _publisher; private readonly ILoggerEventForwardingService _logger; public EventForwardingService( ChannelEventEnvelope channel, IExternalEventPublisher publisher, ILoggerEventForwardingService logger) { _channel channel; _publisher publisher; _logger logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { await foreach (var envelope in _channel.Reader.ReadAllAsync(stoppingToken)) { try { await _publisher.PublishAsync(envelope, stoppingToken); } catch (Exception ex) { _logger.LogError(ex, Failed to publish workflow event {EventId} {EventType}, envelope.EventId, envelope.EventType); } } } }这里有一个我要专门提醒的坑不能因为一个事件发送失败就把整个后台服务搞挂。PublishAsync 抛异常时我先把异常记录下来然后继续消费下一个信封。至于这条失败的事件要不要重试、要不要进死信队列取决于你的发布器实现。简单场景下记录日志就够了但如果是核心业务事件建议在发布器内部做指数退避重试并把最终仍失败的投递到本地持久化队列。3.5 对接 RabbitMQ用路由键承载事件类型RabbitMQ 发布器是项目里最常用的通道。我用 RabbitMQ.Client 实现了一个极简发布器Exchange 用 Topic 类型路由键直接复用 EventType这样下游可以用通配符灵活订阅自己关心的事件。public sealed class RabbitMqEventPublisher : IExternalEventPublisher { private readonly IConnection _connection; private readonly IModel _channel; private readonly string _exchange; public RabbitMqEventPublisher(IOptionsExternalEventOptions options) { var rmq options.Value.RabbitMq; var factory new ConnectionFactory { HostName rmq.Host, Port rmq.Port, UserName rmq.UserName, Password rmq.Password }; _connection factory.CreateConnection(); _channel _connection.CreateModel(); _exchange rmq.Exchange; _channel.ExchangeDeclare(_exchange, ExchangeType.Topic, durable: true); } public Task PublishAsync(EventEnvelope envelope, CancellationToken cancellationToken default) { var body JsonSerializer.SerializeToUtf8Bytes(envelope); var props _channel.CreateBasicProperties(); props.MessageId envelope.EventId; props.Timestamp new AmqpTimestamp( new DateTimeOffset(envelope.Timestamp).ToUnixTimeSeconds()); props.ContentType application/json; props.DeliveryMode 2; _channel.BasicPublish( exchange: _exchange, routingKey: envelope.EventType, mandatory: true, basicProperties: props, body: body); return Task.CompletedTask; } }注意 BasicPublish 的 mandatory 参数。设为 true 时如果路由键匹配不到任何队列RabbitMQ 会把消息退回给生产者需要额外处理退回回调。如果不需要这个保障mandatory 改成 false 更省事。生产环境我建议把 mandatory 设为 true再监听 BasicReturn 事件打日志这样能及时发现下游队列配置错误不至于消息发出去就消失了。3.6 对接 Webhook用 HttpClient 加指数退避有些外部系统只暴露 HTTP 接口没有消息队列那就需要 Webhook 发布器。实现一点都不复杂但重试机制必须做否则网络抖动一次事件就这么丢了。public sealed class WebhookEventPublisher : IExternalEventPublisher { private readonly HttpClient _httpClient; private readonly string _endpoint; public WebhookEventPublisher(IOptionsExternalEventOptions options) { _httpClient new HttpClient(); _endpoint options.Value.Webhook.Endpoint; } public async Task PublishAsync(EventEnvelope envelope, CancellationToken cancellationToken default) { var response await _httpClient.PostAsJsonAsync(_endpoint, envelope, cancellationToken); response.EnsureSuccessStatusCode(); } }这是最简版本。在实际项目里我会加一个重试策略第一次失败后等 1 秒第二次失败等 3 秒第三次失败等 9 秒最多重试 5 次。用 .NET 自带的 ResiliencePipeline.NET 8或者 Polly 都可以核心是避免在 PublishAsync 里无限制重试因为后台服务是串行处理事件的一个事件卡住重试后面所有事件都会排队积压。3.7 注册扩展方法一个方法搞定所有装配上面的所有部件需要一个装配方法。我把它写成 ServiceCollection 的扩展方法这样业务侧只需要一行注册代码。public static class WorkflowCoreExternalEventExtensions { public static IServiceCollection AddWorkflowCoreExternalEvents( this IServiceCollection services, ActionExternalEventOptions configure) { services.Configure(configure); services.AddSingletonEventEnvelopeFactory(); services.AddSingletonChannelEventEnvelope(_ Channel.CreateUnboundedEventEnvelope( new UnboundedChannelOptions { SingleReader false, SingleWriter false })); services.AddSingletonIExternalEventPublisher, RabbitMqEventPublisher(); services.AddSingletonIExternalEventPublisher, WebhookEventPublisher(); services.AddSingletonIExternalEventPublisher, CompositeEventPublisher(); services.AddHostedServiceWorkflowEventBridge(); services.AddHostedServiceEventForwardingService(); return services; } }有几个细节说一下。CompositeEventPublisher 注册的是解析 IExternalEventPublisher 时的最终实现它会拿到容器里所有 IExternalEventPublisher 实现包括它自身的问题需要避免通常在 Composite 中排除自身或者让 EventForwardingService 直接注入 CompositeEventPublisher。最干净的做法是 EventForwardingService 直接注入 CompositeEventPublisher而其他实现不注册成 IExternalEventPublisher或者 Composite 里过滤掉类型为 CompositeEventPublisher 的实例。注册之后启动一个测试工作流观察日志或者消费 RabbitMQ 里的消息就能看到 workflow.started、step.started、step.completed、workflow.completed 这些事件按顺序出现。看到这个队列基本就说明管道通了。4. 进阶外抛事件与 WaitFor 配合实现外部回调恢复流程外抛事件解决了“流程状态主动同步出去”的问题。但只做单向外抛还不够很多场景需要闭环工作流挂起等待外部处理事件抛出去之后外部系统处理完再把结果送回工作流流程继续跑。这正是外抛事件和 WaitFor 配合的经典用法。4.1 典型场景人工审批拿审批流举例。一个订单审批流程走到“审批”节点时工作流不能继续推进得停下来等人处理。这时候要把这个“等待”状态主动通知给审批系统否则审批系统根本不知道有单子要审。完整链路是这样提交订单 → 创建审批任务 → 工作流挂起并等待“ApprovalEvent” → 外抛“等待审批”事件给审批系统 → 审批人操作 → 审批系统调用 PublishEvent(“ApprovalEvent”) → 工作流被唤醒 → 处理审批结果 → 继续后续步骤外抛事件和 WaitFor 在这里形成了一个双向通道工作流向外部抛事件告诉外部系统“我在等你”外部系统处理完又通过 PublishEvent 把结果推进工作流。这个闭环是 workflow-core 在审批流、人审、异步补偿场景下的核心用法。4.2 工作流定义中用 WaitFor 挂起节点先在流程定义里加上等待节点public sealed class ApprovalWorkflow : IWorkflowApprovalData { public string Id ApprovalWorkflow; public int Version 1; public void Build(IWorkflowBuilderApprovalData builder) { builder .StartWithSubmitOrderStep() .WaitFor(ApprovalEvent, context context.Workflow.WorkflowInstanceId) .Output(data data.ApprovalResult, step step.EventData) .ThenHandleApprovalResultStep(); } }这里 WaitFor 的 eventKey 用的是当前工作流实例 ID。这样外部回传事件时只需要把 WorkflowInstanceId 带回来就能精确唤醒对应的工作流实例。4.3 外抛事件时把回调信息带出去问题来了工作流挂起之后审批系统是怎么收到通知的WaitFor 节点本身不会自动通知外部系统它只是静默等待。所以必须在流程进入等待之前主动向外抛出事件。public sealed class ApprovalWorkflow : IWorkflowApprovalData { public void Build(IWorkflowBuilderApprovalData builder) { builder .StartWithSubmitOrderStep() .ThenNotifyApprovalSystemStep() // 外抛等待审批事件 .WaitFor(ApprovalEvent, context context.Workflow.WorkflowInstanceId) .Output(data data.ApprovalResult, step step.EventData) .ThenHandleApprovalResultStep(); } }在 NotifyApprovalSystemStep 里就可以调用前面建好的外抛管道public sealed class NotifyApprovalSystemStep : StepBodyAsync { private readonly EventEnvelopeFactory _factory; private readonly ChannelEventEnvelope _channel; public override async TaskExecutionResult RunAsync(IStepExecutionContext context) { var envelope _factory.CreateForExternalWait( workflowInstanceId: context.Workflow.WorkflowInstanceId, eventName: ApprovalEvent, eventKey: context.Workflow.WorkflowInstanceId, data: new { ApproveUrl $/api/approvals/{context.Workflow.WorkflowInstanceId} }); _channel.Writer.TryWrite(envelope); return ExecutionResult.Next(); } }这里有一个关键设计信封里带了 EventName 和 EventKey。外部审批系统不需要理解 workflow-core 的 WaitFor 机制它只需要在回调时原样返回这两个值工作流就能精确匹配。4.4 外部系统处理完用 PublishEvent 回传结果审批人在界面上点了“通过”之后审批后端要做的事情很简单await host.PublishEvent( eventName: ApprovalEvent, eventKey: workflowInstanceId, eventData: new ApprovalResult { Approved true, Comment 同意 });PublishEvent 会拿着 eventName 和 eventKey 去所有挂起的工作流实例里找匹配项找到就把 eventData 写入 WaitFor 节点的输出。此时工作流被唤醒继续执行 HandleApprovalResultStep。整套闭环跑通后你会意识到一个问题审批系统调 PublishEvent 时它怎么知道 eventName 是什么答案是这个消息就在最初外抛的事件信封里。所以外抛事件不仅是“状态同步”还相当于把回传所需的“协议信息”一起送给了外部系统。设计事件信封时Data 里除了业务数据一定要带上回调所需的 eventName 和 eventKey这样外部系统才能无脑配合不需要跟工作流定义耦合。5. 常见问题与避坑实录实现这套扩展的过程中我踩过不少坑也帮同事排查过不少问题。列几个最典型的给后面要做的朋友提个醒。5.1 事件 handler 里做重活工作流被拖慢几秒这个问题我前面提过一次但值得单独再讲。WorkflowEvents 是在引擎线程里同步调用的如果你在 handler 里发 HTTP 请求或写数据库每触发一次事件就阻塞一次引擎。我第一版就是这么干的直接在 handler 里调了 RabbitMQ结果一个流程跑完多了几十秒。排查方法其实很简单在看日志的时候发现事件 handler 的耗时明显集中在网络调用上基本就是这个问题。解决方式也很直接handler 里只做 TryWrite发送交给后台服务。5.2 事件重复投递消费端必须做幂等外抛这件事标准语义就是“至少一次”不是“恰好一次”。RabbitMQ 确认机制、Webhook 重试都会导致同一个事件被投递两次。所以下游消费端必须按 EventId 做幂等处理过的直接跳过。我的建议是消费端维护一张去重表主键就是 envelope.EventId或者用 Redis SETNX。投递端要保证同一个事件重试时 EventId 不变而不能每次重试都 new 一个这就是 EventEnvelope 里 EventId 的最大价值。5.3 序列化时类型丢失下游拿不到强类型System.Text.Json 序列化 EventEnvelope 中的 Data 字段时如果 Data 是匿名对象问题不大但如果你把 Data 定义为 Dictionarystring, object序列化后反序列化回来object 的值会变成 JsonElement。下游拿到的不是 string、int而是必须再套一层 JsonElement 解析很多人第一次遇到都会懵。解决办法有两个一是 Data 字段用强类型 DTO这样编译期就能保证类型二是接收端统一按 JsonElement 解析不要指望从字典里直接 ToString 出友好内容。这个设计要尽早定等下游一堆服务依赖了某个格式再改全是泪。5.4 外抛事件名和 PublishEvent 事件名撞车造成死循环我最开始给外抛事件起的类型名和 WaitFor 等待的入站事件名是同一个。结果工作流很快出现一个诡异的现场流程跑起来后异常频繁日志里同一个事件反复出现。原因在于命名冲突。外抛事件发到了外部系统外部系统处理后又调用 PublishEvent结果这个 PublishEvent 又把工作流里对应的 WaitFor 唤醒了两个机制用同一个事件名串了。解决方式是把两套事件命名空间分开。外抛事件统一用ext.前缀比如ext.workflow.started、ext.approval.waiting入站事件统一用in.前缀比如in.approval.result。这样既能一眼看出方向也从根本上避开了循环触发。5.5 宿主重启后订阅丢失事件管道静默失效WorkflowEvents 的订阅是内存态的。如果你在控制台程序里手动host.WorkflowEvents handler进程重启后这个订阅就没了而且不会报错。生产环境的表现就是日志看着正常流程跑得很顺利但外部系统收不到任何事件。我的做法是把这个订阅放进 WorkflowEventBridge并且注册为 IHostedService。这样应用的启动和停止必然带着订阅的挂上和摘下。如果你在 ASP.NET Core 里用尤其注意不要只在临时作用域里订阅一定要跟着应用生命周期走。这个扩展在我这边的订单履约系统已经跑了小半年最直观的收益是把实时状态同步从轮询数据库改成了事件驱动。但真要说起来扩展本身的技术实现只是其中一半另一半是对事件语义的设计哪些事件值得抛、事件命名怎么规划、下游怎么幂等消费这些才是真正决定这个扩展好不好用的地方。如果只让我留一条经验我会说从第一版开始就给每个事件带上全局唯一的 EventId并且从一开始就把ext.和in.两套命名空间分开。这两个决定越早做后面省的事越多这个坑我替你踩过了。
返回列表