From 16bb3e2f62b4f1bb4ec1ffca62b214a23914300e Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 23 Aug 2019 18:53:32 -0400 Subject: [PATCH] Testing topic provisioning properties for KStream Resolves #684 --- .../StreamToTableJoinIntegrationTests.java | 12 ++++++++++++ 1 file changed, 12 insertions(+) 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 694e258af..cece90df4 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 @@ -54,6 +54,7 @@ import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder; import org.springframework.cloud.stream.binder.ProducerProperties; import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsConsumerProperties; +import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsProducerProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.CleanupConfig; @@ -111,6 +112,7 @@ public class StreamToTableJoinIntegrationTests { "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + "=StreamToTableJoinIntegrationTests-abc", "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.topic.properties.cleanup.policy=compact", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.topic.properties.cleanup.policy=compact", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString()); try { @@ -128,6 +130,16 @@ public class StreamToTableJoinIntegrationTests { assertThat(cleanupPolicyX).isEqualTo("compact"); + Binder kStreamBinder = binderFactory + .getBinder("kstream", KStream.class); + + KafkaStreamsProducerProperties producerProperties = (KafkaStreamsProducerProperties) ((ExtendedPropertiesBinder) kStreamBinder) + .getExtendedProducerProperties("output"); + + String cleanupPolicyOutput = producerProperties.getTopic().getProperties().get("cleanup.policy"); + + assertThat(cleanupPolicyOutput).isEqualTo("compact"); + // Input 1: Region per user (multiple records allowed per user). List> userRegions = Arrays.asList(new KeyValue<>( "alice", "asia"), /* Alice lived in Asia originally... */