【问题标题】:Getting hold of Kafka Stream in process() method在 process() 方法中获取 Kafka Stream
【发布时间】:2021-01-29 16:09:58
【问题描述】:

我正在开发一个 POC 来演示一些 Spring Cloud Kafka 功能。我正在使用带有 SpringBoot 2.4.1 的 Java 11。我的 build.gradle 有以下库

    implementation 'org.apache.kafka:kafka-streams'
    implementation 'org.springframework.cloud:spring-cloud-stream'
    implementation 'org.springframework.cloud:spring-cloud-stream-binder-kafka-streams'
    testImplementation 'org.springframework.boot:spring-boot-starter-test'

我已经创建了一个如下图所示的SpringBoot组件

public class FirstConsumer implements Processor<String, AirSupportRequest> {
    private static Logger logger = LoggerFactory.getLogger(FirstConsumer.class);
    @Override
    public void init(ProcessorContext context) {

    }

    @Override
    public void process(String key, LocationSupport value) {
        logger.info("Key : " + key + value.toString());
        Predicate<String, LocationSupport> northAmerica = (k, v) -> v.getLocation().equalsIgnoreCase("NORTHAMERICA");
        Predicate<String, LocationSupport> asia = (k, v) -> v.getLocation().equalsIgnoreCase("ASIA");

        KStream<String, LocationSupport>[] branches = <How do I get reference to the stream instantiated by Spring Cloud Stream>

我的 application.yml 有以下代码

spring:
  cloud:
    stream:
      kafka:
        streams:
          binder:
            functions:
              process:
                applicationId: asr
      bindings:
        process-in-0:
          destination: location_stream

根据我的配置,Spring Cloud 已经实例化了一个流。我的问题是如何在我的代码中引用该流来过滤和创建流分支。

谢谢!

【问题讨论】:

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


    【解决方案1】:

    您在上面显示的处理器是 Kafka Streams 中低级处理器 API 的实现。使用 Kafka Streams binder 建立绑定的方式是提供一个函数式 bean。例如,

    @Bean
    public Consumer<KStream<String, AirSupportRequest>> proess() {
    
      return ks -> {}
    }
    

    然后process-in-0.destination从Kafka主题消费,并将数据交给消费者。

    您可以将处理器 API 与 DSL 方法混合搭配。请参阅文档中的 this section

    以下是 Kafka Streams 中 branching 的一些信息,尤其是在 Spring Cloud Stream 中使用它时。

    【讨论】:

    • 将高级 DSL 与低级 API 混合的参考很有帮助。我之前错过了那部分。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-03-16
    • 2018-05-14
    • 1970-01-01
    • 2023-04-05
    • 1970-01-01
    • 2017-11-24
    • 1970-01-01
    相关资源
    最近更新 更多