如何在MassTransit中记录Routing Slip Activity的异常?
Routing Slip 活动异常日志记录方案
我们的Routing Slip Activity抛出异常时,没有任何日志记录。虽然可以监听RoutingSlipFaulted消息查看异常,但希望能把错误记录到日志,同时纳入OTL数据。想知道能不能在活动的处理管道中添加组件,像普通消费者那样捕获并记录错误?
可行方案:通过自定义活动执行中间件实现
MassTransit的活动执行管道支持添加自定义中间件,和普通消费者的管道逻辑一致,你可以通过中间件捕获活动执行时的异常,完成日志记录和OTL数据上报,同时不破坏原有故障流程。
1. 编写异常日志中间件
创建实现IActivityExecutionFilter接口的中间件类,在执行流程中捕获异常并处理:
public class ActivityExceptionLoggingFilter<TActivity, TArguments> : IActivityExecutionFilter<TActivity, TArguments> where TActivity : class, IExecuteActivity<TArguments> where TArguments : class { private readonly ILogger<ActivityExceptionLoggingFilter<TActivity, TArguments>> _logger; private readonly Tracer _tracer; public ActivityExceptionLoggingFilter(ILogger<ActivityExceptionLoggingFilter<TActivity, TArguments>> logger, Tracer tracer) { _logger = logger; _tracer = tracer; } public async Task Send(ActivityExecutionContext<TActivity, TArguments> context, IPipe<ActivityExecutionContext<TActivity, TArguments>> next) { try { await next.Send(context); } catch (Exception ex) { // 记录包含活动信息的错误日志 _logger.LogError(ex, "Routing Slip活动执行失败,活动类型:{ActivityType},路由单ID:{RoutingSlipId}", typeof(TActivity).Name, context.RoutingSlipId); // 上报异常到OTL链路追踪系统 using var span = _tracer.StartActiveSpan($"Activity-{typeof(TActivity).Name}-Fault"); span.Status = Status.Error.WithDescription(ex.Message); span.RecordException(ex); // 重新抛出异常,保证MassTransit原有故障消息(RoutingSlipFaulted等)正常发送 throw; } } public void Probe(ProbeContext context) { context.CreateFilterScope("activity-exception-logging"); } }
2. 注册中间件到活动端点
在配置活动接收端点时,将自定义中间件添加到执行管道:
services.AddMassTransit(x => { x.AddActivity<Activity3, Activity3Arguments>(); x.UsingInMemory((context, cfg) => { cfg.ReceiveEndpoint("Activity3_execute", e => { e.UseActivityExecute<Activity3, Activity3Arguments>(context, cfg => { cfg.UseFilter(new ActivityExceptionLoggingFilter<Activity3, Activity3Arguments>( context.GetRequiredService<ILogger<ActivityExceptionLoggingFilter<Activity3, Activity3Arguments>>>(), context.GetRequiredService<Tracer>())); }); }); }); });
3. 全局生效(可选)
如果需要所有活动都应用该异常处理逻辑,可以注册全局活动执行过滤器:
services.AddMassTransit(x => { x.AddActivity<Activity3, Activity3Arguments>(); // 注册全局过滤器,对所有活动生效 x.AddActivityExecutionFilter(typeof(ActivityExceptionLoggingFilter<,>)); x.UsingInMemory((context, cfg) => { cfg.ReceiveEndpoint("Activity3_execute", e => { e.UseActivityExecute<Activity3, Activity3Arguments>(); }); }); });
效果说明
- 异常发生时会立即记录详细日志,包含活动类型、路由单ID和完整异常堆栈
- 异常数据会同步上报到OTL系统,便于链路追踪分析
- 不会干扰MassTransit原生故障处理流程,
RoutingSlipActivityFaulted和RoutingSlipFaulted消息仍会正常发送
以下是故障活动的运行日志:
16:31:32.516-D Create send transport: loopback://localhost/activity3_execute 16:31:32.518-D SEND loopback://localhost/activity3_execute dc340000-5ece-a04a-f20d-08db5d351d93 MassTransit.Courier.Contracts.RoutingSlip 16:31:32.518-D RECEIVE loopback://localhost/Activity2_execute dc340000-5ece-a04a-f20d-08db5d351d93 MassTransit.Courier.Contracts.RoutingSlip Activity2(00:00:00.5259605) 16:31:32.523-D Execute Activity: dc340000-5ece-a04a-d752-08db5d351d92 (Activity3, loopback://localhost/Activity3_execute) Done Activity 3 16:31:33.091-D Create send transport: loopback://localhost/urn:message:MassTransit.Courier.Contracts:RoutingSlipActivityFaulted 16:31:33.096-D SEND loopback://localhost/urn:message:MassTransit.Courier.Contracts:RoutingSlipActivityFaulted dc340000-5ece-a04a-0382-08db5d351e97 MassTransit.Courier.Contracts.RoutingSlipActivityFaulted 16:31:33.097-D Create send transport: loopback://localhost/urn:message:MassTransit.Courier.Contracts:RoutingSlipFaulted 16:31:33.102-D SEND loopback://localhost/urn:message:MassTransit.Courier.Contracts:RoutingSlipFaulted dc340000-5ece-a04a-ffe3-08db5d351e97 MassTransit.Courier.Contracts.RoutingSlipFaulted 16:31:33.103-D RECEIVE loopback://localhost/Activity3_execute dc340000-5ece-a04a-f20d-08db5d351d93 MassTransit.Courier.Contracts.RoutingSlip Activity3(00:00:00.5843245)
内容的提问来源于stack exchange,提问作者Kevin Smith
相关产品推荐
相关产品推荐

