From 0d0cf8dcb7992e261d260294c79371bc2347b902 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 7 Aug 2018 10:50:12 -0400 Subject: [PATCH] Update to spring-cloud-build 2.1.0 snapshot Spring Boot 2.1.0 snapshot Spring Kafka 2.2.0/3.1.0 snapshots Apache Kafka client 2.0.0 Fixing tests Removing deprecations and removals Polishing Resolves #424 --- pom.xml | 8 ++-- ...KafkaStreamsMessageConversionDelegate.java | 6 --- ...StreamListenerSetupMethodOrchestrator.java | 2 +- ...reamsInteractiveQueryIntegrationTests.java | 3 +- ...afkaStreamsStateStoreIntegrationTests.java | 5 -- .../kafka/KafkaMessageChannelBinder.java | 16 +++---- .../stream/binder/kafka/AdminConfigTests.java | 3 +- .../kafka/KafkaBinderHealthIndicatorTest.java | 8 ++-- .../binder/kafka/KafkaBinderMetricsTest.java | 20 ++++---- .../stream/binder/kafka/KafkaBinderTests.java | 48 +++++++++---------- .../binder/kafka/KafkaBinderUnitTests.java | 43 ++++++++++------- .../binder/kafka/KafkaTransactionTests.java | 6 +-- .../bootstrap/KafkaBinderBootstrapTest.java | 8 ++-- .../integration/KafkaBinderActuatorTests.java | 9 ++-- 14 files changed, 89 insertions(+), 96 deletions(-) diff --git a/pom.xml b/pom.xml index 87533f8b3..0508dd1ac 100644 --- a/pom.xml +++ b/pom.xml @@ -7,14 +7,14 @@ org.springframework.cloud spring-cloud-build - 2.0.2.RELEASE + 2.1.0.BUILD-SNAPSHOT 1.8 - 2.1.7.RELEASE - 3.0.3.RELEASE - 1.1.0 + 2.2.0.BUILD-SNAPSHOT + 3.1.0.BUILD-SNAPSHOT + 2.0.0 2.1.0.BUILD-SNAPSHOT diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java index 192004c68..714377407 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java @@ -182,12 +182,6 @@ public class KafkaStreamsMessageConversionDelegate { } } - @SuppressWarnings("deprecation") - @Override - public void punctuate(long timestamp) { - - } - @Override public void close() { 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 f348d669a..5969b8b2c 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 @@ -26,10 +26,10 @@ import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.utils.Bytes; -import org.apache.kafka.streams.Consumed; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.errors.DeserializationExceptionHandler; +import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.KTable; import org.apache.kafka.streams.kstream.Materialized; diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java index 84b9f658c..12482dafb 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java @@ -25,6 +25,7 @@ import org.apache.kafka.common.serialization.IntegerSerializer; import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.kstream.KStream; +import org.apache.kafka.streams.kstream.Materialized; import org.apache.kafka.streams.kstream.Serialized; import org.apache.kafka.streams.state.HostInfo; import org.apache.kafka.streams.state.QueryableStoreTypes; @@ -138,7 +139,7 @@ public class KafkaStreamsInteractiveQueryIntegrationTests { .filter((key, product) -> product.getId() == 123) .map((key, value) -> new KeyValue<>(value.id, value)) .groupByKey(Serialized.with(new Serdes.IntegerSerde(), new JsonSerde<>(Product.class))) - .count("prod-id-count-store") + .count(Materialized.as("prod-id-count-store")) .toStream() .map((key, value) -> new KeyValue<>(null, "Count for product with ID 123: " + value)); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsStateStoreIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsStateStoreIntegrationTests.java index a900e85dd..2c4b2d6aa 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsStateStoreIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsStateStoreIntegrationTests.java @@ -116,11 +116,6 @@ public class KafkaStreamsStateStoreIntegrationTests { processed = true; } - @Override - public void punctuate(long l) { - - } - @Override public void close() { if (state != null) { diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index d9d8a5514..eb909fc32 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -75,27 +75,25 @@ import org.springframework.cloud.stream.provisioning.ProducerDestination; import org.springframework.context.Lifecycle; import org.springframework.expression.common.LiteralExpression; import org.springframework.expression.spel.standard.SpelExpressionParser; +import org.springframework.integration.StaticMessageHeaderAccessor; +import org.springframework.integration.acks.AcknowledgmentCallback; import org.springframework.integration.channel.ChannelInterceptorAware; import org.springframework.integration.core.MessageProducer; import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter; import org.springframework.integration.kafka.inbound.KafkaMessageSource; import org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler; import org.springframework.integration.kafka.support.RawRecordHeaderErrorMessageStrategy; -import org.springframework.integration.support.AcknowledgmentCallback; -import org.springframework.integration.support.AcknowledgmentCallback.Status; import org.springframework.integration.support.ErrorMessageStrategy; import org.springframework.integration.support.MessageBuilder; -import org.springframework.integration.support.StaticMessageHeaderAccessor; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.AbstractMessageListenerContainer; -import org.springframework.kafka.listener.AbstractMessageListenerContainer.AckMode; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.ConsumerAwareRebalanceListener; -import org.springframework.kafka.listener.config.ContainerProperties; +import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.support.DefaultKafkaHeaderMapper; import org.springframework.kafka.support.KafkaHeaderMapper; import org.springframework.kafka.support.KafkaHeaders; @@ -417,14 +415,14 @@ public class KafkaMessageChannelBinder extends // end of these won't be needed... if (!extendedConsumerProperties.getExtension().isAutoCommitOffset()) { messageListenerContainer.getContainerProperties() - .setAckMode(AbstractMessageListenerContainer.AckMode.MANUAL); + .setAckMode(ContainerProperties.AckMode.MANUAL); messageListenerContainer.getContainerProperties().setAckOnError(false); } else { messageListenerContainer.getContainerProperties() .setAckOnError(isAutoCommitOnError(extendedConsumerProperties)); if (extendedConsumerProperties.getExtension().isAckEachRecord()) { - messageListenerContainer.getContainerProperties().setAckMode(AckMode.RECORD); + messageListenerContainer.getContainerProperties().setAckMode(ContainerProperties.AckMode.RECORD); } } if (this.logger.isDebugEnabled()) { @@ -783,10 +781,10 @@ public class KafkaMessageChannelBinder extends ((MessagingException) message.getPayload()).getFailedMessage()); if (ack != null) { if (isAutoCommitOnError(properties)) { - ack.acknowledge(Status.REJECT); + ack.acknowledge(AcknowledgmentCallback.Status.REJECT); } else { - ack.acknowledge(Status.REQUEUE); + ack.acknowledge(AcknowledgmentCallback.Status.REQUEUE); } } } diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AdminConfigTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AdminConfigTests.java index a6b8c61e3..ec36c153d 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AdminConfigTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AdminConfigTests.java @@ -46,7 +46,8 @@ import static org.assertj.core.api.Assertions.assertThat; @TestPropertySource(properties = { "spring.cloud.stream.kafka.bindings.input.consumer.admin.replication-factor=2", "spring.cloud.stream.kafka.bindings.input.consumer.admin.replicas-assignments.0=0,1", - "spring.cloud.stream.kafka.bindings.input.consumer.admin.configuration.message.format.version=0.9.0.0" }) + "spring.cloud.stream.kafka.bindings.input.consumer.admin.configuration.message.format.version=0.9.0.0", + "spring.main.allow-bean-definition-overriding=true"}) @EnableIntegration public class AdminConfigTests { diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicatorTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicatorTest.java index 2b3f101ed..bec846086 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicatorTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicatorTest.java @@ -73,7 +73,7 @@ public class KafkaBinderHealthIndicatorTest { @Test public void kafkaBinderIsUp() { final List partitions = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group", partitions)); + topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group1-healthIndicator", partitions)); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); Health health = indicator.health(); assertThat(health.getStatus()).isEqualTo(Status.UP); @@ -82,7 +82,7 @@ public class KafkaBinderHealthIndicatorTest { @Test public void kafkaBinderIsDown() { final List partitions = partitions(new Node(-1, null, 0)); - topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group", partitions)); + topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group2-healthIndicator", partitions)); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); Health health = indicator.health(); assertThat(health.getStatus()).isEqualTo(Status.DOWN); @@ -91,7 +91,7 @@ public class KafkaBinderHealthIndicatorTest { @Test(timeout = 5000) public void kafkaBinderDoesNotAnswer() { final List partitions = partitions(new Node(-1, null, 0)); - topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group", partitions)); + topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group3-healthIndicator", partitions)); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willAnswer(new Answer() { @Override @@ -110,7 +110,7 @@ public class KafkaBinderHealthIndicatorTest { @Test public void createsConsumerOnceWhenInvokedMultipleTimes() { final List partitions = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group", partitions)); + topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group4-healthIndicator", partitions)); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); indicator.health(); diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java index 394c2f75b..3926f9b07 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java @@ -84,11 +84,11 @@ public class KafkaBinderMetricsTest { public void shouldIndicateLag() { org.mockito.BDDMockito.given(consumer.committed(ArgumentMatchers.any(TopicPartition.class))).willReturn(new OffsetAndMetadata(500)); List partitions = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, new TopicInformation("group", partitions)); + topicsInUse.put(TEST_TOPIC, new TopicInformation("group1-metrics", partitions)); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); - assertThat(meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group").tag("topic", TEST_TOPIC).timeGauge() + assertThat(meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group1-metrics").tag("topic", TEST_TOPIC).timeGauge() .value(TimeUnit.MILLISECONDS)).isEqualTo(500.0); } @@ -100,22 +100,22 @@ public class KafkaBinderMetricsTest { org.mockito.BDDMockito.given(consumer.endOffsets(ArgumentMatchers.anyCollection())).willReturn(endOffsets); org.mockito.BDDMockito.given(consumer.committed(ArgumentMatchers.any(TopicPartition.class))).willReturn(new OffsetAndMetadata(500)); List partitions = partitions(new Node(0, null, 0), new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, new TopicInformation("group", partitions)); + topicsInUse.put(TEST_TOPIC, new TopicInformation("group2-metrics", partitions)); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); - assertThat(meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group").tag("topic", TEST_TOPIC).timeGauge() + assertThat(meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group2-metrics").tag("topic", TEST_TOPIC).timeGauge() .value(TimeUnit.MILLISECONDS)).isEqualTo(1000.0); } @Test public void shouldIndicateFullLagForNotCommittedGroups() { List partitions = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, new TopicInformation("group", partitions)); + topicsInUse.put(TEST_TOPIC, new TopicInformation("group3-metrics", partitions)); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); - assertThat(meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group").tag("topic", TEST_TOPIC).timeGauge() + assertThat(meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group3-metrics").tag("topic", TEST_TOPIC).timeGauge() .value(TimeUnit.MILLISECONDS)).isEqualTo(1000.0); } @@ -130,11 +130,11 @@ public class KafkaBinderMetricsTest { @Test public void createsConsumerOnceWhenInvokedMultipleTimes() { final List partitions = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, new TopicInformation("group", partitions)); + topicsInUse.put(TEST_TOPIC, new TopicInformation("group4-metrics", partitions)); metrics.bindTo(meterRegistry); - TimeGauge gauge = meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group").tag("topic", TEST_TOPIC).timeGauge(); + TimeGauge gauge = meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group4-metrics").tag("topic", TEST_TOPIC).timeGauge(); gauge.value(TimeUnit.MILLISECONDS); assertThat(gauge.value(TimeUnit.MILLISECONDS)).isEqualTo(1000.0); @@ -147,11 +147,11 @@ public class KafkaBinderMetricsTest { .willReturn(consumer); final List partitions = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, new TopicInformation("group", partitions)); + topicsInUse.put(TEST_TOPIC, new TopicInformation("group5-metrics", partitions)); metrics.bindTo(meterRegistry); - TimeGauge gauge = meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group").tag("topic", TEST_TOPIC).timeGauge(); + TimeGauge gauge = meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group5-metrics").tag("topic", TEST_TOPIC).timeGauge(); assertThat(gauge.value(TimeUnit.MILLISECONDS)).isEqualTo(0); assertThat(gauge.value(TimeUnit.MILLISECONDS)).isEqualTo(1000.0); diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index 9024e0cbc..83ca0b920 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -35,7 +35,6 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import com.fasterxml.jackson.databind.ObjectMapper; - import org.apache.kafka.clients.admin.AdminClient; import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.admin.CreateTopicsResult; @@ -104,8 +103,8 @@ import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.AbstractMessageListenerContainer; -import org.springframework.kafka.listener.AbstractMessageListenerContainer.AckMode; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; +import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.listener.MessageListenerContainer; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.KafkaHeaders; @@ -113,7 +112,7 @@ import org.springframework.kafka.support.SendResult; import org.springframework.kafka.support.TopicPartitionInitialOffset; import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.kafka.test.core.BrokerAddress; -import org.springframework.kafka.test.rule.KafkaEmbedded; +import org.springframework.kafka.test.rule.EmbeddedKafkaRule; import org.springframework.kafka.test.utils.KafkaTestUtils; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; @@ -127,7 +126,6 @@ import org.springframework.messaging.support.ErrorMessage; import org.springframework.messaging.support.GenericMessage; import org.springframework.messaging.support.MessageBuilder; import org.springframework.util.Assert; -import org.springframework.util.MimeType; import org.springframework.util.MimeTypeUtils; import org.springframework.util.concurrent.ListenableFuture; import org.springframework.util.concurrent.SettableListenableFuture; @@ -154,7 +152,7 @@ public class KafkaBinderTests extends private final String CLASS_UNDER_TEST_NAME = KafkaMessageChannelBinder.class.getSimpleName(); @ClassRule - public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, 10, "error.pollableDlq.group"); + public static EmbeddedKafkaRule embeddedKafka = new EmbeddedKafkaRule(1, true, 10, "error.pollableDlq.group-pcWithDlq"); private KafkaTestBinder binder; @@ -214,7 +212,7 @@ public class KafkaBinderTests extends private KafkaBinderConfigurationProperties createConfigurationProperties() { KafkaBinderConfigurationProperties binderConfiguration = new KafkaBinderConfigurationProperties( new TestKafkaProperties()); - BrokerAddress[] brokerAddresses = embeddedKafka.getBrokerAddresses(); + BrokerAddress[] brokerAddresses = embeddedKafka.getEmbeddedKafka().getBrokerAddresses(); List bAddresses = new ArrayList<>(); for (BrokerAddress bAddress : brokerAddresses) { bAddresses.add(bAddress.toString()); @@ -247,7 +245,7 @@ public class KafkaBinderTests extends timeoutMultiplier = Double.parseDouble(multiplier); } - BrokerAddress[] brokerAddresses = embeddedKafka.getBrokerAddresses(); + BrokerAddress[] brokerAddresses = embeddedKafka.getEmbeddedKafka().getBrokerAddresses(); List bAddresses = new ArrayList<>(); for (BrokerAddress bAddress : brokerAddresses) { bAddresses.add(bAddress.toString()); @@ -339,9 +337,9 @@ public class KafkaBinderTests extends Assertions.assertThat(inboundMessageRef.get().getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE)).isNull(); Assertions.assertThat(inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE)) .isEqualTo(MimeTypeUtils.TEXT_PLAIN); - Assertions.assertThat(inboundMessageRef.get().getHeaders().get("foo")).isInstanceOf(MimeType.class); - MimeType actual = (MimeType) inboundMessageRef.get().getHeaders().get("foo"); - Assertions.assertThat(actual).isEqualTo(MimeTypeUtils.TEXT_PLAIN); + Assertions.assertThat(inboundMessageRef.get().getHeaders().get("foo")).isInstanceOf(String.class); + String actual = (String) inboundMessageRef.get().getHeaders().get("foo"); + Assertions.assertThat(actual).isEqualTo(MimeTypeUtils.TEXT_PLAIN.toString()); producerBinding.unbind(); consumerBinding.unbind(); } @@ -498,7 +496,7 @@ public class KafkaBinderTests extends assertThat(receivedMessage.getHeaders().get(KafkaMessageChannelBinder.X_ORIGINAL_TOPIC)) .isEqualTo("foo.bar".getBytes(StandardCharsets.UTF_8)); assertThat(new String((byte[]) receivedMessage.getHeaders().get(KafkaMessageChannelBinder.X_EXCEPTION_MESSAGE))) - .startsWith("failed to send Message to channel 'input'"); + .startsWith("Dispatcher failed to deliver Message; nested exception is java.lang.RuntimeException: fail"); assertThat(receivedMessage.getHeaders().get(KafkaMessageChannelBinder.X_EXCEPTION_STACKTRACE)) .isNotNull(); binderBindUnbindLatency(); @@ -680,7 +678,7 @@ public class KafkaBinderTests extends .isEqualTo(TimestampType.CREATE_TIME.toString()); assertThat(((String) receivedMessage.getHeaders().get(KafkaMessageChannelBinder.X_EXCEPTION_MESSAGE))) - .startsWith("failed to send Message to channel 'input'"); + .startsWith("Dispatcher failed to deliver Message; nested exception is java.lang.RuntimeException: fail"); assertThat(receivedMessage.getHeaders().get(KafkaMessageChannelBinder.X_EXCEPTION_STACKTRACE)) .isNotNull(); assertThat(receivedMessage.getHeaders().get(KafkaMessageChannelBinder.X_EXCEPTION_FQCN)).isNotNull(); @@ -704,7 +702,7 @@ public class KafkaBinderTests extends .isEqualTo(TimestampType.CREATE_TIME.toString().getBytes()); assertThat(new String((byte[]) receivedMessage.getHeaders().get(KafkaMessageChannelBinder.X_EXCEPTION_MESSAGE))) - .startsWith("failed to send Message to channel 'input'"); + .startsWith("Dispatcher failed to deliver Message; nested exception is java.lang.RuntimeException: fail"); assertThat(receivedMessage.getHeaders().get(KafkaMessageChannelBinder.X_EXCEPTION_STACKTRACE)) .isNotNull(); @@ -1167,7 +1165,7 @@ public class KafkaBinderTests extends AbstractMessageListenerContainer container = TestUtils.getPropertyValue(consumerBinding, "lifecycle.messageListenerContainer", AbstractMessageListenerContainer.class); - assertThat(container.getContainerProperties().getAckMode()).isEqualTo(AckMode.BATCH); + assertThat(container.getContainerProperties().getAckMode()).isEqualTo(ContainerProperties.AckMode.BATCH); String testPayload1 = "foo" + UUID.randomUUID().toString(); Message message1 = org.springframework.integration.support.MessageBuilder.withPayload( @@ -1214,7 +1212,7 @@ public class KafkaBinderTests extends AbstractMessageListenerContainer container = TestUtils.getPropertyValue(consumerBinding2, "lifecycle.messageListenerContainer", AbstractMessageListenerContainer.class); - assertThat(container.getContainerProperties().getAckMode()).isEqualTo(AckMode.RECORD); + assertThat(container.getContainerProperties().getAckMode()).isEqualTo(ContainerProperties.AckMode.RECORD); Message receivedMessage1 = receive(inbound1); assertThat(receivedMessage1).isNotNull(); @@ -2466,7 +2464,7 @@ public class KafkaBinderTests extends assertThat(inbound.getHeaders().get(BinderHeaders.NATIVE_HEADERS_PRESENT)).isNull(); Map consumerProps = KafkaTestUtils.consumerProps("testSendAndReceiveWithMixedMode", "false", - embeddedKafka); + embeddedKafka.getEmbeddedKafka()); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); @@ -2503,9 +2501,9 @@ public class KafkaBinderTests extends PollableSource inboundBindTarget = new DefaultPollableMessageSource(this.messageConverter); ExtendedConsumerProperties consumerProps = createConsumerProperties(); consumerProps.setMultiplex(true); - Binding> binding = binder.bindPollableConsumer("pollable,anotherOne", "group", + Binding> binding = binder.bindPollableConsumer("pollable,anotherOne", "group-polledConsumer", inboundBindTarget, consumerProps); - Map producerProps = KafkaTestUtils.producerProps(embeddedKafka); + Map producerProps = KafkaTestUtils.producerProps(embeddedKafka.getEmbeddedKafka()); KafkaTemplate template = new KafkaTemplate(new DefaultKafkaProducerFactory<>(producerProps)); template.send("pollable", "testPollable"); boolean polled = inboundBindTarget.poll(m -> { @@ -2544,8 +2542,8 @@ public class KafkaBinderTests extends properties.setMaxAttempts(2); properties.setBackOffInitialInterval(0); properties.getExtension().setEnableDlq(true); - Map producerProps = KafkaTestUtils.producerProps(embeddedKafka); - Binding> binding = binder.bindPollableConsumer("pollableDlq", "group", + Map producerProps = KafkaTestUtils.producerProps(embeddedKafka.getEmbeddedKafka()); + Binding> binding = binder.bindPollableConsumer("pollableDlq", "group-pcWithDlq", inboundBindTarget, properties); KafkaTemplate template = new KafkaTemplate(new DefaultKafkaProducerFactory<>(producerProps)); template.send("pollableDlq", "testPollableDLQ"); @@ -2561,12 +2559,12 @@ public class KafkaBinderTests extends catch (MessageHandlingException e) { assertThat(e.getCause().getMessage()).isEqualTo("test DLQ"); } - Map consumerProps = KafkaTestUtils.consumerProps("dlq", "false", embeddedKafka); + Map consumerProps = KafkaTestUtils.consumerProps("dlq", "false", embeddedKafka.getEmbeddedKafka()); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); ConsumerFactory cf = new DefaultKafkaConsumerFactory<>(consumerProps); Consumer consumer = cf.createConsumer(); - embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "error.pollableDlq.group"); - ConsumerRecord deadLetter = KafkaTestUtils.getSingleRecord(consumer, "error.pollableDlq.group"); + embeddedKafka.getEmbeddedKafka().consumeFromAnEmbeddedTopic(consumer, "error.pollableDlq.group-pcWithDlq"); + ConsumerRecord deadLetter = KafkaTestUtils.getSingleRecord(consumer, "error.pollableDlq.group-pcWithDlq"); assertThat(deadLetter).isNotNull(); assertThat(deadLetter.value()).isEqualTo("testPollableDLQ"); binding.unbind(); @@ -2577,7 +2575,7 @@ public class KafkaBinderTests extends @Test public void testTopicPatterns() throws Exception { try (AdminClient admin = AdminClient.create(Collections.singletonMap(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, - embeddedKafka.getBrokersAsString()))) { + embeddedKafka.getEmbeddedKafka().getBrokersAsString()))) { admin.createTopics(Collections.singletonList(new NewTopic("topicPatterns.1", 1, (short) 1))).all().get(); Binder binder = getBinder(); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); @@ -2592,7 +2590,7 @@ public class KafkaBinderTests extends Binding consumerBinding = binder.bindConsumer("topicPatterns\\..*", "testTopicPatterns", moduleInputChannel, consumerProperties); DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory( - KafkaTestUtils.producerProps(embeddedKafka)); + KafkaTestUtils.producerProps(embeddedKafka.getEmbeddedKafka())); KafkaTemplate template = new KafkaTemplate(pf); template.send("topicPatterns.1", "foo"); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderUnitTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderUnitTests.java index 48f3e6b4f..7347c2eeb 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderUnitTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderUnitTests.java @@ -36,6 +36,7 @@ import org.apache.kafka.common.TopicPartition; import org.junit.Test; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; +import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; @@ -82,25 +83,25 @@ public class KafkaBinderUnitTests { method.setAccessible(true); // test default for anon - Object factory = method.invoke(binder, true, "foo", ecp); + Object factory = method.invoke(binder, true, "foo-1", ecp); Map configs = TestUtils.getPropertyValue(factory, "configs", Map.class); assertThat(configs.get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)).isEqualTo("latest"); // test default for named - factory = method.invoke(binder, false, "foo", ecp); + factory = method.invoke(binder, false, "foo-2", ecp); configs = TestUtils.getPropertyValue(factory, "configs", Map.class); assertThat(configs.get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)).isEqualTo("earliest"); // binder level setting binderConfigurationProperties.setConfiguration( Collections.singletonMap(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest")); - factory = method.invoke(binder, false, "foo", ecp); + factory = method.invoke(binder, false, "foo-3", ecp); configs = TestUtils.getPropertyValue(factory, "configs", Map.class); assertThat(configs.get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)).isEqualTo("latest"); // consumer level setting consumerProps.setConfiguration(Collections.singletonMap(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest")); - factory = method.invoke(binder, false, "foo", ecp); + factory = method.invoke(binder, false, "foo-4", ecp); configs = TestUtils.getPropertyValue(factory, "configs", Map.class); assertThat(configs.get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)).isEqualTo("earliest"); } @@ -131,38 +132,38 @@ public class KafkaBinderUnitTests { @Test public void testOffsetResetWithGroupManagementEarliest() throws Exception { - testOffsetResetWithGroupManagement(true, true); + testOffsetResetWithGroupManagement(true, true, "foo-100", "testOffsetResetWithGroupManagementEarliest"); } @Test public void testOffsetResetWithGroupManagementLatest() throws Throwable { - testOffsetResetWithGroupManagement(false, true); + testOffsetResetWithGroupManagement(false, true, "foo-101", "testOffsetResetWithGroupManagementLatest"); } @Test public void testOffsetResetWithManualAssignmentEarliest() throws Exception { - testOffsetResetWithGroupManagement(true, false); + testOffsetResetWithGroupManagement(true, false, "foo-102", "testOffsetResetWithManualAssignmentEarliest"); } @Test public void testOffsetResetWithGroupManualAssignmentLatest() throws Throwable { - testOffsetResetWithGroupManagement(false, false); + testOffsetResetWithGroupManagement(false, false, "foo-103", "testOffsetResetWithGroupManualAssignmentLatest"); } - private void testOffsetResetWithGroupManagement(final boolean earliest, boolean groupManage) throws Exception { + private void testOffsetResetWithGroupManagement(final boolean earliest, boolean groupManage, String topic, String group) throws Exception { final List partitions = new ArrayList<>(); - partitions.add(new TopicPartition("foo", 0)); - partitions.add(new TopicPartition("foo", 1)); + partitions.add(new TopicPartition(topic, 0)); + partitions.add(new TopicPartition(topic, 1)); KafkaBinderConfigurationProperties configurationProperties = new KafkaBinderConfigurationProperties( new TestKafkaProperties()); KafkaTopicProvisioner provisioningProvider = mock(KafkaTopicProvisioner.class); ConsumerDestination dest = mock(ConsumerDestination.class); - given(dest.getName()).willReturn("foo"); + given(dest.getName()).willReturn(topic); given(provisioningProvider.provisionConsumerDestination(anyString(), anyString(), any())).willReturn(dest); final AtomicInteger part = new AtomicInteger(); willAnswer(i -> { return partitions.stream() - .map(p -> new PartitionInfo("foo", part.getAndIncrement(), null, null, null)) + .map(p -> new PartitionInfo(topic, part.getAndIncrement(), null, null, null)) .collect(Collectors.toList()); }).given(provisioningProvider).getPartitionsForTopic(anyInt(), anyBoolean(), any()); @SuppressWarnings("unchecked") @@ -183,7 +184,7 @@ public class KafkaBinderUnitTests { latch.countDown(); latch.countDown(); return null; - }).given(consumer).subscribe(eq(Collections.singletonList("foo")), + }).given(consumer).subscribe(eq(Collections.singletonList(topic)), any(org.apache.kafka.clients.consumer.ConsumerRebalanceListener.class)); willAnswer(i -> { latch.countDown(); @@ -193,7 +194,7 @@ public class KafkaBinderUnitTests { @Override protected ConsumerFactory createKafkaConsumerFactory(boolean anonymous, String consumerGroup, - ExtendedConsumerProperties consumerProperties) { + ExtendedConsumerProperties consumerProperties) { return new ConsumerFactory() { @@ -212,6 +213,11 @@ public class KafkaBinderUnitTests { return consumer; } + @Override + public Consumer createConsumer(String groupId, String clientIdPrefix, String clientIdSuffix) { + return consumer; + } + @Override public boolean isAutoCommit() { return false; @@ -222,7 +228,7 @@ public class KafkaBinderUnitTests { Map props = new HashMap<>(); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, earliest ? "earliest" : "latest"); - props.put(ConsumerConfig.GROUP_ID_CONFIG, "bar"); + props.put(ConsumerConfig.GROUP_ID_CONFIG, group); return props; } @@ -240,7 +246,7 @@ public class KafkaBinderUnitTests { ExtendedConsumerProperties consumerProperties = new ExtendedConsumerProperties( extension); consumerProperties.setInstanceCount(1); - binder.bindConsumer("foo", "bar", channel, consumerProperties); + Binding messageChannelBinding = binder.bindConsumer(topic, group, channel, consumerProperties); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); if (groupManage) { if (earliest) { @@ -260,7 +266,8 @@ public class KafkaBinderUnitTests { verify(consumer).seek(partitions.get(1), Long.MAX_VALUE); } } + messageChannelBinding.unbind(); } -} +} \ No newline at end of file diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaTransactionTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaTransactionTests.java index 762b3110b..4995b54f8 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaTransactionTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaTransactionTests.java @@ -34,7 +34,7 @@ import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProv import org.springframework.context.support.GenericApplicationContext; import org.springframework.integration.channel.DirectChannel; import org.springframework.kafka.core.DefaultKafkaProducerFactory; -import org.springframework.kafka.test.rule.KafkaEmbedded; +import org.springframework.kafka.test.rule.EmbeddedKafkaRule; import org.springframework.messaging.support.GenericMessage; import org.springframework.retry.support.RetryTemplate; @@ -53,13 +53,13 @@ import static org.mockito.Mockito.spy; public class KafkaTransactionTests { @ClassRule - public static final KafkaEmbedded embeddedKafka = new KafkaEmbedded(1); + public static final EmbeddedKafkaRule embeddedKafka = new EmbeddedKafkaRule(1); @SuppressWarnings({ "rawtypes", "unchecked" }) @Test public void testProducerRunsInTx() { KafkaProperties kafkaProperties = new TestKafkaProperties(); - kafkaProperties.setBootstrapServers(Collections.singletonList(embeddedKafka.getBrokersAsString())); + kafkaProperties.setBootstrapServers(Collections.singletonList(embeddedKafka.getEmbeddedKafka().getBrokersAsString())); KafkaBinderConfigurationProperties configurationProperties = new KafkaBinderConfigurationProperties(kafkaProperties); configurationProperties.getTransaction().setTransactionIdPrefix("foo-"); diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/bootstrap/KafkaBinderBootstrapTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/bootstrap/KafkaBinderBootstrapTest.java index fbc29be4e..a4de2c2ac 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/bootstrap/KafkaBinderBootstrapTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/bootstrap/KafkaBinderBootstrapTest.java @@ -23,7 +23,7 @@ import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.kafka.test.rule.KafkaEmbedded; +import org.springframework.kafka.test.rule.EmbeddedKafkaRule; /** * @author Marius Bogoevici @@ -31,14 +31,14 @@ import org.springframework.kafka.test.rule.KafkaEmbedded; public class KafkaBinderBootstrapTest { @ClassRule - public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, 10); + public static EmbeddedKafkaRule embeddedKafka = new EmbeddedKafkaRule(1, true, 10); @Test public void testKafkaBinderConfiguration() throws Exception { ConfigurableApplicationContext applicationContext = new SpringApplicationBuilder(SimpleApplication.class) .web(WebApplicationType.NONE) - .run("--spring.cloud.stream.kafka.binder.brokers=" + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString()); + .run("--spring.cloud.stream.kafka.binder.brokers=" + embeddedKafka.getEmbeddedKafka().getBrokersAsString(), + "--spring.cloud.stream.kafka.binder.zkNodes=" + embeddedKafka.getEmbeddedKafka().getZookeeperConnectionString()); applicationContext.close(); } diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java index daba1b48e..237a64dbd 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java @@ -21,7 +21,6 @@ import java.util.Map; import io.micrometer.core.instrument.MeterRegistry; import io.micrometer.core.instrument.binder.MeterBinder; - import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.ClassRule; @@ -43,7 +42,7 @@ import org.springframework.cloud.stream.messaging.Sink; import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.listener.AbstractMessageListenerContainer; -import org.springframework.kafka.test.rule.KafkaEmbedded; +import org.springframework.kafka.test.rule.EmbeddedKafkaRule; import org.springframework.messaging.MessageChannel; import org.springframework.test.context.junit4.SpringRunner; @@ -61,16 +60,16 @@ import static org.assertj.core.api.Assertions.assertThat; properties = "spring.cloud.stream.bindings.input.group=" + KafkaBinderActuatorTests.TEST_CONSUMER_GROUP) public class KafkaBinderActuatorTests { - static final String TEST_CONSUMER_GROUP = "testGroup"; + static final String TEST_CONSUMER_GROUP = "testGroup-actuatorTests"; private static final String KAFKA_BROKERS_PROPERTY = "spring.kafka.bootstrap-servers"; @ClassRule - public static KafkaEmbedded kafkaEmbedded = new KafkaEmbedded(1, true); + public static EmbeddedKafkaRule kafkaEmbedded = new EmbeddedKafkaRule(1, true); @BeforeClass public static void setup() { - System.setProperty(KAFKA_BROKERS_PROPERTY, kafkaEmbedded.getBrokersAsString()); + System.setProperty(KAFKA_BROKERS_PROPERTY, kafkaEmbedded.getEmbeddedKafka().getBrokersAsString()); } @AfterClass