* GH-210: Upgrade test-embedded-kafka to JUnit5 Resolves https://github.com/spring-cloud/spring-cloud-stream-samples/issues/210 * Avoid a rebalance by using a different group in the test; use `assign()` instead of `subscribe()`.
101 lines
4.1 KiB
Java
101 lines
4.1 KiB
Java
/*
|
|
* Copyright 2017-2018 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.
|
|
* You may obtain a copy of the License at
|
|
*
|
|
* https://www.apache.org/licenses/LICENSE-2.0
|
|
*
|
|
* Unless required by applicable law or agreed to in writing, software
|
|
* distributed under the License is distributed on an "AS IS" BASIS,
|
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
* See the License for the specific language governing permissions and
|
|
* limitations under the License.
|
|
*/
|
|
/*
|
|
* Copyright 2017-2018 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.
|
|
* You may obtain a copy of the License at
|
|
*
|
|
* https://www.apache.org/licenses/LICENSE-2.0
|
|
*
|
|
* Unless required by applicable law or agreed to in writing, software
|
|
* distributed under the License is distributed on an "AS IS" BASIS,
|
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
* See the License for the specific language governing permissions and
|
|
* limitations under the License.
|
|
*/
|
|
|
|
package demo;
|
|
|
|
import org.apache.kafka.clients.consumer.Consumer;
|
|
import org.apache.kafka.clients.consumer.ConsumerConfig;
|
|
import org.apache.kafka.clients.consumer.ConsumerRecords;
|
|
import org.apache.kafka.common.TopicPartition;
|
|
import org.apache.kafka.common.serialization.ByteArrayDeserializer;
|
|
import org.apache.kafka.common.serialization.ByteArraySerializer;
|
|
import org.junit.jupiter.api.Test;
|
|
|
|
import org.springframework.beans.factory.annotation.Autowired;
|
|
import org.springframework.boot.test.context.SpringBootTest;
|
|
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
|
|
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
|
|
import org.springframework.kafka.core.KafkaTemplate;
|
|
import org.springframework.kafka.test.EmbeddedKafkaBroker;
|
|
import org.springframework.kafka.test.context.EmbeddedKafka;
|
|
import org.springframework.kafka.test.utils.KafkaTestUtils;
|
|
|
|
import java.time.Duration;
|
|
import java.util.Collections;
|
|
import java.util.Map;
|
|
|
|
import static org.assertj.core.api.Assertions.assertThat;
|
|
|
|
/**
|
|
* Test class demonstrating how to use an embedded kafka service with the
|
|
* kafka binder.
|
|
*
|
|
* @author Gary Russell
|
|
* @author Soby Chacko
|
|
*
|
|
*/
|
|
@SpringBootTest
|
|
@EmbeddedKafka(topics = { EmbeddedKafkaApplicationTests.INPUT_TOPIC, EmbeddedKafkaApplicationTests.OUTPUT_TOPIC },
|
|
partitions = 1,
|
|
bootstrapServersProperty = "spring.kafka.bootstrap-servers")
|
|
public class EmbeddedKafkaApplicationTests {
|
|
|
|
public static final String INPUT_TOPIC = "testEmbeddedIn";
|
|
public static final String OUTPUT_TOPIC = "testEmbeddedOut";
|
|
private static final String GROUP_NAME = "embeddedKafkaApplicationTest";
|
|
|
|
@Test
|
|
void testSendReceive(@Autowired EmbeddedKafkaBroker embeddedKafka) {
|
|
Map<String, Object> senderProps = KafkaTestUtils.producerProps(embeddedKafka);
|
|
senderProps.put("key.serializer", ByteArraySerializer.class);
|
|
senderProps.put("value.serializer", ByteArraySerializer.class);
|
|
DefaultKafkaProducerFactory<byte[], byte[]> pf = new DefaultKafkaProducerFactory<>(senderProps);
|
|
KafkaTemplate<byte[], byte[]> template = new KafkaTemplate<>(pf, true);
|
|
template.setDefaultTopic(INPUT_TOPIC);
|
|
template.sendDefault("foo".getBytes());
|
|
|
|
Map<String, Object> consumerProps = KafkaTestUtils.consumerProps(GROUP_NAME, "false", embeddedKafka);
|
|
consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
|
|
consumerProps.put("key.deserializer", ByteArrayDeserializer.class);
|
|
consumerProps.put("value.deserializer", ByteArrayDeserializer.class);
|
|
DefaultKafkaConsumerFactory<byte[], byte[]> cf = new DefaultKafkaConsumerFactory<>(consumerProps);
|
|
|
|
Consumer<byte[], byte[]> consumer = cf.createConsumer();
|
|
consumer.assign(Collections.singleton(new TopicPartition(OUTPUT_TOPIC, 0)));
|
|
ConsumerRecords<byte[], byte[]> records = consumer.poll(Duration.ofSeconds(10));
|
|
consumer.commitSync();
|
|
|
|
assertThat(records.count()).isEqualTo(1);
|
|
assertThat(new String(records.iterator().next().value())).isEqualTo("FOO");
|
|
}
|
|
|
|
}
|