【发布时间】:2017-06-15 17:34:50
【问题描述】:
我正在尝试使用带有自动确认功能的 Reactor Kafka 来实现 Kafka 主题分区的并发处理。这里的文档使这看起来是可能的:
http://projectreactor.io/docs/kafka/milestone/reference/#concurrent-ordered
这与我正在尝试的唯一区别是我使用的是自动确认。
我有如下代码(相关方法为receiveAuto):
public class KafkaFluxFactory<K, V> {
private final Map<String, Object> properties;
public KafkaFluxFactory(Map<String, Object> properties) {
this.properties = properties;
}
public Flux<ConsumerRecord<K, V>> receiveAuto(Collection<String> topics, Scheduler scheduler) {
return KafkaReceiver.create(ReceiverOptions.create(properties).subscription(topics))
.receiveAutoAck()
.flatMap(flux -> flux.groupBy(this::extractTopicPartition))
.flatMap(topicPartitionFlux -> topicPartitionFlux.publishOn(scheduler));
}
private TopicPartition extractTopicPartition(ConsumerRecord<K, V> record) {
return new TopicPartition(record.topic(), record.partition());
}
}
当我使用它通过并行调度程序 (Schedulers.newParallel("debug", 10)) 从 Kafka 创建消费者记录通量时,我看到它们最终都在同一个线程上得到处理。
对我可能做错了什么有什么想法吗?
【问题讨论】:
标签: apache-kafka rx-java reactive-programming kafka-consumer-api project-reactor