From e8d202404bc79fbd68b5cf6254ed7223a6d73cc9 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 23 Oct 2019 15:27:39 -0400 Subject: [PATCH] Custom timestamp extractor per binding Currenlty there is no way to pass a custom timestamp extractor per consumer binding in Kafka Streams binder. Adding this ability. Resolves #640 --- .../AbstractKafkaStreamsBinderProcessor.java | 62 +++++++++++------ .../KafkaStreamsFunctionProcessor.java | 2 +- ...StreamListenerSetupMethodOrchestrator.java | 7 +- .../KafkaStreamsConsumerProperties.java | 12 ++++ .../StreamToGlobalKTableFunctionTests.java | 69 +++++++++++++++++++ 5 files changed, 127 insertions(+), 25 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java index 2c1bb1b5b..dc966e035 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java @@ -32,6 +32,7 @@ import org.apache.kafka.streams.kstream.GlobalKTable; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.KTable; import org.apache.kafka.streams.kstream.Materialized; +import org.apache.kafka.streams.processor.TimestampExtractor; import org.apache.kafka.streams.state.KeyValueStore; import org.apache.kafka.streams.state.StoreBuilder; @@ -123,7 +124,7 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application if (parameterType.isAssignableFrom(KTable.class)) { String materializedAs = extendedConsumerProperties.getMaterializedAs(); String bindingDestination = this.bindingServiceProperties.getBindingDestination(input); - KTable table = getKTable(streamsBuilder, keySerde, valueSerde, materializedAs, + KTable table = getKTable(extendedConsumerProperties, streamsBuilder, keySerde, valueSerde, materializedAs, bindingDestination, autoOffsetReset); KTableBoundElementFactory.KTableWrapper kTableWrapper = (KTableBoundElementFactory.KTableWrapper) targetBean; @@ -135,7 +136,7 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application else if (parameterType.isAssignableFrom(GlobalKTable.class)) { String materializedAs = extendedConsumerProperties.getMaterializedAs(); String bindingDestination = this.bindingServiceProperties.getBindingDestination(input); - GlobalKTable table = getGlobalKTable(streamsBuilder, keySerde, valueSerde, materializedAs, + GlobalKTable table = getGlobalKTable(extendedConsumerProperties, streamsBuilder, keySerde, valueSerde, materializedAs, bindingDestination, autoOffsetReset); GlobalKTableBoundElementFactory.GlobalKTableWrapper globalKTableWrapper = (GlobalKTableBoundElementFactory.GlobalKTableWrapper) targetBean; @@ -244,16 +245,15 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application } } - protected KStream getKStream(String inboundName, BindingProperties bindingProperties, StreamsBuilder streamsBuilder, - Serde keySerde, Serde valueSerde, Topology.AutoOffsetReset autoOffsetReset) { + protected KStream getKStream(String inboundName, BindingProperties bindingProperties, KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, + StreamsBuilder streamsBuilder, Serde keySerde, Serde valueSerde, Topology.AutoOffsetReset autoOffsetReset) { addStateStoreBeans(streamsBuilder); String[] bindingTargets = StringUtils.commaDelimitedListToStringArray( this.bindingServiceProperties.getBindingDestination(inboundName)); - + final Consumed consumed = getConsumed(kafkaStreamsConsumerProperties, keySerde, valueSerde, autoOffsetReset); KStream stream = streamsBuilder.stream(Arrays.asList(bindingTargets), - Consumed.with(keySerde, valueSerde) - .withOffsetResetPolicy(autoOffsetReset)); + consumed); final boolean nativeDecoding = this.bindingServiceProperties .getConsumerProperties(inboundName).isUseNativeDecoding(); if (nativeDecoding) { @@ -305,9 +305,11 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application } private KTable materializedAs(StreamsBuilder streamsBuilder, String destination, String storeName, - Serde k, Serde v, Topology.AutoOffsetReset autoOffsetReset) { + Serde k, Serde v, Topology.AutoOffsetReset autoOffsetReset, KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties) { + + final Consumed consumed = getConsumed(kafkaStreamsConsumerProperties, k, v, autoOffsetReset); return streamsBuilder.table(this.bindingServiceProperties.getBindingDestination(destination), - Consumed.with(k, v).withOffsetResetPolicy(autoOffsetReset), getMaterialized(storeName, k, v)); + consumed, getMaterialized(storeName, k, v)); } private Materialized> getMaterialized( @@ -318,32 +320,50 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application private GlobalKTable materializedAsGlobalKTable( StreamsBuilder streamsBuilder, String destination, String storeName, - Serde k, Serde v, Topology.AutoOffsetReset autoOffsetReset) { + Serde k, Serde v, Topology.AutoOffsetReset autoOffsetReset, KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties) { + final Consumed consumed = getConsumed(kafkaStreamsConsumerProperties, k, v, autoOffsetReset); return streamsBuilder.globalTable( this.bindingServiceProperties.getBindingDestination(destination), - Consumed.with(k, v).withOffsetResetPolicy(autoOffsetReset), + consumed, getMaterialized(storeName, k, v)); } - private GlobalKTable getGlobalKTable(StreamsBuilder streamsBuilder, + private GlobalKTable getGlobalKTable(KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, + StreamsBuilder streamsBuilder, Serde keySerde, Serde valueSerde, String materializedAs, String bindingDestination, Topology.AutoOffsetReset autoOffsetReset) { + final Consumed consumed = getConsumed(kafkaStreamsConsumerProperties, keySerde, valueSerde, autoOffsetReset); return materializedAs != null ? materializedAsGlobalKTable(streamsBuilder, bindingDestination, - materializedAs, keySerde, valueSerde, autoOffsetReset) + materializedAs, keySerde, valueSerde, autoOffsetReset, kafkaStreamsConsumerProperties) : streamsBuilder.globalTable(bindingDestination, - Consumed.with(keySerde, valueSerde) - .withOffsetResetPolicy(autoOffsetReset)); + consumed); } - private KTable getKTable(StreamsBuilder streamsBuilder, Serde keySerde, - Serde valueSerde, String materializedAs, String bindingDestination, - Topology.AutoOffsetReset autoOffsetReset) { + private KTable getKTable(KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, + StreamsBuilder streamsBuilder, Serde keySerde, + Serde valueSerde, String materializedAs, String bindingDestination, + Topology.AutoOffsetReset autoOffsetReset) { + final Consumed consumed = getConsumed(kafkaStreamsConsumerProperties, keySerde, valueSerde, autoOffsetReset); return materializedAs != null ? materializedAs(streamsBuilder, bindingDestination, materializedAs, - keySerde, valueSerde, autoOffsetReset) + keySerde, valueSerde, autoOffsetReset, kafkaStreamsConsumerProperties) : streamsBuilder.table(bindingDestination, - Consumed.with(keySerde, valueSerde) - .withOffsetResetPolicy(autoOffsetReset)); + consumed); + } + + private Consumed getConsumed(KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, + Serde keySerde, Serde valueSerde, Topology.AutoOffsetReset autoOffsetReset) { + TimestampExtractor timestampExtractor = null; + if (kafkaStreamsConsumerProperties.getTimestampExtractorBeanName() != null) { + timestampExtractor = applicationContext.getBean(kafkaStreamsConsumerProperties.getTimestampExtractorBeanName(), + TimestampExtractor.class); + } + final Consumed consumed = Consumed.with(keySerde, valueSerde) + .withOffsetResetPolicy(autoOffsetReset); + if (timestampExtractor != null) { + consumed.withTimestampExtractor(timestampExtractor); + } + return consumed; } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java index 6408c564d..4d1e4ac50 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java @@ -311,7 +311,7 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro final Topology.AutoOffsetReset autoOffsetReset = getAutoOffsetReset(input, extendedConsumerProperties); if (parameterType.isAssignableFrom(KStream.class)) { - KStream stream = getKStream(input, bindingProperties, streamsBuilder, keySerde, valueSerde, autoOffsetReset); + KStream stream = getKStream(input, bindingProperties, extendedConsumerProperties, streamsBuilder, keySerde, valueSerde, autoOffsetReset); KStreamBoundElementFactory.KStreamWrapper kStreamWrapper = (KStreamBoundElementFactory.KStreamWrapper) targetBean; //wrap the proxy created during the initial target type binding with real object (KStream) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java index e07f5da36..50a5d096c 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java @@ -270,7 +270,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator extends AbstractKafkaStr if (parameterType.isAssignableFrom(KStream.class)) { KStream stream = getkStream(inboundName, spec, - bindingProperties, streamsBuilder, keySerde, valueSerde, + bindingProperties, extendedConsumerProperties, streamsBuilder, keySerde, valueSerde, autoOffsetReset); KStreamBoundElementFactory.KStreamWrapper kStreamWrapper = (KStreamBoundElementFactory.KStreamWrapper) targetBean; // wrap the proxy created during the initial target type binding @@ -361,7 +361,8 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator extends AbstractKafkaStr private KStream getkStream(String inboundName, KafkaStreamsStateStoreProperties storeSpec, - BindingProperties bindingProperties, StreamsBuilder streamsBuilder, + BindingProperties bindingProperties, + KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, StreamsBuilder streamsBuilder, Serde keySerde, Serde valueSerde, Topology.AutoOffsetReset autoOffsetReset) { if (storeSpec != null) { @@ -371,7 +372,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator extends AbstractKafkaStr LOG.info("state store " + storeBuilder.name() + " added to topology"); } } - return getKStream(inboundName, bindingProperties, streamsBuilder, keySerde, valueSerde, autoOffsetReset); + return getKStream(inboundName, bindingProperties, kafkaStreamsConsumerProperties, streamsBuilder, keySerde, valueSerde, autoOffsetReset); } private void validateStreamListenerMethod(StreamListener streamListener, diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsConsumerProperties.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsConsumerProperties.java index 717ef0da7..debfd169c 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsConsumerProperties.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsConsumerProperties.java @@ -43,6 +43,11 @@ public class KafkaStreamsConsumerProperties extends KafkaConsumerProperties { */ private String materializedAs; + /** + * {@link org.apache.kafka.streams.processor.TimestampExtractor} bean name to use for this consumer. + */ + private String timestampExtractorBeanName; + public String getApplicationId() { return this.applicationId; } @@ -75,4 +80,11 @@ public class KafkaStreamsConsumerProperties extends KafkaConsumerProperties { this.materializedAs = materializedAs; } + public String getTimestampExtractorBeanName() { + return timestampExtractorBeanName; + } + + public void setTimestampExtractorBeanName(String timestampExtractorBeanName) { + this.timestampExtractorBeanName = timestampExtractorBeanName; + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java index 7af1eddf6..84648e36b 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java @@ -32,12 +32,17 @@ import org.apache.kafka.common.serialization.LongSerializer; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.kstream.GlobalKTable; import org.apache.kafka.streams.kstream.KStream; +import org.apache.kafka.streams.kstream.KTable; +import org.apache.kafka.streams.processor.TimestampExtractor; +import org.apache.kafka.streams.processor.WallclockTimestampExtractor; import org.junit.ClassRule; import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBindingProperties; +import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; @@ -176,6 +181,56 @@ public class StreamToGlobalKTableFunctionTests { } } + @Test + public void testTimeExtractor() throws Exception { + SpringApplication app = new SpringApplication(OrderEnricherApplication.class); + app.setWebApplicationType(WebApplicationType.NONE); + + try (ConfigurableApplicationContext context = app.run( + "--server.port=0", + "--spring.jmx.enabled=false", + "--spring.cloud.stream.function.definition=forTimeExtractorTest", + "--spring.cloud.stream.bindings.forTimeExtractorTest-in-0.destination=orders", + "--spring.cloud.stream.bindings.forTimeExtractorTest-in-1.destination=customers", + "--spring.cloud.stream.bindings.forTimeExtractorTest-in-2.destination=products", + "--spring.cloud.stream.bindings.forTimeExtractorTest-out-0.destination=enriched-order", + "--spring.cloud.stream.kafka.streams.bindings.forTimeExtractorTest-in-0.consumer.timestampExtractorBeanName" + + "=timestampExtractor", + "--spring.cloud.stream.kafka.streams.bindings.forTimeExtractorTest-in-1.consumer.timestampExtractorBeanName" + + "=timestampExtractor", + "--spring.cloud.stream.kafka.streams.bindings.forTimeExtractorTest-in-2.consumer.timestampExtractorBeanName" + + "=timestampExtractor", + "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=10000", + "--spring.cloud.stream.kafka.streams.bindings.order.consumer.applicationId=" + + "testTimeExtractor-abc", + "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) { + + final KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties = + context.getBean(KafkaStreamsExtendedBindingProperties.class); + + final Map bindings = kafkaStreamsExtendedBindingProperties.getBindings(); + + final KafkaStreamsBindingProperties kafkaStreamsBindingProperties0 = bindings.get("forTimeExtractorTest-in-0"); + final String timestampExtractorBeanName0 = kafkaStreamsBindingProperties0.getConsumer().getTimestampExtractorBeanName(); + final TimestampExtractor timestampExtractor0 = context.getBean(timestampExtractorBeanName0, TimestampExtractor.class); + assertThat(timestampExtractor0).isNotNull(); + + final KafkaStreamsBindingProperties kafkaStreamsBindingProperties1 = bindings.get("forTimeExtractorTest-in-1"); + final String timestampExtractorBeanName1 = kafkaStreamsBindingProperties1.getConsumer().getTimestampExtractorBeanName(); + final TimestampExtractor timestampExtractor1 = context.getBean(timestampExtractorBeanName1, TimestampExtractor.class); + assertThat(timestampExtractor1).isNotNull(); + + final KafkaStreamsBindingProperties kafkaStreamsBindingProperties2 = bindings.get("forTimeExtractorTest-in-2"); + final String timestampExtractorBeanName2 = kafkaStreamsBindingProperties2.getConsumer().getTimestampExtractorBeanName(); + final TimestampExtractor timestampExtractor2 = context.getBean(timestampExtractorBeanName2, TimestampExtractor.class); + assertThat(timestampExtractor2).isNotNull(); + } + } + @EnableAutoConfiguration public static class OrderEnricherApplication { @@ -204,6 +259,20 @@ public class StreamToGlobalKTableFunctionTests { ) ); } + + @Bean + public Function, + Function, + Function, KStream>>> forTimeExtractorTest() { + return orderStream -> + customers -> + products -> orderStream; + } + + @Bean + public TimestampExtractor timestampExtractor() { + return new WallclockTimestampExtractor(); + } } static class Order {