From daedfb92108fe94699579d1845e487450922df81 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 13 Feb 2019 16:05:01 -0500 Subject: [PATCH] Upgrade dependencies and fix issues --- .../inbound/MessageDrivenAdapterTests.java | 4 +- .../kafka/dsl/KafkaDslKotlinTests.kt | 86 ++++++++++--------- 2 files changed, 47 insertions(+), 43 deletions(-) diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java index f385e2b5a0..cdeba539cd 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2018 the original author or authors. + * Copyright 2016-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. @@ -445,7 +445,7 @@ public class MessageDrivenAdapterTests { public void testPauseResume() throws Exception { ConsumerFactory cf = mock(ConsumerFactory.class); Consumer consumer = mock(Consumer.class); - given(cf.createConsumer(isNull(), eq("clientId"), isNull())).willReturn(consumer); + given(cf.createConsumer(isNull(), eq("clientId"), isNull(), any())).willReturn(consumer); final Map>> records = new HashMap<>(); records.put(new TopicPartition("foo", 0), Arrays.asList( new ConsumerRecord<>("foo", 0, 0L, 1, "foo"), diff --git a/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/KafkaDslKotlinTests.kt b/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/KafkaDslKotlinTests.kt index 4abb133a80..4d18c31e27 100644 --- a/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/KafkaDslKotlinTests.kt +++ b/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/KafkaDslKotlinTests.kt @@ -1,5 +1,5 @@ /* - * Copyright 2018 the original author or authors. + * Copyright 2018-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. @@ -16,8 +16,14 @@ package org.springframework.integration.kafka.dsl -import assertk.assert -import assertk.assertions.* +import assertk.assertThat +import assertk.assertions.contains +import assertk.assertions.isEqualTo +import assertk.assertions.isInstanceOf +import assertk.assertions.isNotNull +import assertk.assertions.isNull +import assertk.assertions.isSameAs +import assertk.assertions.isTrue import assertk.catch import org.apache.kafka.clients.consumer.ConsumerConfig import org.apache.kafka.clients.consumer.ConsumerRebalanceListener @@ -44,7 +50,11 @@ import org.springframework.integration.kafka.support.RawRecordHeaderErrorMessage import org.springframework.integration.support.MessageBuilder import org.springframework.integration.test.util.TestUtils import org.springframework.kafka.annotation.EnableKafka -import org.springframework.kafka.core.* +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.ContainerProperties import org.springframework.kafka.listener.GenericMessageListenerContainer import org.springframework.kafka.listener.KafkaMessageListenerContainer @@ -143,80 +153,74 @@ class KafkaDslKotlinTests { fun testKafkaAdapters() { val exception = catch { this.sendToKafkaFlowInput.send(GenericMessage("foo")) } - assert(exception!!.message).isNotNull { - it.contains("10 is not in the range") - } + assertThat(exception!!.message).isNotNull().contains("10 is not in the range") this.kafkaProducer1.setPartitionIdExpression(ValueExpression(0)) this.kafkaProducer2.setPartitionIdExpression(ValueExpression(0)) this.sendToKafkaFlowInput.send(GenericMessage("foo", hashMapOf("foo" to "bar"))) - assert(TestUtils.getPropertyValue(this.kafkaProducer1, "headerMapper")).isSameAs(this.mapper) + assertThat(TestUtils.getPropertyValue(this.kafkaProducer1, "headerMapper")).isSameAs(this.mapper) for (i in 0..99) { val receive = this.listeningFromKafkaResults1.receive(20000) - assert(receive).isNotNull() - assert(receive!!.payload).isEqualTo("FOO") + assertThat(receive).isNotNull() + assertThat(receive!!.payload).isEqualTo("FOO") val headers = receive.headers - assert(headers.containsKey(KafkaHeaders.ACKNOWLEDGMENT)).isTrue() + assertThat(headers.containsKey(KafkaHeaders.ACKNOWLEDGMENT)).isTrue() val acknowledgment = headers.get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment::class.java) acknowledgment?.acknowledge() - assert(headers[KafkaHeaders.RECEIVED_TOPIC]).isEqualTo(TEST_TOPIC1) - assert(headers[KafkaHeaders.RECEIVED_MESSAGE_KEY]).isEqualTo(i + 1) - assert(headers[KafkaHeaders.RECEIVED_PARTITION_ID]).isEqualTo(0) - assert(headers[KafkaHeaders.OFFSET]).isEqualTo(i.toLong()) - assert(headers[KafkaHeaders.TIMESTAMP_TYPE]).isEqualTo("CREATE_TIME") - assert(headers[KafkaHeaders.RECEIVED_TIMESTAMP]).isEqualTo(1487694048633L) - assert(headers["foo"]).isEqualTo("bar") + assertThat(headers[KafkaHeaders.RECEIVED_TOPIC]).isEqualTo(TEST_TOPIC1) + assertThat(headers[KafkaHeaders.RECEIVED_MESSAGE_KEY]).isEqualTo(i + 1) + assertThat(headers[KafkaHeaders.RECEIVED_PARTITION_ID]).isEqualTo(0) + assertThat(headers[KafkaHeaders.OFFSET]).isEqualTo(i.toLong()) + assertThat(headers[KafkaHeaders.TIMESTAMP_TYPE]).isEqualTo("CREATE_TIME") + assertThat(headers[KafkaHeaders.RECEIVED_TIMESTAMP]).isEqualTo(1487694048633L) + assertThat(headers["foo"]).isEqualTo("bar") } for (i in 0..99) { val receive = this.listeningFromKafkaResults2.receive(20000) - assert(receive).isNotNull() - assert(receive!!.payload).isEqualTo("FOO") + assertThat(receive).isNotNull() + assertThat(receive!!.payload).isEqualTo("FOO") val headers = receive.headers - assert(headers.containsKey(KafkaHeaders.ACKNOWLEDGMENT)).isTrue() + assertThat(headers.containsKey(KafkaHeaders.ACKNOWLEDGMENT)).isTrue() val acknowledgment = headers.get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment::class.java) acknowledgment?.acknowledge() - assert(headers[KafkaHeaders.RECEIVED_TOPIC]).isEqualTo(TEST_TOPIC2) - assert(headers[KafkaHeaders.RECEIVED_MESSAGE_KEY]).isEqualTo(i + 1) - assert(headers[KafkaHeaders.RECEIVED_PARTITION_ID]).isEqualTo(0) - assert(headers[KafkaHeaders.OFFSET]).isEqualTo(i.toLong()) - assert(headers[KafkaHeaders.TIMESTAMP_TYPE]).isEqualTo("CREATE_TIME") - assert(headers[KafkaHeaders.RECEIVED_TIMESTAMP]).isEqualTo(1487694048644L) + assertThat(headers[KafkaHeaders.RECEIVED_TOPIC]).isEqualTo(TEST_TOPIC2) + assertThat(headers[KafkaHeaders.RECEIVED_MESSAGE_KEY]).isEqualTo(i + 1) + assertThat(headers[KafkaHeaders.RECEIVED_PARTITION_ID]).isEqualTo(0) + assertThat(headers[KafkaHeaders.OFFSET]).isEqualTo(i.toLong()) + assertThat(headers[KafkaHeaders.TIMESTAMP_TYPE]).isEqualTo("CREATE_TIME") + assertThat(headers[KafkaHeaders.RECEIVED_TIMESTAMP]).isEqualTo(1487694048644L) } val message = MessageBuilder.withPayload("BAR").setHeader(KafkaHeaders.TOPIC, TEST_TOPIC2).build() this.sendToKafkaFlowInput.send(message) - assert(this.listeningFromKafkaResults1.receive(10)).isNull() + assertThat(this.listeningFromKafkaResults1.receive(10)).isNull() val error = this.errorChannel.receive(10000) - assert(error).isNotNull() { - it.isInstanceOf(ErrorMessage::class.java) - } + assertThat(error).isNotNull().isInstanceOf(ErrorMessage::class.java) val payload = error?.payload - assert(payload).isNotNull() { - it.isInstanceOf(MessageRejectedException::class.java) - } + assertThat(payload).isNotNull().isInstanceOf(MessageRejectedException::class.java) - assert(this.messageListenerContainer).isNotNull() - assert(this.kafkaTemplateTopic1).isNotNull() - assert(this.kafkaTemplateTopic2).isNotNull() + assertThat(this.messageListenerContainer).isNotNull() + assertThat(this.kafkaTemplateTopic1).isNotNull() + assertThat(this.kafkaTemplateTopic2).isNotNull() this.kafkaTemplateTopic1.send(TEST_TOPIC3, "foo") - assert(this.config.sourceFlowLatch.await(10, TimeUnit.SECONDS)).isTrue() - assert(this.config.fromSource).isEqualTo("foo") + assertThat(this.config.sourceFlowLatch.await(10, TimeUnit.SECONDS)).isTrue() + assertThat(this.config.fromSource).isEqualTo("foo") } @Test fun testGateways() { - assert(this.config.replyContainerLatch.await(30, TimeUnit.SECONDS)) - assert(this.gate.exchange(TEST_TOPIC4, "foo")).isEqualTo("FOO") + assertThat(this.config.replyContainerLatch.await(30, TimeUnit.SECONDS)) + assertThat(this.gate.exchange(TEST_TOPIC4, "foo")).isEqualTo("FOO") } @Configuration