Spring Reactive Kafka生产者多主题差异化配置咨询:仅指定主题启用消息压缩
针对不同Kafka主题配置差异化生产者的解决方案
嘿,作为Kafka新手碰到这种需要给特定主题做差异化配置的需求太常见了!虽然Reactive Kafka没有直接提供和RoutingKafkaTemplate完全对应的开箱即用组件,但我们可以通过两种实用思路来实现仅对T1主题的消息做压缩的目标,而且完全适配你的Spring WebFlux+Reactive Kafka技术栈:
方案一:创建多个ReactiveKafkaProducerTemplate实例(简单直接)
这种思路最容易上手——为需要特殊配置的主题单独创建一个ProducerTemplate,剩下的用默认配置的模板。
步骤1:配置两个ProducerTemplate Bean
你可以在配置类里分别定义针对T1的压缩配置模板,以及默认模板:
@Configuration public class KafkaProducerConfig { // 默认生产者模板(无压缩,用于T2等主题) @Bean("defaultReactiveKafkaProducerTemplate") public ReactiveKafkaProducerTemplate<String, Object> defaultReactiveKafkaProducerTemplate(KafkaProperties properties) { Map<String, Object> props = properties.buildProducerProperties(); return new ReactiveKafkaProducerTemplate<>(SenderOptions.create(props)); } // 针对T1主题的压缩生产者模板 @Bean("t1ReactiveKafkaProducerTemplate") public ReactiveKafkaProducerTemplate<String, Object> t1ReactiveKafkaProducerTemplate(KafkaProperties properties) { Map<String, Object> props = properties.buildProducerProperties(); // 设置压缩类型,可选gzip/snappy/lz4/zstd,按需选择 props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "gzip"); return new ReactiveKafkaProducerTemplate<>(SenderOptions.create(props)); } }
步骤2:在业务代码中根据主题选择对应模板
在你的业务服务里,注入两个模板,发送消息时根据目标主题选择:
@Service public class KafkaMessageService { private final ReactiveKafkaProducerTemplate<String, Object> defaultTemplate; private final ReactiveKafkaProducerTemplate<String, Object> t1Template; // 构造函数注入 public KafkaMessageService(@Qualifier("defaultReactiveKafkaProducerTemplate") ReactiveKafkaProducerTemplate<String, Object> defaultTemplate, @Qualifier("t1ReactiveKafkaProducerTemplate") ReactiveKafkaProducerTemplate<String, Object> t1Template) { this.defaultTemplate = defaultTemplate; this.t1Template = t1Template; } public Mono<Void> sendMessage(String topic, Object message) { if ("T1".equals(topic)) { return t1Template.send(topic, message).then(); } else { return defaultTemplate.send(topic, message).then(); } } }
这个方案的优点是简单易懂、易于维护,适合主题数量不多的场景;缺点是如果后续需要给更多主题加差异化配置,就得不断新增模板Bean。
方案二:自定义路由式的生产者模板(灵活扩展)
如果你担心未来会有更多主题需要差异化配置,可以自己封装一个类似RoutingKafkaTemplate的Reactive版本,根据主题动态选择对应的SenderOptions。
步骤1:封装路由模板
@Component public class RoutingReactiveKafkaProducerTemplate { private final Map<String, ReactiveKafkaProducerTemplate<String, Object>> templateMap; private final ReactiveKafkaProducerTemplate<String, Object> defaultTemplate; // 构造函数注入所有模板和默认模板 public RoutingReactiveKafkaProducerTemplate( @Qualifier("defaultReactiveKafkaProducerTemplate") ReactiveKafkaProducerTemplate<String, Object> defaultTemplate, @Qualifier("t1ReactiveKafkaProducerTemplate") ReactiveKafkaProducerTemplate<String, Object> t1Template) { this.defaultTemplate = defaultTemplate; this.templateMap = new HashMap<>(); templateMap.put("T1", t1Template); // 后续新增主题配置时,直接往map里加对应的模板即可 } public Mono<Void> send(String topic, Object message) { ReactiveKafkaProducerTemplate<String, Object> template = templateMap.getOrDefault(topic, defaultTemplate); return template.send(topic, message).then(); } }
步骤2:业务代码使用路由模板
@Service public class KafkaMessageService { private final RoutingReactiveKafkaProducerTemplate routingTemplate; public KafkaMessageService(RoutingReactiveKafkaProducerTemplate routingTemplate) { this.routingTemplate = routingTemplate; } public Mono<Void> sendMessage(String topic, Object message) { return routingTemplate.send(topic, message); } }
这个方案的优势是扩展性强,后续新增主题配置只需要添加对应的ProducerTemplate Bean并注册到路由map里,业务代码无需修改。
额外提醒
- 压缩类型的选择:
gzip压缩率高但CPU开销大,snappy/lz4兼顾压缩率和性能,zstd是较新的算法,性能和压缩率都不错,你可以根据自己的业务场景选择。 - 如果你用的是Spring Boot 3.x+的Reactive Kafka版本,也可以通过
SenderOptions.withProducerCustomizer来做更细粒度的配置,但本质上和创建多个模板的思路一致。
内容的提问来源于stack exchange,提问作者vinilpj
相关产品推荐
相关产品推荐

