【问题标题】:Can not retrieve KafkaStreams object in a Spring Boot application that uses Spring cloud stream binder无法在使用 Spring 云流绑定器的 Spring Boot 应用程序中检索 KafkaStreams 对象
【发布时间】:2020-07-08 13:10:17
【问题描述】:

所以我的问题是我的属性文件中定义了一些 Kafka 主题,我可以从该主题中读取 KafkaStream<String, String> 在我的 SpringBoot 应用程序中没有问题。但我想访问KafkaStreams 对象,以便我可以打印我的KafkaStreams 拓扑,这对开发很有用。 在我的@StreamListener 之一中,我尝试检索stream-builder-process bean,以便我可以通过这种方式获取底层KafkaStreams 对象(如此处所述:https://cloud.spring.io/spring-cloud-stream-binder-kafka/spring-cloud-stream-binder-kafka.html#_accessing_the_underlying_kafkastreams_object)但不幸的是它不起作用。 代码如下:

@StreamListener
public void processEvent(@Input("order-paid-stream") KStream<String, String> inputStream) {
    StreamsBuilderFactoryBean streamsBuilderFactoryBean = applicationContext.getBean("&stream-builder-process", StreamsBuilderFactoryBean.class);
    KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams();
    System.out.println(kafkaStreams.toString());
    inputStream.foreach(this::handleMessage);
}

应用程序启动时,我收到以下消息:

在应用程序启动后出现类似错误(未找到该名称的 bean)后,我还尝试在我的一个 REST 控制器方法上以相同的方式检索 KafkaStreams 对象。

有什么帮助吗?

【问题讨论】:

  • 以前没有人遇到过这个问题吗?

标签: java spring-boot apache-kafka


【解决方案1】:

我遇到了类似的问题,并发现我还需要包含定义流程函数的类名,例如stream-builder-MyStreamProcessor-process.

【讨论】:

    【解决方案2】:

    从延迟的单独线程开始,以确保创建 kafkaStream 对象并完成拓扑。

    streamsBuilderFactoryBean.getKafkaStreams() 返回 KafkaStream 对象 streamBuilderFactoryBean.getSingletonInstance().topology 返回拓扑对象 streamBuilderFactoryBean.getStreamsConfiguration() 返回所有设置的 kafka-streams 配置

    class StreamsListener {
        @StreamListener
        @SendTo("output")
        public KStream<String, String> process(@Input("input') KStream<String,String> rawCloudEventKStream {
    
            new Thread(() -> {
                try { Thread.sleep(30000); } catch (InterruptedException e) { e.printStackTrace(); }
    
                StreamsBuilderFactoryBean streamsBuilderFactoryBean = context.getBean("&stream-builder-StreamsListener-process", StreamsBuilderFactoryBean.class);
                System.out.println("KafkaStreams configs: " + streamsBuilderFactoryBean.getStreamsConfiguration());
                KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams();
            }).start();
    
            ...
        }
    }
    

    【讨论】:

      猜你喜欢
      • 2017-03-20
      • 1970-01-01
      • 2021-08-07
      • 1970-01-01
      • 1970-01-01
      • 2019-05-11
      • 2018-12-11
      • 1970-01-01
      • 2022-08-18
      相关资源
      最近更新 更多