使用MassTransit Saga的EventActivityBinder.Produce发送Kafka消息键的方法问询
问题
我正在开发一个结合MassTransit、Saga与Kafka的项目,遇到了一个问题:我尝试使用EventActivityBinder.Produce方法向Kafka主题发送消息,但已注册的生产者要求必须指定消息键,我却无法找到在该方法中设置消息键的方式,因此触发了“生产者未注册”的异常。
我的问题是:能否通过上述Produce方法发送消息键?
我的实现方法:
public static EventActivityBinder<OrderRequestSagaInstance, ErrorMessageEvent> NotifySourceSystem( this EventActivityBinder<OrderRequestSagaInstance, ErrorMessageEvent> binder) { var @event = binder.Produce(context => { context.Saga.UpdatedAt = DateTime.Now; var @event = new { Success = false, context.Saga.Reason, __Header_Reason = context.Saga.Reason, }; return context.Init<ResponseWrapper<OrderResponseEvent>>(@event); }); return @event; }
生产者注册代码:
rider.AddProducer<string, ResponseWrapper<OrderResponseEvent>>(kafkaTopics.SourceSystemTopic);
解决方案
可以通过EventActivityBinder.Produce方法设置Kafka消息键,你需要利用带发送上下文配置的重载版本来指定键值。
方法一:在消息初始化后设置键
修改你的NotifySourceSystem方法,在初始化消息后,通过SendContext的SetMessageKey方法指定键:
public static EventActivityBinder<OrderRequestSagaInstance, ErrorMessageEvent> NotifySourceSystem( this EventActivityBinder<OrderRequestSagaInstance, ErrorMessageEvent> binder) { return binder.Produce(context => { context.Saga.UpdatedAt = DateTime.Now; var message = new { Success = false, context.Saga.Reason, __Header_Reason = context.Saga.Reason, }; var sendContext = context.Init<ResponseWrapper<OrderResponseEvent>>(message); // 设置Kafka消息键,这里用Saga中的唯一标识(比如OrderId)作为键值 sendContext.SetMessageKey(context.Saga.OrderId.ToString()); return sendContext; }); }
方法二:单独配置发送上下文
也可以在Produce的第二个委托参数中专门配置发送上下文,分离消息内容和Kafka参数:
public static EventActivityBinder<OrderRequestSagaInstance, ErrorMessageEvent> NotifySourceSystem( this EventActivityBinder<OrderRequestSagaInstance, ErrorMessageEvent> binder) { return binder.Produce(context => { context.Saga.UpdatedAt = DateTime.Now; var message = new { Success = false, context.Saga.Reason, __Header_Reason = context.Saga.Reason, }; return context.Init<ResponseWrapper<OrderResponseEvent>>(message); }, (sendContext, sagaContext) => { // 在这里设置消息键,类型要和注册生产者时的键类型一致(你注册的是string类型) sendContext.SetMessageKey(sagaContext.Saga.OrderId.ToString()); }); }
关键说明
- 你注册的生产者键类型为
string(AddProducer<string, ResponseWrapper<OrderResponseEvent>>),所以SetMessageKey必须传入string类型的值,比如Saga实例中的订单ID、唯一标识等。 - 之前触发“生产者未注册”异常,是因为MassTransit找不到匹配的生产者:当你注册了带键类型的生产者,发送时必须提供对应类型的键,否则MassTransit会尝试使用无键生产者,而你并未注册这类生产者,因此报错。
内容的提问来源于stack exchange,提问作者Stanley_Marsh
相关产品推荐
相关产品推荐

