【问题标题】:Pause Kafka Consumer with spring-cloud-stream and Functional Style使用 spring-cloud-stream 和功能样式暂停 Kafka Consumer
【发布时间】:2021-02-09 16:56:31
【问题描述】:

我正在尝试为我的 kafka 流应用程序实现重试机制。我的想法是我会从输入主题中获取消费者和分区 ID 以及主题名称,然后在存储在有效负载中的持续时间内暂停消费者。

我搜索了文档和示例,但我发现的都是基于 spring-cloud-stream 提供的经典绑定的示例。我正在尝试查看是否有办法以功能样式访问这些信息。

例如,下面的代码可以让我以经典的绑定样式访问消费者。

@StreamListener(Sink.INPUT)
public void in(String in, @Header(KafkaHeaders.CONSUMER) Consumer<?, ?> consumer) {
    System.out.println(in);
    consumer.pause(Collections.singleton(new TopicPartition("myTopic", 0)));
}

如何获得与 Functional Style 的等价性?

我尝试使用以下代码,但出现异常,提示找不到此类绑定。

@Bean
public Function<Message<?>, KStream<String, String>> process() {
    message -> {
        Consumer<?, ?> consumer = message.getHeaders().get(KafkaHeaders.Consumer, Consumer.class);
        String topic = message.getHeaders().get(KafkaHeaders.Topic, String.class);
        Integer partitionId = message.getHeaders().get(KafkaHeaders.RECEIVED_PARTITION_ID, Integer.class);
        CustomPayload payload = (CustomPayload) message.getPayload();
        if (payload.getRetryTime() < System.currentTimeMillis()) {
            consumer.pause(Collections.singleton(new TopicPartition(topic, partitionId)));
        }
    }
}

我得到了异常

Caused by: java.lang.IllegalStateException: No factory found for binding target type: org.springframework.messaging.Message among registered factories: channelFactory,messageSourceFactory,kStreamBoundElementFactory,kTableBoundElementFactory,globalKTableBoundElementFactory
    at org.springframework.cloud.stream.binding.AbstractBindableProxyFactory.getBindingTargetFactory(AbstractBindableProxyFactory.java:82)
    at org.springframework.cloud.stream.binder.kafka.streams.function.KafkaStreamsBindableProxyFactory.bindInput(KafkaStreamsBindableProxyFactory.java:191)
    at org.springframework.cloud.stream.binder.kafka.streams.function.KafkaStreamsBindableProxyFactory.afterPropertiesSet(KafkaStreamsBindableProxyFactory.java:111)
    at org.springframework.beans.factory.support.AbstractAutowireCapableBeanFactory.invokeInitMethods(AbstractAutowireCapableBeanFactory.java:1853)
    at org.springframework.beans.factory.support.AbstractAutowireCapableBeanFactory.initializeBean(AbstractAutowireCapableBeanFactory.java:1790)
    ... 96 more

【问题讨论】:

    标签: apache-kafka functional-programming spring-cloud-stream


    【解决方案1】:

    在您的函数式 bean 示例中,您混合了 Message 和 KStream。这就是该特定例外的原因。功能bean可以重写如下。

    @Bean
    public java.util.function.Consumer<Message<?>> process() {
        return message -> {
            Consumer<?, ?> consumer = message.getHeaders().get(KafkaHeaders.Consumer, Consumer.class);
            String topic = message.getHeaders().get(KafkaHeaders.Topic, String.class);
            Integer partitionId = message.getHeaders().get(KafkaHeaders.RECEIVED_PARTITION_ID, Integer.class);
            CustomPayload payload = (CustomPayload) message.getPayload();
            if (payload.getRetryTime() < System.currentTimeMillis()) {
                consumer.pause(Collections.singleton(new TopicPartition(topic, partitionId)));
            }
        }
    }
    

    【讨论】:

    • 显然就像你提到的那样,我不能混合使用 Message 和 KStream。这现在就像一个魅力。非常感谢您的帮助!
    猜你喜欢
    • 2019-04-27
    • 2019-07-23
    • 2020-04-10
    • 1970-01-01
    • 1970-01-01
    • 2022-01-06
    • 2020-10-24
    • 2018-04-28
    • 1970-01-01
    相关资源
    最近更新 更多