使用Kafka Parallel Consumer的ReactorProcessor生产事件失败求助
Kafka Parallel Consumer ReactorProcessor无法生产事件的解决方案
问题核心
你的ReactorProcessor代码无法生产事件,根源在于返回的流结构不符合框架预期,同时需要确保ReactorProcessor的配置与ParallelStreamProcessor保持一致。
修正后的代码示例
1. 正确的ReactorProcessor消费生产逻辑
pConsumer.react { context -> // 使用context原生的Flux流处理每条消息,避免转换为List批量处理 context.flux() .doOnNext { record -> println("Consuming ${record.partition()}:${record.offset()}") } // 为每条消费消息映射对应的生产记录 .map { record -> ProducerRecord<String, JsonObject>( "output", record.key(), JsonObject(mapOf("someTest" to record.offset())) ) } }
2. 确保ReactorProcessor的创建配置正确
必须在创建ReactorProcessor时传入Producer实例,否则框架没有生产客户端无法发送消息:
private fun createReactorPConsumer(): ReactorProcessor<String, JsonObject> { val producer = KafkaProducerBuilder.getProducer(kafkaConsumerConfig) val options = ParallelConsumerOptions.builder<String, JsonObject>() .ordering(ParallelConsumerOptions.ProcessingOrder.KEY) .maxConcurrency(parallelConsumerConfig.maxConcurrency) .batchSize(parallelConsumerConfig.batchSize) .consumer(buildConsumer(kafkaConsumerConfig)) .producer(producer) // 必须配置Producer实例 .build() return ReactorProcessor.createEosReactorProcessor(options) }
原因解释
- 流结构错误:原代码将所有生产记录打包成List,用
Mono.just(results)返回,框架无法识别每条消费消息对应的生产动作,因此不会触发生产逻辑。ReactorProcessor要求返回的是对应每条消费消息的生产记录流(Flux),而非批量集合。 - 配置缺失:如果创建ReactorProcessor时未传入Producer实例,框架没有可用的生产客户端,自然无法完成消息发送。
可选:监听生产结果
如果需要像ParallelStreamProcessor那样监听生产成功的回调,可以在流中添加扩展处理:
pConsumer.react { context -> context.flux() .doOnNext { record -> println("Consuming ${record.partition()}:${record.offset()}") } .map { record -> ProducerRecord<String, JsonObject>( "output", record.key(), JsonObject(mapOf("someTest" to record.offset())) ) } .doOnNext { producerRecord -> println("Message ${producerRecord.key()} sent successfully") } }
内容的提问来源于stack exchange,提问作者Ehud Lev
相关产品推荐
相关产品推荐

