【问题标题】:Sping Cloud Stream Kafka - Consume messages in batch mode and send out as single processed messagesSpring Cloud Stream Kafka - 以批处理模式消费消息并作为单个处理的消息发送
【发布时间】:2019-11-08 11:45:12
【问题描述】:

我尝试为下一个 Spring Cloud Stream 版本准备我们的应用程序。 (目前使用 3.0.0.RC1)。 使用 Kafka 活页夹。

现在我们收到一条消息,对其进行处理并将其重新发送到另一个主题。单独处理每条消息会导致对我们数据库的大量单个请求。

在 3.0.0 版本中,我们希望将消息作为批量处理,因此我们可以在批量更新中保存数据。

在当前版本中我们使用@EnableBinding、@StreamListener

@StreamListener( ExchangeableItemProcessor.STOCK_INPUT )
public void processExchangeableStocks( final ExchangeableStock item ) {
    publishItems( exchangeableItemProcessor.stocks(), articleService.updateStockInformation( Collections.singletonList( item ) ) );
}

void publishItems( final MessageChannel messageChannel, final List<? extends ExchangeableItem> item ) {
    for ( final ExchangeableItem exchangeableItem : item ) {
        final Message<ExchangeableItem> message = MessageBuilder.withPayload( exchangeableItem )
                            .setHeader( "partitionKey", exchangeableItem.getId() )
                            .build();
        messageChannel.send( message )
    }
}

我已将使用者属性设置为“批处理模式”并将签名更改为List&lt;&gt;,但这样做会导致收到List&lt;byte[]&gt; 而不是预期的List&lt;ExchangeableStock&gt;。 Ofc 之后可以进行转换,但这感觉就像“meh”,我认为这应该在调用 Listener 之前发生。

然后我尝试了(新的)功能版本,并且消费工作正常。 我也喜欢这种简单的处理方式

@Bean
public Function<List<ExchangeableStock>, List<ExchangeableStock>> stocks() {
    return articleService::updateStockInformation;
}

但是输出主题现在接收到一个对象列表作为一条消息,并且后续消费者无法正常工作。

我想我错过了什么……

我是否需要添加某种 MessageConverter(对于注释驱动版本)或者是否有办法通过功能版本实现所需的行为?

【问题讨论】:

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


    【解决方案1】:

    IIRC,批处理模式仅支持函数。

    你不能像现在在 StreamListener 中那样使用Consumer&lt;List&lt; ExchangeableStock&gt;&gt; 并将消息发送到频道吗?

    【讨论】:

    • 我试过了,但是当我添加@EnableBinding( ExchangeableItemProcessor.class ) 时,应用程序无法启动。 org.springframework.beans.factory.NoSuchBeanDefinitionException: No qualifying bean of type 'org.springframework.cloud.stream.binding.MessageConverterConfigurer' available: expected at least 1 bean which qualifies as autowire candidate. Dependency annotations: {}
    • 如果我删除功能样式,应用程序将使用相同的绑定接口运行,但是我有 List&lt;byte[]&gt; 消息。
    【解决方案2】:

    我已经成功了:

    @Bean
    @Measure
    public Consumer<List<ExchangeableStock>> stocks() {
        return items -> {
            for ( final ExchangeableStock exchangeableItem : articleService. updateStockInformation( items ) ) {
                final Message<?> message = MessageBuilder.withPayload( exchangeableItem )
                                .setHeader( "partitionKey", exchangeableItem.getId() )
                                .setHeader( KafkaHeaders.TOPIC, "stocks-stg" )
                                .build();
    
                processor.onNext( message );
            }
        };
    }
    
    private final TopicProcessor<Message<?>> processor = TopicProcessor.create();
    
    @Bean
    @Measure
    public Supplier<Flux<?>> source() {
        return () -> processor;
    }
    

    但是动态目的地解析对我不起作用。 我尝试使用KafkaHeaders.TOPICspring.cloud.stream.sendto.destination 作为标头,并设置Kafka 绑定生产者属性use-topic-header: true(用于绑定source-out-0

    如果我为source-out-0 设置目标,它可以工作,但这样做会导致很多TopicProceessors 和Suppliers - 我们有大约10 种不同的消息类型。

    也许我错过了一些小东西来让动态目的地解析工作......

    【讨论】:

      猜你喜欢
      • 2021-04-27
      • 2021-06-12
      • 2017-02-09
      • 2019-12-15
      • 2019-06-09
      • 1970-01-01
      • 2021-02-20
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多