【问题标题】:How do I concurrently process Reactor Kafka Streams by Topic and Partition with Auto Acknowledgement?如何通过自动确认按主题和分区同时处理 Reactor Kafka 流?
【发布时间】: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


    【解决方案1】:

    经过相当多的反复试验以及对我想要完成的任务的重新思考后,我意识到我试图用一段代码解决两个问题。

    我需要的两件事是:

    1. 按顺序处理 Kafka 分区
    2. 能够并行处理每个分区

    在尝试使用这段代码解决这两个问题时,我限制了下游用户配置并行化级别的能力。因此,我更改了方法以返回 GroupedFluxes 的 Flux,它为下游用户提供了确定可并行化的正确粒度:

    public Flux<GroupedFlux<TopicPartition, ConsumerRecord<K, V>>> receiveAuto(Collection<String> topics) {
        return KafkaReceiver.create(createReceiverOptions(topics))
            .receiveAutoAck()
            .flatMap(flux -> flux.groupBy(this::extractTopicPartition));
    }
    

    在下游,用户可以使用他们希望的任何调度器并行化每个发出的 GroupedFlux:

    public <V> void work(Flux<GroupedFlux<TopicPartition, V>> flux) {
        flux.doOnNext(groupPublisher -> groupPublisher
                .publishOn(Schedulers.elastic())
                .subscribe(this::doWork))
            .subscribe();
    }
    

    这具有按顺序处理每个 TopicPartition-GroupedFlux 并与其他 GroupedFlux 并行处理的所需行为。

    【讨论】:

      【解决方案2】:

      我猜它至少在你的消费者中是按顺序执行的。要进行并行消耗,您应该将通量转换为ParallelFlux

      public ParallelFlux<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.parallel().runOn(Schedulers.parallel()));
          }
      

      在你的消费者函数之后,如果你想以并行方式消费,你应该使用如下方法:

      void subscribe(Consumer<? super T> onNext, Consumer<? super Throwable>
                  onError, Runnable onComplete, Consumer<? super Subscription> onSubscribe)
      

      或任何其他带有Consumer&lt;T super T&gt; onNext 参数的重载方法。 如果您只使用以下方法,您将按顺序消耗通量

      void subscribe(Subscriber<? super T> s)
      

      【讨论】:

      • @Mr.E.Gas 你也可以在这里查看我的答案stackoverflow.com/a/44588357/2055854 以更好地了解Flux 的工作原理
      • 感谢您对此的意见;它并没有完全解决我的问题,但我找到了一种多方面的答案来解决我想要完成的事情,我会用它来回答我的问题。您的回答不能完全解决我的问题的原因是它最终并行处理了 GroupedFlux,因此可能导致无序处理,这对我的用例来说是不可取的。
      猜你喜欢
      • 1970-01-01
      • 2016-04-18
      • 2016-10-15
      • 1970-01-01
      • 2019-12-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多