From eab86f81a8b34569f8075c352b4bb9602150fbc5 Mon Sep 17 00:00:00 2001 From: "WOODMARK\\r.wiedmann.extern" Date: Wed, 20 Mar 2024 08:05:10 +0100 Subject: [PATCH] GH-2922: Timestamp extractor - Kafka Streams 3.7.0 * Address immutability changes for the call to `Consumed#withTimestampExtractor`. * In 3.7.0, this call returns a new instance of `Consumed` as oppposed to mutating the existing instance in the previous versions. Address this change in behavior in the Kafka Streams binder. Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2922 --- .../streams/AbstractKafkaStreamsBinderProcessor.java | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java index 3be5102e9..2f289d316 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java @@ -1,5 +1,5 @@ /* - * Copyright 2019-2023 the original author or authors. + * Copyright 2019-2024 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -90,6 +90,7 @@ import org.springframework.util.StringUtils; /** * @author Soby Chacko + * @author Ralf Wiedmann * @since 3.0.0 */ public abstract class AbstractKafkaStreamsBinderProcessor implements ApplicationContextAware { @@ -622,13 +623,13 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application timestampExtractor = applicationContext.getBean(kafkaStreamsConsumerProperties.getTimestampExtractorBeanName(), TimestampExtractor.class); } - final Consumed consumed = Consumed.with(keySerde, valueSerde) + Consumed consumed = Consumed.with(keySerde, valueSerde) .withOffsetResetPolicy(autoOffsetReset); if (timestampExtractor != null) { - consumed.withTimestampExtractor(timestampExtractor); + consumed = consumed.withTimestampExtractor(timestampExtractor); } if (StringUtils.hasText(kafkaStreamsConsumerProperties.getConsumedAs())) { - consumed.withName(kafkaStreamsConsumerProperties.getConsumedAs()); + consumed = consumed.withName(kafkaStreamsConsumerProperties.getConsumedAs()); } return consumed; }