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 2b2ea1054..91aa4605c 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 @@ -70,10 +70,10 @@ public class StreamToGlobalKTableJoinIntegrationTests { interface CustomGlobalKTableProcessor extends KafkaStreamsProcessor { - @Input("inputX") + @Input("input-x") GlobalKTable inputX(); - @Input("inputY") + @Input("input-y") GlobalKTable inputY(); } @@ -85,8 +85,8 @@ public class StreamToGlobalKTableJoinIntegrationTests { @StreamListener @SendTo("output") public KStream process(@Input("input") KStream ordersStream, - @Input("inputX") GlobalKTable customers, - @Input("inputY") GlobalKTable products) { + @Input("input-x") GlobalKTable customers, + @Input("input-y") GlobalKTable products) { KStream customerOrdersStream = ordersStream.join(customers, (orderId, order) -> order.getCustomerId(), @@ -112,19 +112,19 @@ public class StreamToGlobalKTableJoinIntegrationTests { try (ConfigurableApplicationContext ignored = app.run("--server.port=0", "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=orders", - "--spring.cloud.stream.bindings.inputX.destination=customers", - "--spring.cloud.stream.bindings.inputY.destination=products", + "--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.inputX.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.inputY.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.integration.StreamToGlobalKTableJoinIntegrationTests$OrderSerde", - "--spring.cloud.stream.kafka.streams.bindings.inputX.consumer.keySerde=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.inputX.consumer.valueSerde=org.springframework.cloud.stream.binder.kafka.streams.integration.StreamToGlobalKTableJoinIntegrationTests$CustomerSerde", - "--spring.cloud.stream.kafka.streams.bindings.inputY.consumer.keySerde=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.inputY.consumer.valueSerde=org.springframework.cloud.stream.binder.kafka.streams.integration.StreamToGlobalKTableJoinIntegrationTests$ProductSerde", + "--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.integration.StreamToGlobalKTableJoinIntegrationTests$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.integration.StreamToGlobalKTableJoinIntegrationTests$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.integration.StreamToGlobalKTableJoinIntegrationTests$EnrichedOrderSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde", 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 36fb755b8..bb9cea6f3 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 @@ -97,7 +97,7 @@ public class StreamToTableJoinIntegrationTests { @StreamListener @SendTo("output") public KStream process(@Input("input") KStream userClicksStream, - @Input("inputX") KTable userRegionsTable) { + @Input("input-x") KTable userRegionsTable) { return userClicksStream .leftJoin(userRegionsTable, (clicks, region) -> new RegionWithClicks(region == null ? "UNKNOWN" : region, clicks), @@ -111,7 +111,7 @@ public class StreamToTableJoinIntegrationTests { interface KafkaStreamsProcessorX extends KafkaStreamsProcessor { - @Input("inputX") + @Input("input-x") KTable inputX(); } @@ -123,10 +123,10 @@ public class StreamToTableJoinIntegrationTests { try (ConfigurableApplicationContext ignored = app.run("--server.port=0", "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=user-clicks", - "--spring.cloud.stream.bindings.inputX.destination=user-regions", + "--spring.cloud.stream.bindings.input-x.destination=user-regions", "--spring.cloud.stream.bindings.output.destination=output-topic", "--spring.cloud.stream.bindings.input.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.inputX.consumer.useNativeDecoding=true", + "--spring.cloud.stream.bindings.input-x.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$StringSerde", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde=org.apache.kafka.common.serialization.Serdes$LongSerde", diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderExtendedPropertiesTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderExtendedPropertiesTest.java index 2c18b7c04..c55551be0 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderExtendedPropertiesTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderExtendedPropertiesTest.java @@ -59,7 +59,9 @@ properties = {"spring.cloud.stream.kafka.bindings.output.producer.configuration. "spring.cloud.stream.kafka.default.consumer.configuration.key.serializer=BarSerializer.class", "spring.cloud.stream.kafka.default.consumer.configuration.value.serializer=BarSerializer.class", "spring.cloud.stream.kafka.default.producer.configuration.foo=bar", - "spring.cloud.stream.kafka.bindings.output.producer.configuration.foo=bindingSpecificPropertyShouldWinOverDefault"}) + "spring.cloud.stream.kafka.bindings.output.producer.configuration.foo=bindingSpecificPropertyShouldWinOverDefault", + "spring.cloud.stream.kafka.default.consumer.ackEachRecord=true", + "spring.cloud.stream.kafka.bindings.custom-in.consumer.ackEachRecord=false"}) public class KafkaBinderExtendedPropertiesTest { private static final String KAFKA_BROKERS_PROPERTY = "spring.cloud.stream.kafka.binder.brokers"; @@ -122,6 +124,9 @@ public class KafkaBinderExtendedPropertiesTest { //binding "input" gets BarSerializer and BarSerializer for ker.serializer/value.serializer through default properties. assertThat(customKafkaConsumerProperties.getConfiguration().get("key.serializer")).isEqualTo("BarSerializer.class"); assertThat(customKafkaConsumerProperties.getConfiguration().get("value.serializer")).isEqualTo("BarSerializer.class"); + + assertThat(kafkaConsumerProperties.isAckEachRecord()).isEqualTo(true); + assertThat(customKafkaConsumerProperties.isAckEachRecord()).isEqualTo(false); } @EnableBinding(CustomBindingForExtendedPropertyTesting.class) @@ -130,13 +135,13 @@ public class KafkaBinderExtendedPropertiesTest { @StreamListener(Sink.INPUT) @SendTo(Processor.OUTPUT) - public String process(String payload) throws InterruptedException { + public String process(String payload) { return payload; } @StreamListener("custom-in") @SendTo("custom-out") - public String processCustom(String payload) throws InterruptedException { + public String processCustom(String payload) { return payload; }