为Rebus总线实现自定义ISerializer时,如何获取目标主题名称?
在Rebus自定义ISerializer中获取目标主题名称的方案
要在Rebus的自定义ISerializer的Serialize方法中获取消息要发送到的目标主题,核心是读取消息头中的Headers.Destination字段——Rebus会自动将目标主题(无论是通过ITopicNameConvention自动路由生成,还是手动指定的主题)存入这个头字段,因此这是通用的解决方案。
具体实现步骤
- 在自定义序列化器的
Serialize方法中,从传入的Message对象的Headers集合中提取Headers.Destination值。 - 根据获取到的目标主题,选择对应的Confluent Schema Registry序列化策略(比如匹配不同的Schema ID、主题对应的Schema规则等)。
代码示例
using Rebus.Messages; using Rebus.Serialization; public class SchemaRegistryAwareSerializer : ISerializer { private readonly ISchemaRegistryClient _schemaRegistryClient; public SchemaRegistryAwareSerializer(ISchemaRegistryClient schemaRegistryClient) { _schemaRegistryClient = schemaRegistryClient; } public async Task<TransportMessage> Serialize(Message message) { // 获取目标主题 if (!message.Headers.TryGetValue(Headers.Destination, out var targetTopic)) { throw new InvalidOperationException("无法获取消息的目标主题"); } // 根据目标主题匹配对应的Schema ID var schemaId = await GetSchemaIdForTopic(targetTopic); // 结合Schema ID执行序列化逻辑 var messageBody = await SerializeWithSchema(message.Body, schemaId); return new TransportMessage(message.Headers, messageBody); } // 自定义逻辑:按主题映射Schema ID private async Task<int> GetSchemaIdForTopic(string topic) { if (topic.StartsWith("order-")) { return await _schemaRegistryClient.GetSchemaIdAsync("order-value"); } else if (topic.StartsWith("payment-")) { return await _schemaRegistryClient.GetSchemaIdAsync("payment-value"); } throw new NotSupportedException($"未配置主题{topic}对应的Schema"); } // 自定义逻辑:基于Schema ID完成序列化 private async Task<byte[]> SerializeWithSchema(object messageBody, int schemaId) { // 实现与Confluent Schema Registry兼容的序列化逻辑 // 例如使用AvroSerializer结合schemaId处理消息体 // ... } }
关键说明
- 无论消息是通过自动路由(依赖
ITopicNameConvention)还是手动指定主题发送(比如await bus.Advanced.Routing.Send("custom-topic", message);),Headers.Destination都会被正确填充为目标主题名称。 - 这种方式不依赖
ITopicNameConvention,完全通过Rebus内置的消息头传递主题信息,避免了手动指定主题时的失效问题。
内容的提问来源于stack exchange,提问作者Mentor
相关产品推荐
相关产品推荐

