【问题标题】:How to forward records to multiple Kafka Stream child Processors at once?如何一次将记录转发到多个 Kafka Stream 子处理器?
【发布时间】:2019-09-20 06:11:37
【问题描述】:

在 Kafka Stream API 中,是否可以一次将多个记录转发到不同的子处理器?例如,假设我们有一个称为 Processor-Parent 的父处理器和两个子处理器 Child-1、Child-2。

当 Processor-Parent 收到要处理的记录时,我想执行以下操作。

new_record = create_new_record(current_record)
context.forward(new_record, To(Child-1))
context.forward(old_record, To(Child-2))

这样转发记录是个好习惯吗?

【问题讨论】:

    标签: java apache-kafka apache-kafka-streams stream-processing


    【解决方案1】:

    这不是最佳做法。相反,使用一个父处理器和多个子处理器创建您的拓扑。

    builder = new TopologyBuilder();
        builder.addSource(SOURCE, kafkaTopic)
    .addProcessor("child1", () -> new child1(),SOURCE)
    .addProcessor("child2", () -> new child2(),SOURCE);
    

    通过这种方式,kafka 流确保到达源的每条消息都到达两个子处理器。

    【讨论】:

    • 看来你的问题搞错了。整个想法是使用不同的处理器以不同的方式处理消息。作为一个例子,你可以想象多模式主题和消费者,它们应该根据事件类型处理不同的数据。
    【解决方案2】:

    这取决于您的要求:

    • 如果您的逻辑很简单,您甚至可以使用 Kafka Streams DSL。

    • 如果它稍微复杂一点,并且您需要 Procesor API,但您想将相同的记录传递给两个处理器,您可以像 @Sameer Killamsetty 提到的那样做。

    builder = new TopologyBuilder();
        builder.addSource(SOURCE, kafkaTopic)
    .addProcessor("child1", () -> new child1(), SOURCE)
    .addProcessor("child2", () -> new child2(), SOURCE);
    
    • 如果它更复杂并且取决于处理器中的某些逻辑,您希望将消息传递到不同的处理器节点,您可以这样做。
    builder = new TopologyBuilder();
        builder.addSource(SOURCE, kafkaTopic)
    .addProcessor("InputProcessor", () -> new InputProcessor(), SOURCE)
    .addProcessor("child1", () -> new child1(), "InputProcessor")
    .addProcessor("child2", () -> new child2(), "InputProcessor");
    
    public class InputProcessor extends AbstractProcessor<String, String> {
        @Override
        public void process(String key, String value) {
            try {
                context().forward(key, Integer.parseInt(value), To.child("child1"));
                context().forward(key, value, To.child("child2"));
            }
            catch (NumberFormatException nfe) {
                context().forward(key, value, To.child("child2"));
            }
        }
    }
    

    【讨论】:

      猜你喜欢
      • 2015-12-10
      • 2016-03-14
      • 2017-07-22
      • 1970-01-01
      • 1970-01-01
      • 2016-04-28
      • 1970-01-01
      • 2019-03-15
      • 1970-01-01
      相关资源
      最近更新 更多