【问题标题】:How to do this topology in Spring Cloud Kafka Streams in function style?如何在 Spring Cloud Kafka Streams 中以函数样式执行此拓扑?
【发布时间】:2021-09-27 17:24:43
【问题描述】:
var streamsBuilder = new StreamsBuilder();
    KStream<String, String> inputStream = streamsBuilder.stream("input_topic");

KStream<String, String> upperCaseString =
        inputStream.mapValues((ValueMapper<String, String>) String::toUpperCase);

upperCaseString.mapValues(v -> v + "postfix").to("with_postfix_topic");
upperCaseString.mapValues(v -> "prefix" + v).to("with_prefix_topic");

Topology topology = streamsBuilder.build();

我可以写三个函数bean。第一个 bean 将大小写并将结果写入某个主题('upper_case_topic')。其他bean 将使用这个结果(来自'upper_case_topic')并添加前缀/后缀。但是如何在不写中间主题('upper_case_topic')的情况下做到这一点?

更新: 这是我可能的解决方案:

@Bean
public Consumer<KStream<String, String>> process() {
    return input -> {
        KStream<String, String> upperCaseStream =
                input.mapValues((ValueMapper<String, String>) String::toUpperCase);

        upperCaseStream.mapValues(v -> v + " 111").to("new_topic_1");

        upperCaseStream.mapValues(v -> v + " 222").to("new_topic_2");
    };
}

【问题讨论】:

    标签: java apache-kafka-streams spring-kafka spring-cloud-stream spring-cloud-stream-binder-kafka


    【解决方案1】:

    这里有一些您可以尝试的选项。

    选项 1 - 使用 Spring Cloud Stream 中的 StreamBridge

    @Bean
    public Consumer<KStream<String, String>> process() {
    
      KStream<String, String> upperCaseString =
            inputStream.mapValues((ValueMapper<String, String>) 
        String::toUpperCase);
    
      upperCaseString.foreach((key, value) -> {
                    streamBridge.send("with_postfix_topic", value + "postfix");
                    streamBridge.send("with_prefix_topic", "prefix" + value);
                });
    }
    

    上述方法的一方面是您需要来自 Spring Cloud Stream 的 Kafka 和 Kafka Streams 绑定器才能使其工作。另一个问题是,当您在业务逻辑中直接发送到 Kafka 主题时,您会失去 Kafka Streams 原生提供的端到端语义。根据您的用例,这种方法可能没问题。

    选项 2 - 在出站时使用 KStream[]

    您通常使用 Kafka Streams API 的分支功能在出站上使用 KStream[],但您可以利用 Spring Cloud Stream 在分支功能之上构建的输出绑定功能作为您的用例的解决方法。这是一个您可以尝试的想法。

    @Bean
    public Function<KStream<String, String>, KStream<String, String>[]> process() {
      return inputStream -> {
                    KStream<String, String> upperCaseString =
                            inputStream.mapValues((ValueMapper<String, String>)
                                    String::toUpperCase);
                    KStream<String, String>[] kStreams = new KStream[2];
                    kStreams[0] = upperCaseString.mapValues(v -> v + "postfix");
                    kStreams[1] = upperCaseString.mapValues(v -> v + "postfix");
                    return kStreams;
                };
    }
    

    然后你可以定义你的目的地如下:

    spring.cloud.stream.bindings.process-in-0.destination: <input-topic-name>
    spring.cloud.stream.bindings.process-out-0.destination: <output-topic-name>
    spring.cloud.stream.bindings.process-out-1.destination: <output-topic-name>
    

    使用这种方法,您可以从 Kafka Streams 获得端到端语义,因为发送到 Kafka 主题是通过 Kafka Streams 处理的,方法是由 binder 在后台调用 KStream 上的 to 方法。

    选项 3 - 使用函数组合

    另一个选项是 Kafka Streams binder 中的函数组合。请记住,此功能尚未发布,但活页夹的最新 3.1.x/3.2.x 快照具有此功能。有了它,您可以定义如下三个简单的函数。

    @Bean
    public Function<KStream<String, String>, KStream<String, String>> uppercase() {
      return inputStream -> inputStream.mapValues((ValueMapper<String, String>) String::toUpperCase);
    }
    
    @Bean
    public Function<KStream<String, String>, KStream<String, String>> postfixed() {
      return inputStream -> inputStream.mapValues(v -> v + "postfix");
    }
    
    @Bean
    public Function<KStream<String, String>, KStream<String, String>> prefixed() { 
      return inputStream -> inputStream.mapValues(v -> "prefix" + v);
    }
    
    

    然后你可以有如下两个函数组合流程:

    spring.cloud.function.definition: uppercase|postfixed;uppercase|prefixed
    

    您可以在每个组合函数绑定上设置输入主题,如下所示。

    spring.cloud.stream.bindings.uppercasepostfixed-in-0.destination=<input-topic>
    spring.cloud.stream.bindings.uppercaseprefixed-in-0.destination=<input-topic>
    
    

    使用这种函数组合方法,您可以从 Kafka Streams 获得端到端语义,并且避免了额外的中间主题。不过这里的缺点是uppercase 函数将为每个传入记录调用两次。

    上述方法可行,但在将它们用于您的用例之前请考虑权衡。

    【讨论】:

    • 我为您的第一个选项提供了替代方案。你怎么看待这件事? (我更新了我的问题并添加了一个解决方案)。顺便说一句,我认为你的第三个选项是最好的选择。
    • 您的解决方案也有效。您正在使用 Kafka Streams 直接发送到在这种情况下完美的主题。
    猜你喜欢
    • 2017-06-09
    • 1970-01-01
    • 2019-10-23
    • 2021-09-28
    • 2023-03-10
    • 2019-09-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多