diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java index b4c65f63e..3bdcb1428 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java @@ -31,6 +31,7 @@ import org.springframework.beans.factory.ObjectProvider; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.boot.context.properties.EnableConfigurationProperties; @@ -238,18 +239,23 @@ public class KafkaStreamsBinderSupportAutoConfiguration { cleanupConfig.getIfUnique()); } - public KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor(BindingServiceProperties bindingServiceProperties, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, - KeyValueSerdeResolver keyValueSerdeResolver, KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, - KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate, - ObjectProvider cleanupConfig, - FunctionCatalog functionCatalog, BindableProxyFactory bindableProxyFactory){ + @Bean + @ConditionalOnProperty("spring.cloud.stream.kafka.streams.function.definition") + public KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor(BindingServiceProperties bindingServiceProperties, + KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, + KeyValueSerdeResolver keyValueSerdeResolver, + KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, + KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate, + ObjectProvider cleanupConfig, + FunctionCatalog functionCatalog, BindableProxyFactory bindableProxyFactory) { return new KafkaStreamsFunctionProcessor(bindingServiceProperties, kafkaStreamsExtendedBindingProperties, keyValueSerdeResolver, kafkaStreamsBindingInformationCatalogue, kafkaStreamsMessageConversionDelegate, cleanupConfig.getIfUnique(), functionCatalog, bindableProxyFactory); } @Bean - public KafkaStreamsMessageConversionDelegate messageConversionDelegate(CompositeMessageConverterFactory compositeMessageConverterFactory, + public KafkaStreamsMessageConversionDelegate messageConversionDelegate( + CompositeMessageConverterFactory compositeMessageConverterFactory, SendToDlqAndContinue sendToDlqAndContinue, KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue, KafkaStreamsBinderConfigurationProperties binderConfigurationProperties) { 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 925c35635..76b043401 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 @@ -1,5 +1,5 @@ /* - * Copyright 2019 the original author or authors. + * Copyright 2019-2019 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. @@ -48,6 +48,7 @@ import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.cloud.function.context.FunctionCatalog; +import org.springframework.cloud.function.core.FluxedFunction; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsConsumerProperties; @@ -88,8 +89,10 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { private ConfigurableApplicationContext applicationContext; - public KafkaStreamsFunctionProcessor(BindingServiceProperties bindingServiceProperties, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, - KeyValueSerdeResolver keyValueSerdeResolver, KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, + public KafkaStreamsFunctionProcessor(BindingServiceProperties bindingServiceProperties, + KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, + KeyValueSerdeResolver keyValueSerdeResolver, + KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate, CleanupConfig cleanupConfig, FunctionCatalog functionCatalog, @@ -116,7 +119,8 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { while (iterator.hasNext() && generic != null) { if (generic.getRawClass() != null && - (generic.getRawClass().equals(Function.class) || generic.getRawClass().equals(Consumer.class))) { + (generic.getRawClass().equals(Function.class) || + generic.getRawClass().equals(Consumer.class))) { map.put(iterator.next(), generic.getGeneric(0)); } generic = generic.getGeneric(1); @@ -142,15 +146,22 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { if (resolvableType.getRawClass() != null && resolvableType.getRawClass().equals(Consumer.class)) { Consumer consumer = functionCatalog.lookup(Consumer.class, functionName); consumer.accept(adaptedInboundArguments[0]); - } else { + } + else { Function function = functionCatalog.lookup(Function.class, functionName); + Object target = null; + if (function instanceof FluxedFunction) { + target = ((FluxedFunction) function).getTarget(); + } + function = (Function) target; Object result = function.apply(adaptedInboundArguments[0]); int i = 1; while (result instanceof Function || result instanceof Consumer) { if (result instanceof Function) { result = ((Function) result).apply(adaptedInboundArguments[i]); - } else { + } + else { ((Consumer) result).accept(adaptedInboundArguments[i]); result = null; } @@ -160,7 +171,8 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { if (result.getClass().isArray()) { Assert.isTrue(methodAnnotatedOutboundNames.length == ((Object[]) result).length, "Result does not match with the number of declared outbounds"); - } else { + } + else { Assert.isTrue(methodAnnotatedOutboundNames.length == 1, "Result does not match with the number of declared outbounds"); } @@ -174,7 +186,8 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { boundElement = (KStreamBoundElementFactory.KStreamWrapper) targetBean; boundElement.wrap((KStream) outboundKStream); } - } else { + } + else { Object targetBean = this.applicationContext.getBean(methodAnnotatedOutboundNames[0]); KStreamBoundElementFactory.KStreamWrapper @@ -183,13 +196,15 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { } } } - } catch (Exception ex) { + } + catch (Exception ex) { throw new BeanInitializationException("Cannot setup StreamListener for foobar", ex); } } @SuppressWarnings({"unchecked"}) - private Object[] adaptAndRetrieveInboundArguments(Map stringResolvableTypeMap, String functionName) { + private Object[] adaptAndRetrieveInboundArguments(Map stringResolvableTypeMap, + String functionName) { Object[] arguments = new Object[stringResolvableTypeMap.size()]; int i = 0; for (String input : stringResolvableTypeMap.keySet()) { @@ -206,13 +221,16 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { buildStreamsBuilderAndRetrieveConfig(functionName, applicationContext, input); } try { - StreamsBuilderFactoryBean streamsBuilderFactoryBean = this.methodStreamsBuilderFactoryBeanMap.get(functionName); + StreamsBuilderFactoryBean streamsBuilderFactoryBean = + this.methodStreamsBuilderFactoryBeanMap.get(functionName); StreamsBuilder streamsBuilder = streamsBuilderFactoryBean.getObject(); - KafkaStreamsConsumerProperties extendedConsumerProperties = this.kafkaStreamsExtendedBindingProperties.getExtendedConsumerProperties(input); + KafkaStreamsConsumerProperties extendedConsumerProperties = + this.kafkaStreamsExtendedBindingProperties.getExtendedConsumerProperties(input); //get state store spec //KafkaStreamsStateStoreProperties spec = buildStateStoreSpec(method); Serde keySerde = this.keyValueSerdeResolver.getInboundKeySerde(extendedConsumerProperties); - Serde valueSerde = this.keyValueSerdeResolver.getInboundValueSerde(bindingProperties.getConsumer(), extendedConsumerProperties); + Serde valueSerde = this.keyValueSerdeResolver.getInboundValueSerde( + bindingProperties.getConsumer(), extendedConsumerProperties); final KafkaConsumerProperties.StartOffset startOffset = extendedConsumerProperties.getStartOffset(); Topology.AutoOffsetReset autoOffsetReset = null; @@ -236,18 +254,23 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { if (parameterType.isAssignableFrom(KStream.class)) { KStream stream = getkStream(input, bindingProperties, streamsBuilder, keySerde, valueSerde, autoOffsetReset); - KStreamBoundElementFactory.KStreamWrapper kStreamWrapper = (KStreamBoundElementFactory.KStreamWrapper) targetBean; + KStreamBoundElementFactory.KStreamWrapper kStreamWrapper = + (KStreamBoundElementFactory.KStreamWrapper) targetBean; //wrap the proxy created during the initial target type binding with real object (KStream) kStreamWrapper.wrap((KStream) stream); this.kafkaStreamsBindingInformationCatalogue.addStreamBuilderFactory(streamsBuilderFactoryBean); if (KStream.class.isAssignableFrom(stringResolvableTypeMap.get(input).getRawClass())) { - final Class valueClass = (stringResolvableTypeMap.get(input).getGeneric(1).getRawClass() != null) + final Class valueClass = + (stringResolvableTypeMap.get(input).getGeneric(1).getRawClass() != null) ? (stringResolvableTypeMap.get(input).getGeneric(1).getRawClass()) : Object.class; - if (this.kafkaStreamsBindingInformationCatalogue.isUseNativeDecoding((KStream) kStreamWrapper)) { + if (this.kafkaStreamsBindingInformationCatalogue.isUseNativeDecoding( + (KStream) kStreamWrapper)) { arguments[i] = stream; - } else { - arguments[i] = this.kafkaStreamsMessageConversionDelegate.deserializeOnInbound(valueClass, stream); + } + else { + arguments[i] = this.kafkaStreamsMessageConversionDelegate.deserializeOnInbound( + valueClass, stream); } } @@ -255,69 +278,83 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { arguments[i] = stream; } Assert.notNull(arguments[i], "problems.."); - } else if (parameterType.isAssignableFrom(KTable.class)) { + } + else if (parameterType.isAssignableFrom(KTable.class)) { String materializedAs = extendedConsumerProperties.getMaterializedAs(); String bindingDestination = this.bindingServiceProperties.getBindingDestination(input); KTable table = getKTable(streamsBuilder, keySerde, valueSerde, materializedAs, bindingDestination, autoOffsetReset); - KTableBoundElementFactory.KTableWrapper kTableWrapper = (KTableBoundElementFactory.KTableWrapper) targetBean; + KTableBoundElementFactory.KTableWrapper kTableWrapper = + (KTableBoundElementFactory.KTableWrapper) targetBean; //wrap the proxy created during the initial target type binding with real object (KTable) kTableWrapper.wrap((KTable) table); this.kafkaStreamsBindingInformationCatalogue.addStreamBuilderFactory(streamsBuilderFactoryBean); arguments[i] = table; - } else if (parameterType.isAssignableFrom(GlobalKTable.class)) { + } + else if (parameterType.isAssignableFrom(GlobalKTable.class)) { String materializedAs = extendedConsumerProperties.getMaterializedAs(); String bindingDestination = this.bindingServiceProperties.getBindingDestination(input); GlobalKTable table = getGlobalKTable(streamsBuilder, keySerde, valueSerde, materializedAs, bindingDestination, autoOffsetReset); - GlobalKTableBoundElementFactory.GlobalKTableWrapper globalKTableWrapper = (GlobalKTableBoundElementFactory.GlobalKTableWrapper) targetBean; + GlobalKTableBoundElementFactory.GlobalKTableWrapper globalKTableWrapper = + (GlobalKTableBoundElementFactory.GlobalKTableWrapper) targetBean; //wrap the proxy created during the initial target type binding with real object (KTable) globalKTableWrapper.wrap((GlobalKTable) table); this.kafkaStreamsBindingInformationCatalogue.addStreamBuilderFactory(streamsBuilderFactoryBean); arguments[i] = table; } i++; - } catch (Exception ex) { + } + catch (Exception ex) { throw new IllegalStateException(ex); } - } else { + } + else { throw new IllegalStateException(StreamListenerErrorMessages.INVALID_DECLARATIVE_METHOD_PARAMETERS); } } return arguments; } - private GlobalKTable getGlobalKTable(StreamsBuilder streamsBuilder, Serde keySerde, Serde valueSerde, String materializedAs, + private GlobalKTable getGlobalKTable(StreamsBuilder streamsBuilder, + Serde keySerde, Serde valueSerde, String materializedAs, String bindingDestination, Topology.AutoOffsetReset autoOffsetReset) { return materializedAs != null ? - materializedAsGlobalKTable(streamsBuilder, bindingDestination, materializedAs, keySerde, valueSerde, autoOffsetReset) : + materializedAsGlobalKTable(streamsBuilder, bindingDestination, materializedAs, + keySerde, valueSerde, autoOffsetReset) : streamsBuilder.globalTable(bindingDestination, Consumed.with(keySerde, valueSerde).withOffsetResetPolicy(autoOffsetReset)); } - private KTable getKTable(StreamsBuilder streamsBuilder, Serde keySerde, Serde valueSerde, String materializedAs, + private KTable getKTable(StreamsBuilder streamsBuilder, Serde keySerde, Serde valueSerde, + String materializedAs, String bindingDestination, Topology.AutoOffsetReset autoOffsetReset) { return materializedAs != null ? - materializedAs(streamsBuilder, bindingDestination, materializedAs, keySerde, valueSerde, autoOffsetReset) : + materializedAs(streamsBuilder, bindingDestination, materializedAs, keySerde, valueSerde, + autoOffsetReset) : streamsBuilder.table(bindingDestination, Consumed.with(keySerde, valueSerde).withOffsetResetPolicy(autoOffsetReset)); } - private KTable materializedAs(StreamsBuilder streamsBuilder, String destination, String storeName, Serde k, Serde v, + private KTable materializedAs(StreamsBuilder streamsBuilder, String destination, + String storeName, Serde k, Serde v, Topology.AutoOffsetReset autoOffsetReset) { return streamsBuilder.table(this.bindingServiceProperties.getBindingDestination(destination), Consumed.with(k, v).withOffsetResetPolicy(autoOffsetReset), getMaterialized(storeName, k, v)); } - private GlobalKTable materializedAsGlobalKTable(StreamsBuilder streamsBuilder, String destination, String storeName, Serde k, Serde v, + private GlobalKTable materializedAsGlobalKTable(StreamsBuilder streamsBuilder, + String destination, String storeName, + Serde k, Serde v, Topology.AutoOffsetReset autoOffsetReset) { return streamsBuilder.globalTable(this.bindingServiceProperties.getBindingDestination(destination), Consumed.with(k, v).withOffsetResetPolicy(autoOffsetReset), getMaterialized(storeName, k, v)); } - private Materialized> getMaterialized(String storeName, Serde k, Serde v) { + private Materialized> getMaterialized(String storeName, + Serde k, Serde v) { return Materialized.>as(storeName) .withKeySerde(k) .withValueSerde(v); @@ -334,11 +371,15 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { streamsBuilder.stream(Arrays.asList(bindingTargets), Consumed.with(keySerde, valueSerde) .withOffsetResetPolicy(autoOffsetReset)); - final boolean nativeDecoding = this.bindingServiceProperties.getConsumerProperties(inboundName).isUseNativeDecoding(); + final boolean nativeDecoding = this.bindingServiceProperties.getConsumerProperties(inboundName) + .isUseNativeDecoding(); if (nativeDecoding) { - LOG.info("Native decoding is enabled for " + inboundName + ". Inbound deserialization done at the broker."); - } else { - LOG.info("Native decoding is disabled for " + inboundName + ". Inbound message conversion done by Spring Cloud Stream."); + LOG.info("Native decoding is enabled for " + inboundName + ". " + + "Inbound deserialization done at the broker."); + } + else { + LOG.info("Native decoding is disabled for " + inboundName + ". " + + "Inbound message conversion done by Spring Cloud Stream."); } stream = stream.mapValues((value) -> { @@ -347,7 +388,8 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { if (value != null && !StringUtils.isEmpty(contentType) && !nativeDecoding) { returnValue = MessageBuilder.withPayload(value) .setHeader(MessageHeaders.CONTENT_TYPE, contentType).build(); - } else { + } + else { returnValue = value; } return returnValue; @@ -370,9 +412,11 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { String inboundName) { ConfigurableListableBeanFactory beanFactory = this.applicationContext.getBeanFactory(); - Map streamConfigGlobalProperties = applicationContext.getBean("streamConfigGlobalProperties", Map.class); + Map streamConfigGlobalProperties = applicationContext.getBean("streamConfigGlobalProperties", + Map.class); - KafkaStreamsConsumerProperties extendedConsumerProperties = this.kafkaStreamsExtendedBindingProperties.getExtendedConsumerProperties(inboundName); + KafkaStreamsConsumerProperties extendedConsumerProperties = this.kafkaStreamsExtendedBindingProperties + .getExtendedConsumerProperties(inboundName); streamConfigGlobalProperties.putAll(extendedConsumerProperties.getConfiguration()); String applicationId = extendedConsumerProperties.getApplicationId(); @@ -387,9 +431,11 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { streamConfigGlobalProperties.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, concurrency); } - Map kafkaStreamsDlqDispatchers = applicationContext.getBean("kafkaStreamsDlqDispatchers", Map.class); + Map kafkaStreamsDlqDispatchers = applicationContext.getBean( + "kafkaStreamsDlqDispatchers", Map.class); - KafkaStreamsConfiguration kafkaStreamsConfiguration = new KafkaStreamsConfiguration(streamConfigGlobalProperties) { + KafkaStreamsConfiguration kafkaStreamsConfiguration = + new KafkaStreamsConfiguration(streamConfigGlobalProperties) { @Override public Properties asProperties() { Properties properties = super.asProperties(); @@ -403,10 +449,13 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { : new StreamsBuilderFactoryBean(kafkaStreamsConfiguration, this.cleanupConfig); streamsBuilder.setAutoStartup(false); BeanDefinition streamsBuilderBeanDefinition = - BeanDefinitionBuilder.genericBeanDefinition((Class) streamsBuilder.getClass(), () -> streamsBuilder) + BeanDefinitionBuilder.genericBeanDefinition( + (Class) streamsBuilder.getClass(), () -> streamsBuilder) .getRawBeanDefinition(); - ((BeanDefinitionRegistry) beanFactory).registerBeanDefinition("stream-builder-" + functionName, streamsBuilderBeanDefinition); - StreamsBuilderFactoryBean streamsBuilderX = applicationContext.getBean("&stream-builder-" + functionName, StreamsBuilderFactoryBean.class); + ((BeanDefinitionRegistry) beanFactory).registerBeanDefinition("stream-builder-" + + functionName, streamsBuilderBeanDefinition); + StreamsBuilderFactoryBean streamsBuilderX = applicationContext.getBean("&stream-builder-" + + functionName, StreamsBuilderFactoryBean.class); this.methodStreamsBuilderFactoryBeanMap.put(functionName, streamsBuilderX); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionAutoConfiguration.java index f1b218c55..c352bc361 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionAutoConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionAutoConfiguration.java @@ -1,5 +1,5 @@ /* - * Copyright 2019 the original author or authors. + * Copyright 2019-2019 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. @@ -32,21 +32,22 @@ import org.springframework.context.annotation.Configuration; @EnableConfigurationProperties(KafkaStreamsFunctionProperties.class) public class KafkaStreamsFunctionAutoConfiguration { - @Autowired - private KafkaStreamsFunctionProperties properties; - @Autowired ConfigurableApplicationContext context; @Bean - public KafkaStreamsFunctionProcessorInvoker kafkaStreamsFunctionProcessorInvoker(KafkaStreamsFunctionBeanPostProcessor kafkaStreamsFunctionBeanPostProcessor, - KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor) { - return new KafkaStreamsFunctionProcessorInvoker(kafkaStreamsFunctionBeanPostProcessor.getResolvableType(), properties.getDefinition(), kafkaStreamsFunctionProcessor); + public KafkaStreamsFunctionProcessorInvoker kafkaStreamsFunctionProcessorInvoker( + KafkaStreamsFunctionBeanPostProcessor kafkaStreamsFunctionBeanPostProcessor, + KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor, + KafkaStreamsFunctionProperties properties) { + return new KafkaStreamsFunctionProcessorInvoker(kafkaStreamsFunctionBeanPostProcessor.getResolvableType(), + properties.getDefinition(), kafkaStreamsFunctionProcessor); } @Bean - public KafkaStreamsFunctionBeanPostProcessor kafkaStreamsFunctionBeanPostProcessor(ConfigurableApplicationContext context) { - return new KafkaStreamsFunctionBeanPostProcessor(this.properties.getDefinition(), context); + public KafkaStreamsFunctionBeanPostProcessor kafkaStreamsFunctionBeanPostProcessor( + ConfigurableApplicationContext context, KafkaStreamsFunctionProperties properties) { + return new KafkaStreamsFunctionBeanPostProcessor(properties, context); } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java index 70f61075d..1e62f6e09 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java @@ -1,5 +1,5 @@ /* - * Copyright 2019 the original author or authors. + * Copyright 2019-2019 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. @@ -18,43 +18,41 @@ package org.springframework.cloud.stream.binder.kafka.streams.function; import java.lang.reflect.Method; -import org.springframework.beans.BeansException; +import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.AnnotatedBeanDefinition; -import org.springframework.beans.factory.config.BeanPostProcessor; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.core.ResolvableType; import org.springframework.util.ClassUtils; -public class KafkaStreamsFunctionBeanPostProcessor implements BeanPostProcessor { +public class KafkaStreamsFunctionBeanPostProcessor implements InitializingBean { - private final String functionName; + private final KafkaStreamsFunctionProperties kafkaStreamsFunctionProperties; private final ConfigurableApplicationContext context; private ResolvableType resolvableType; - public KafkaStreamsFunctionBeanPostProcessor(String functionName, ConfigurableApplicationContext context) { - this.functionName = functionName; + public KafkaStreamsFunctionBeanPostProcessor(KafkaStreamsFunctionProperties properties, + ConfigurableApplicationContext context) { + this.kafkaStreamsFunctionProperties = properties; this.context = context; } - @Override - public final Object postProcessAfterInitialization(Object bean, final String beanName) throws BeansException { - - if (beanName.equals(this.functionName)) { - final Class classObj = ClassUtils.resolveClassName(((AnnotatedBeanDefinition) - this.context.getBeanFactory().getBeanDefinition(beanName)).getMetadata().getClassName(), - ClassUtils.getDefaultClassLoader()); - - try { - Method method = classObj.getMethod(this.functionName, null); - resolvableType = ResolvableType.forMethodReturnType(method, classObj); - } catch (NoSuchMethodException e) { - //ignore - } - } - return bean; - } - public ResolvableType getResolvableType() { return this.resolvableType; } + + @Override + public void afterPropertiesSet() throws Exception { + final Class classObj = ClassUtils.resolveClassName(((AnnotatedBeanDefinition) + context.getBeanFactory().getBeanDefinition(kafkaStreamsFunctionProperties.getDefinition())) + .getMetadata().getClassName(), + ClassUtils.getDefaultClassLoader()); + + try { + Method method = classObj.getMethod(this.kafkaStreamsFunctionProperties.getDefinition(), null); + this.resolvableType = ResolvableType.forMethodReturnType(method, classObj); + } + catch (NoSuchMethodException e) { + //ignore + } + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionProcessorInvoker.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionProcessorInvoker.java index 05da9efd5..8f6a666b1 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionProcessorInvoker.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionProcessorInvoker.java @@ -1,5 +1,5 @@ /* - * Copyright 2019 the original author or authors. + * Copyright 2019-2019 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. @@ -27,7 +27,8 @@ public class KafkaStreamsFunctionProcessorInvoker { private final ResolvableType resolvableType; private final String functionName; - public KafkaStreamsFunctionProcessorInvoker(ResolvableType resolvableType, String functionName, KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor) { + public KafkaStreamsFunctionProcessorInvoker(ResolvableType resolvableType, String functionName, + KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor) { this.kafkaStreamsFunctionProcessor = kafkaStreamsFunctionProcessor; this.resolvableType = resolvableType; this.functionName = functionName; diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionProperties.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionProperties.java index d2ca79aa6..4dbf1e78d 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionProperties.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionProperties.java @@ -1,5 +1,5 @@ /* - * Copyright 2019 the original author or authors. + * Copyright 2019-2019 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. diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionWrapperDetector.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionWrapperDetector.java index e31feb135..8ca071937 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionWrapperDetector.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionWrapperDetector.java @@ -1,5 +1,5 @@ /* - * Copyright 2019 the original author or authors. + * Copyright 2019-2019 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. diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java index 6bb6aab22..e8e3cb253 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2019 the original author or authors. + * Copyright 2019-2019 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. @@ -56,7 +56,8 @@ import static org.assertj.core.api.Assertions.assertThat; public class KafkaStreamsBinderWordCountBranchesFunctionTests { @ClassRule - public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, "counts", "foo", "bar"); + public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, + "counts", "foo", "bar"); private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule.getEmbeddedKafka(); @@ -64,7 +65,8 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { @BeforeClass public static void setUp() throws Exception { - Map consumerProps = KafkaTestUtils.consumerProps("groupx", "false", embeddedKafka); + Map consumerProps = KafkaTestUtils.consumerProps("groupx", "false", + embeddedKafka); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(consumerProps); consumer = cf.createConsumer(); @@ -76,45 +78,6 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { consumer.close(); } - @EnableBinding(KStreamProcessorX.class) - @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) - public static class WordCountProcessorApplication { - - @Bean - @SuppressWarnings("unchecked") - public Function, KStream[]> process() { - - Predicate isEnglish = (k, v) -> v.word.equals("english"); - Predicate isFrench = (k, v) -> v.word.equals("french"); - Predicate isSpanish = (k, v) -> v.word.equals("spanish"); - - return input -> input - .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) - .groupBy((key, value) -> value) - .windowedBy(TimeWindows.of(5000)) - .count(Materialized.as("WordCounts-branch")) - .toStream() - .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, new Date(key.window().start()), new Date(key.window().end())))) - .branch(isEnglish, isFrench, isSpanish); - } - } - - interface KStreamProcessorX { - - @Input("input") - KStream input(); - - @Output("output1") - KStream output1(); - - @Output("output2") - KStream output2(); - - @Output("output3") - KStream output3(); - } - @Test public void testKstreamWordCountWithStringInputAndPojoOuput() throws Exception { SpringApplication app = new SpringApplication(WordCountProcessorApplication.class); @@ -131,16 +94,20 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { "--spring.cloud.stream.bindings.output3.destination=bar", "--spring.cloud.stream.bindings.output3.contentType=application/json", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", - "--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.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.timeWindow.length=5000", "--spring.cloud.stream.kafka.streams.timeWindow.advanceBy=0", - "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId=KafkaStreamsBinderWordCountBranchesFunctionTests-abc", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + + "=KafkaStreamsBinderWordCountBranchesFunctionTests-abc", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), "--spring.cloud.stream.kafka.streams.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString()); try { receiveAndValidate(context); - } finally { + } + finally { context.close(); } } @@ -215,4 +182,44 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { this.end = end; } } + + @EnableBinding(KStreamProcessorX.class) + @EnableAutoConfiguration + @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) + public static class WordCountProcessorApplication { + + @Bean + @SuppressWarnings("unchecked") + public Function, KStream[]> process() { + + Predicate isEnglish = (k, v) -> v.word.equals("english"); + Predicate isFrench = (k, v) -> v.word.equals("french"); + Predicate isSpanish = (k, v) -> v.word.equals("spanish"); + + return input -> input + .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) + .groupBy((key, value) -> value) + .windowedBy(TimeWindows.of(5000)) + .count(Materialized.as("WordCounts-branch")) + .toStream() + .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, + new Date(key.window().start()), new Date(key.window().end())))) + .branch(isEnglish, isFrench, isSpanish); + } + } + + interface KStreamProcessorX { + + @Input("input") + KStream input(); + + @Output("output1") + KStream output1(); + + @Output("output2") + KStream output2(); + + @Output("output3") + KStream output3(); + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java index 02a1a71d5..9dcf0d67e 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2019 the original author or authors. + * Copyright 2019-2019 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. @@ -56,7 +56,8 @@ import static org.assertj.core.api.Assertions.assertThat; public class KafkaStreamsBinderWordCountFunctionTests { @ClassRule - public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, "counts"); + public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, + "counts"); private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule.getEmbeddedKafka(); @@ -64,7 +65,8 @@ public class KafkaStreamsBinderWordCountFunctionTests { @BeforeClass public static void setUp() throws Exception { - Map consumerProps = KafkaTestUtils.consumerProps("group", "false", embeddedKafka); + Map consumerProps = KafkaTestUtils.consumerProps("group", "false", + embeddedKafka); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(consumerProps); consumer = cf.createConsumer(); @@ -76,25 +78,6 @@ public class KafkaStreamsBinderWordCountFunctionTests { consumer.close(); } - @EnableBinding(KafkaStreamsProcessor.class) - @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) - static class WordCountProcessorApplication { - - @Bean - public Function, KStream> process() { - - return input -> input - .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) - .map((key, value) -> new KeyValue<>(value, value)) - .groupByKey(Serialized.with(Serdes.String(), Serdes.String())) - .windowedBy(TimeWindows.of(5000)) - .count(Materialized.as("foo-WordCounts")) - .toStream() - .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, new Date(key.window().start()), new Date(key.window().end())))); - } - } - @Test public void testKstreamWordCountFunction() throws Exception { SpringApplication app = new SpringApplication(WordCountProcessorApplication.class); @@ -108,8 +91,10 @@ public class KafkaStreamsBinderWordCountFunctionTests { "--spring.cloud.stream.bindings.output.contentType=application/json", "--spring.cloud.stream.kafka.streams.default.consumer.application-id=basic-word-count", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", - "--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.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.bindings.output.producer.headerMode=raw", "--spring.cloud.stream.bindings.input.consumer.headerMode=raw", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) { @@ -126,7 +111,8 @@ public class KafkaStreamsBinderWordCountFunctionTests { template.sendDefault("foobar"); ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer, "counts"); assertThat(cr.value().contains("\"word\":\"foobar\",\"count\":1")).isTrue(); - } finally { + } + finally { pf.destroy(); } } @@ -180,4 +166,25 @@ public class KafkaStreamsBinderWordCountFunctionTests { this.end = end; } } + + @EnableBinding(KafkaStreamsProcessor.class) + @EnableAutoConfiguration + @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) + static class WordCountProcessorApplication { + + @Bean + public Function, KStream> process() { + + return input -> input + .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) + .map((key, value) -> new KeyValue<>(value, value)) + .groupByKey(Serialized.with(Serdes.String(), Serdes.String())) + .windowedBy(TimeWindows.of(5000)) + .count(Materialized.as("foo-WordCounts")) + .toStream() + .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, + new Date(key.window().start()), new Date(key.window().end())))); + } + } + } 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 0c3f8df50..83ae82678 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 @@ -1,5 +1,5 @@ /* - * Copyright 2019 the original author or authors. + * Copyright 2019-2019 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. @@ -59,12 +59,152 @@ import static org.assertj.core.api.Assertions.assertThat; public class StreamToGlobalKTableFunctionTests { @ClassRule - public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, "enriched-order"); + public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, + "enriched-order"); private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule.getEmbeddedKafka(); private static Consumer consumer; + @Test + public void testStreamToGlobalKTable() throws Exception { + SpringApplication app = new SpringApplication(OrderEnricherApplication.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run("--server.port=0", + "--spring.jmx.enabled=false", + "--spring.cloud.stream.kafka.streams.function.definition=process", + "--spring.cloud.stream.bindings.input.destination=orders", + "--spring.cloud.stream.bindings.input-x.destination=customers", + "--spring.cloud.stream.bindings.input-y.destination=products", + "--spring.cloud.stream.bindings.output.destination=enriched-order", + "--spring.cloud.stream.bindings.input.consumer.useNativeDecoding=true", + "--spring.cloud.stream.bindings.input-x.consumer.useNativeDecoding=true", + "--spring.cloud.stream.bindings.input-y.consumer.useNativeDecoding=true", + "--spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.keySerde" + + "=org.apache.kafka.common.serialization.Serdes$LongSerde", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde" + + "=org.springframework.cloud.stream.binder.kafka.streams.function" + + ".StreamToGlobalKTableFunctionTests$OrderSerde", + "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.keySerde" + + "=org.apache.kafka.common.serialization.Serdes$LongSerde", + "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.valueSerde" + + "=org.springframework.cloud.stream.binder.kafka.streams.function" + + ".StreamToGlobalKTableFunctionTests$CustomerSerde", + "--spring.cloud.stream.kafka.streams.bindings.input-y.consumer.keySerde" + + "=org.apache.kafka.common.serialization.Serdes$LongSerde", + "--spring.cloud.stream.kafka.streams.bindings.input-y.consumer.valueSerde" + + "=org.springframework.cloud.stream.binder.kafka.streams.function" + + ".StreamToGlobalKTableFunctionTests$ProductSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" + + "=org.apache.kafka.common.serialization.Serdes$LongSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde" + + "=org.springframework.cloud.stream.binder.kafka.streams." + + "function.StreamToGlobalKTableFunctionTests$EnrichedOrderSerde", + "--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.input.consumer.applicationId=" + + "StreamToGlobalKTableJoinFunctionTests-abc", + "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), + "--spring.cloud.stream.kafka.streams.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString())) { + Map senderPropsCustomer = KafkaTestUtils.producerProps(embeddedKafka); + senderPropsCustomer.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); + CustomerSerde customerSerde = new CustomerSerde(); + senderPropsCustomer.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, + customerSerde.serializer().getClass()); + + DefaultKafkaProducerFactory pfCustomer = + new DefaultKafkaProducerFactory<>(senderPropsCustomer); + KafkaTemplate template = new KafkaTemplate<>(pfCustomer, true); + template.setDefaultTopic("customers"); + for (long i = 0; i < 5; i++) { + final Customer customer = new Customer(); + customer.setName("customer-" + i); + template.sendDefault(i, customer); + } + + Map senderPropsProduct = KafkaTestUtils.producerProps(embeddedKafka); + senderPropsProduct.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); + ProductSerde productSerde = new ProductSerde(); + senderPropsProduct.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, productSerde.serializer().getClass()); + + DefaultKafkaProducerFactory pfProduct = + new DefaultKafkaProducerFactory<>(senderPropsProduct); + KafkaTemplate productTemplate = new KafkaTemplate<>(pfProduct, true); + productTemplate.setDefaultTopic("products"); + + for (long i = 0; i < 5; i++) { + final Product product = new Product(); + product.setName("product-" + i); + productTemplate.sendDefault(i, product); + } + + Map senderPropsOrder = KafkaTestUtils.producerProps(embeddedKafka); + senderPropsOrder.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); + OrderSerde orderSerde = new OrderSerde(); + senderPropsOrder.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, orderSerde.serializer().getClass()); + + DefaultKafkaProducerFactory pfOrder = new DefaultKafkaProducerFactory<>(senderPropsOrder); + KafkaTemplate orderTemplate = new KafkaTemplate<>(pfOrder, true); + orderTemplate.setDefaultTopic("orders"); + + for (long i = 0; i < 5; i++) { + final Order order = new Order(); + order.setCustomerId(i); + order.setProductId(i); + orderTemplate.sendDefault(i, order); + } + + Map consumerProps = KafkaTestUtils.consumerProps("group", "false", + embeddedKafka); + consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, LongDeserializer.class); + EnrichedOrderSerde enrichedOrderSerde = new EnrichedOrderSerde(); + consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, + enrichedOrderSerde.deserializer().getClass()); + consumerProps.put(JsonDeserializer.VALUE_DEFAULT_TYPE, + "org.springframework.cloud.stream.binder.kafka.streams." + + "function.StreamToGlobalKTableFunctionTests.EnrichedOrder"); + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(consumerProps); + + consumer = cf.createConsumer(); + embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "enriched-order"); + + int count = 0; + long start = System.currentTimeMillis(); + List> enrichedOrders = new ArrayList<>(); + do { + ConsumerRecords records = KafkaTestUtils.getRecords(consumer); + count = count + records.count(); + for (ConsumerRecord record : records) { + enrichedOrders.add(new KeyValue<>(record.key(), record.value())); + } + } while (count < 5 && (System.currentTimeMillis() - start) < 30000); + + assertThat(count == 5).isTrue(); + assertThat(enrichedOrders.size() == 5).isTrue(); + + enrichedOrders.sort(Comparator.comparing(o -> o.key)); + + for (int i = 0; i < 5; i++) { + KeyValue enrichedOrderKeyValue = enrichedOrders.get(i); + assertThat(enrichedOrderKeyValue.key == i).isTrue(); + EnrichedOrder enrichedOrder = enrichedOrderKeyValue.value; + assertThat(enrichedOrder.getOrder().customerId == i).isTrue(); + assertThat(enrichedOrder.getOrder().productId == i).isTrue(); + assertThat(enrichedOrder.getCustomer().name.equals("customer-" + i)).isTrue(); + assertThat(enrichedOrder.getProduct().name.equals("product-" + i)).isTrue(); + } + pfCustomer.destroy(); + pfProduct.destroy(); + pfOrder.destroy(); + consumer.close(); + } + } + interface CustomGlobalKTableProcessor extends KafkaStreamsProcessor { @Input("input-x") @@ -106,123 +246,6 @@ public class StreamToGlobalKTableFunctionTests { } } - @Test - public void testStreamToGlobalKTable() throws Exception { - SpringApplication app = new SpringApplication(OrderEnricherApplication.class); - app.setWebApplicationType(WebApplicationType.NONE); - try (ConfigurableApplicationContext ignored = app.run("--server.port=0", - "--spring.jmx.enabled=false", - "--spring.cloud.stream.kafka.streams.function.definition=process", - "--spring.cloud.stream.bindings.input.destination=orders", - "--spring.cloud.stream.bindings.input-x.destination=customers", - "--spring.cloud.stream.bindings.input-y.destination=products", - "--spring.cloud.stream.bindings.output.destination=enriched-order", - "--spring.cloud.stream.bindings.input.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.input-x.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.input-y.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", - "--spring.cloud.stream.kafka.streams.bindings.input.consumer.keySerde=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde=org.springframework.cloud.stream.binder.kafka.streams.function.StreamToGlobalKTableFunctionTests$OrderSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.keySerde=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.valueSerde=org.springframework.cloud.stream.binder.kafka.streams.function.StreamToGlobalKTableFunctionTests$CustomerSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-y.consumer.keySerde=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-y.consumer.valueSerde=org.springframework.cloud.stream.binder.kafka.streams.function.StreamToGlobalKTableFunctionTests$ProductSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde=org.springframework.cloud.stream.binder.kafka.streams.function.StreamToGlobalKTableFunctionTests$EnrichedOrderSerde", - "--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.input.consumer.applicationId=StreamToGlobalKTableJoinFunctionTests-abc", - "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString())) { - Map senderPropsCustomer = KafkaTestUtils.producerProps(embeddedKafka); - senderPropsCustomer.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); - CustomerSerde customerSerde = new CustomerSerde(); - senderPropsCustomer.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, customerSerde.serializer().getClass()); - - DefaultKafkaProducerFactory pfCustomer = new DefaultKafkaProducerFactory<>(senderPropsCustomer); - KafkaTemplate template = new KafkaTemplate<>(pfCustomer, true); - template.setDefaultTopic("customers"); - for (long i = 0; i < 5; i++) { - final Customer customer = new Customer(); - customer.setName("customer-" + i); - template.sendDefault(i, customer); - } - - Map senderPropsProduct = KafkaTestUtils.producerProps(embeddedKafka); - senderPropsProduct.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); - ProductSerde productSerde = new ProductSerde(); - senderPropsProduct.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, productSerde.serializer().getClass()); - - DefaultKafkaProducerFactory pfProduct = new DefaultKafkaProducerFactory<>(senderPropsProduct); - KafkaTemplate productTemplate = new KafkaTemplate<>(pfProduct, true); - productTemplate.setDefaultTopic("products"); - - for (long i = 0; i < 5; i++) { - final Product product = new Product(); - product.setName("product-" + i); - productTemplate.sendDefault(i, product); - } - - Map senderPropsOrder = KafkaTestUtils.producerProps(embeddedKafka); - senderPropsOrder.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); - OrderSerde orderSerde = new OrderSerde(); - senderPropsOrder.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, orderSerde.serializer().getClass()); - - DefaultKafkaProducerFactory pfOrder = new DefaultKafkaProducerFactory<>(senderPropsOrder); - KafkaTemplate orderTemplate = new KafkaTemplate<>(pfOrder, true); - orderTemplate.setDefaultTopic("orders"); - - for (long i = 0; i < 5; i++) { - final Order order = new Order(); - order.setCustomerId(i); - order.setProductId(i); - orderTemplate.sendDefault(i, order); - } - - Map consumerProps = KafkaTestUtils.consumerProps("group", "false", embeddedKafka); - consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); - consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, LongDeserializer.class); - EnrichedOrderSerde enrichedOrderSerde = new EnrichedOrderSerde(); - consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, enrichedOrderSerde.deserializer().getClass()); - consumerProps.put(JsonDeserializer.VALUE_DEFAULT_TYPE, "org.springframework.cloud.stream.binder.kafka.streams.function.StreamToGlobalKTableFunctionTests.EnrichedOrder"); - DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(consumerProps); - - consumer = cf.createConsumer(); - embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "enriched-order"); - - int count = 0; - long start = System.currentTimeMillis(); - List> enrichedOrders = new ArrayList<>(); - do { - ConsumerRecords records = KafkaTestUtils.getRecords(consumer); - count = count + records.count(); - for (ConsumerRecord record : records) { - enrichedOrders.add(new KeyValue<>(record.key(), record.value())); - } - } while (count < 5 && (System.currentTimeMillis() - start) < 30000); - - assertThat(count == 5).isTrue(); - assertThat(enrichedOrders.size() == 5).isTrue(); - - enrichedOrders.sort(Comparator.comparing(o -> o.key)); - - for (int i = 0; i < 5; i++) { - KeyValue enrichedOrderKeyValue = enrichedOrders.get(i); - assertThat(enrichedOrderKeyValue.key == i).isTrue(); - EnrichedOrder enrichedOrder = enrichedOrderKeyValue.value; - assertThat(enrichedOrder.getOrder().customerId == i).isTrue(); - assertThat(enrichedOrder.getOrder().productId == i).isTrue(); - assertThat(enrichedOrder.getCustomer().name.equals("customer-" + i)).isTrue(); - assertThat(enrichedOrder.getProduct().name.equals("product-" + i)).isTrue(); - } - pfCustomer.destroy(); - pfProduct.destroy(); - pfOrder.destroy(); - consumer.close(); - } - } - static class Order { long customerId; diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java index 97de70570..1c4dce870 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2019 the original author or authors. + * Copyright 2019-2019 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. @@ -62,62 +62,19 @@ import static org.assertj.core.api.Assertions.assertThat; public class StreamToTableJoinFunctionTests { @ClassRule - public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, "output-topic-1", "output-topic-2"); + public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, + true, "output-topic-1", "output-topic-2"); private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule.getEmbeddedKafka(); - @EnableBinding(KStreamKTableProcessor.class) - @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) - public static class CountClicksPerRegionApplication { - - @Bean - public Function, Function, KStream>> process1() { - return userClicksStream -> (userRegionsTable -> (userClicksStream - .leftJoin(userRegionsTable, (clicks, region) -> new RegionWithClicks(region == null ? "UNKNOWN" : region, clicks), - Joined.with(Serdes.String(), Serdes.Long(), null)) - .map((user, regionWithClicks) -> new KeyValue<>(regionWithClicks.getRegion(), regionWithClicks.getClicks())) - .groupByKey(Serialized.with(Serdes.String(), Serdes.Long())) - .reduce((firstClicks, secondClicks) -> firstClicks + secondClicks) - .toStream())); - } - } - - interface KStreamKTableProcessor { - - /** - * Input binding. - * - * @return {@link Input} binding for {@link KStream} type. - */ - @Input("input-1") - KStream input1(); - - /** - * Input binding. - * - * @return {@link Input} binding for {@link KStream} type. - */ - @Input("input-2") - KTable input2(); - - /** - * Output binding. - * - * @return {@link Output} binding for {@link KStream} type. - */ - @Output("output") - KStream output(); - - } - @Test public void testStreamToTable() throws Exception { SpringApplication app = new SpringApplication(CountClicksPerRegionApplication.class); app.setWebApplicationType(WebApplicationType.NONE); Consumer consumer; - Map consumerProps = KafkaTestUtils.consumerProps("group-1", "false", embeddedKafka); + Map consumerProps = KafkaTestUtils.consumerProps("group-1", + "false", embeddedKafka); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, LongDeserializer.class); @@ -134,16 +91,25 @@ public class StreamToTableJoinFunctionTests { "--spring.cloud.stream.bindings.input-1.consumer.useNativeDecoding=true", "--spring.cloud.stream.bindings.input-2.consumer.useNativeDecoding=true", "--spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", - "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.keySerde=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.valueSerde=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.keySerde=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.valueSerde=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--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.bindings.input-1.consumer.keySerde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.valueSerde" + + "=org.apache.kafka.common.serialization.Serdes$LongSerde", + "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.keySerde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.valueSerde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde" + + "=org.apache.kafka.common.serialization.Serdes$LongSerde", + "--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.input-1.consumer.applicationId=StreamToTableJoinFunctionTests-abc", + "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.applicationId" + + "=StreamToTableJoinFunctionTests-abc", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), "--spring.cloud.stream.kafka.streams.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString())) { @@ -214,7 +180,8 @@ public class StreamToTableJoinFunctionTests { assertThat(count == expectedClicksPerRegion.size()).isTrue(); assertThat(actualClicksPerRegion).hasSameElementsAs(expectedClicksPerRegion); - } finally { + } + finally { consumer.close(); } } @@ -225,7 +192,8 @@ public class StreamToTableJoinFunctionTests { app.setWebApplicationType(WebApplicationType.NONE); Consumer consumer; - Map consumerProps = KafkaTestUtils.consumerProps("group-2", "false", embeddedKafka); + Map consumerProps = KafkaTestUtils.consumerProps("group-2", + "false", embeddedKafka); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, LongDeserializer.class); @@ -272,16 +240,25 @@ public class StreamToTableJoinFunctionTests { "--spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", "--spring.cloud.stream.kafka.streams.binder.configuration.auto.offset.reset=latest", "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.startOffset=earliest", - "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.keySerde=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.valueSerde=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.keySerde=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.valueSerde=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--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.bindings.input-1.consumer.keySerde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.valueSerde" + + "=org.apache.kafka.common.serialization.Serdes$LongSerde", + "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.keySerde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.valueSerde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde" + + "=org.apache.kafka.common.serialization.Serdes$LongSerde", + "--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.input-1.consumer.application-id=StreamToTableJoinFunctionTests-foobar", + "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.application-id" + + "=StreamToTableJoinFunctionTests-foobar", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), "--spring.cloud.stream.kafka.streams.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString())) { Thread.sleep(1000L); @@ -349,7 +326,8 @@ public class StreamToTableJoinFunctionTests { assertThat(count).isEqualTo(expectedClicksPerRegion.size()); assertThat(actualClicksPerRegion).hasSameElementsAs(expectedClicksPerRegion); - } finally { + } + finally { consumer.close(); } } @@ -382,4 +360,51 @@ public class StreamToTableJoinFunctionTests { } } + + @EnableBinding(KStreamKTableProcessor.class) + @EnableAutoConfiguration + @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) + public static class CountClicksPerRegionApplication { + + @Bean + public Function, Function, KStream>> process1() { + return userClicksStream -> (userRegionsTable -> (userClicksStream + .leftJoin(userRegionsTable, (clicks, region) -> new RegionWithClicks(region == null ? + "UNKNOWN" : region, clicks), + Joined.with(Serdes.String(), Serdes.Long(), null)) + .map((user, regionWithClicks) -> new KeyValue<>(regionWithClicks.getRegion(), + regionWithClicks.getClicks())) + .groupByKey(Serialized.with(Serdes.String(), Serdes.Long())) + .reduce((firstClicks, secondClicks) -> firstClicks + secondClicks) + .toStream())); + } + } + + interface KStreamKTableProcessor { + + /** + * Input binding. + * + * @return {@link Input} binding for {@link KStream} type. + */ + @Input("input-1") + KStream input1(); + + /** + * Input binding. + * + * @return {@link Input} binding for {@link KStream} type. + */ + @Input("input-2") + KTable input2(); + + /** + * Output binding. + * + * @return {@link Output} binding for {@link KStream} type. + */ + @Output("output") + KStream output(); + + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToGlobalKTableJoinIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToGlobalKTableJoinIntegrationTests.java index 221571e8e..2bd8a7adb 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToGlobalKTableJoinIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToGlobalKTableJoinIntegrationTests.java @@ -20,7 +20,6 @@ import java.util.ArrayList; import java.util.Comparator; import java.util.List; import java.util.Map; -import java.util.function.Function; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; @@ -45,7 +44,6 @@ import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; @@ -72,73 +70,6 @@ public class StreamToGlobalKTableJoinIntegrationTests { private static Consumer consumer; - interface CustomGlobalKTableProcessor extends KafkaStreamsProcessor { - - @Input("input-x") - GlobalKTable inputX(); - - @Input("input-y") - GlobalKTable inputY(); - } - - @EnableBinding(CustomGlobalKTableProcessor.class) - @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) - public static class OrderEnricherApplication { - - @StreamListener - @SendTo("output") - public KStream process(@Input("input") KStream ordersStream, - @Input("input-x") GlobalKTable customers, - @Input("input-y") GlobalKTable products) { - - KStream customerOrdersStream = ordersStream.join(customers, - (orderId, order) -> order.getCustomerId(), - (order, customer) -> new CustomerOrder(customer, order)); - - return customerOrdersStream.join(products, - (orderId, customerOrder) -> customerOrder - .productId(), - (customerOrder, product) -> { - EnrichedOrder enrichedOrder = new EnrichedOrder(); - enrichedOrder.setProduct(product); - enrichedOrder.setCustomer(customerOrder.customer); - enrichedOrder.setOrder(customerOrder.order); - return enrichedOrder; - }); - } - - @Bean - public Function, - Function, - Function, KStream>>> hello() { - - return orderStream -> ( - customers -> ( - products -> ( - orderStream.join(customers, - (orderId, order) -> order.getCustomerId(), - (order, customer) -> new CustomerOrder(customer, order)) - .join(products, - (orderId, customerOrder) -> customerOrder - .productId(), - (customerOrder, product) -> { - EnrichedOrder enrichedOrder = new EnrichedOrder(); - enrichedOrder.setProduct(product); - enrichedOrder.setCustomer(customerOrder.customer); - enrichedOrder.setOrder(customerOrder.order); - return enrichedOrder; - }) - ) - ) - ); - - - } - - - } - @Test public void testStreamToGlobalKTable() throws Exception { SpringApplication app = new SpringApplication( @@ -299,6 +230,45 @@ public class StreamToGlobalKTableJoinIntegrationTests { } + interface CustomGlobalKTableProcessor extends KafkaStreamsProcessor { + + @Input("input-x") + GlobalKTable inputX(); + + @Input("input-y") + GlobalKTable inputY(); + + } + + @EnableBinding(CustomGlobalKTableProcessor.class) + @EnableAutoConfiguration + @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) + public static class OrderEnricherApplication { + + @StreamListener + @SendTo("output") + public KStream process( + @Input("input") KStream ordersStream, + @Input("input-x") GlobalKTable customers, + @Input("input-y") GlobalKTable products) { + + KStream customerOrdersStream = ordersStream.join( + customers, (orderId, order) -> order.getCustomerId(), + (order, customer) -> new CustomerOrder(customer, order)); + + return customerOrdersStream.join(products, + (orderId, customerOrder) -> customerOrder.productId(), + (customerOrder, product) -> { + EnrichedOrder enrichedOrder = new EnrichedOrder(); + enrichedOrder.setProduct(product); + enrichedOrder.setCustomer(customerOrder.customer); + enrichedOrder.setOrder(customerOrder.order); + return enrichedOrder; + }); + } + + } + static class Order { long customerId; @@ -417,4 +387,5 @@ public class StreamToGlobalKTableJoinIntegrationTests { public static class EnrichedOrderSerde extends JsonSerde { } + } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java index 93083bfa8..715b632ce 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java @@ -190,7 +190,8 @@ public class StreamToTableJoinIntegrationTests { assertThat(count == expectedClicksPerRegion.size()).isTrue(); assertThat(actualClicksPerRegion).hasSameElementsAs(expectedClicksPerRegion); - } finally { + } + finally { consumer.close(); } } @@ -336,7 +337,8 @@ public class StreamToTableJoinIntegrationTests { assertThat(count).isEqualTo(expectedClicksPerRegion.size()); assertThat(actualClicksPerRegion).hasSameElementsAs(expectedClicksPerRegion); - } finally { + } + finally { consumer.close(); } }