import org.apache.kafka.common.utils.Bytes;
import org.junit.jupiter.api.*;
import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
+import org.springframework.boot.test.autoconfigure.data.mongo.AutoConfigureDataMongo;
import org.springframework.boot.test.context.ConfigDataApplicationContextInitializer;
import org.springframework.boot.test.context.TestConfiguration;
import org.springframework.context.annotation.Import;
properties = {
"consumer.bootstrap-server=${spring.embedded.kafka.brokers}",
"consumer.topic=" + TOPIC,
- "consumer.commit-interval=1s" })
+ "consumer.commit-interval=1s",
+ "spring.mongodb.embedded.version=4.4.13" })
@EmbeddedKafka(topics = TOPIC, partitions = PARTITIONS)
+@EnableAutoConfiguration
+@AutoConfigureDataMongo
@Slf4j
abstract class GenericApplicationTests<K, V>
{
ApplicationProperties properties;
@Autowired
ExecutorService executor;
+ @Autowired
+ PollIntervalAwareConsumerRebalanceListener rebalanceListener;
+ @Autowired
+ RecordHandler<K, V> recordHandler;
KafkaProducer<Bytes, Bytes> testRecordProducer;
KafkaConsumer<Bytes, Bytes> offsetConsumer;
assertThatExceptionOfType(IllegalStateException.class)
.isThrownBy(() -> endlessConsumer.exitStatus())
.describedAs("Consumer should still be running");
+
+ recordGenerator.assertBusinessLogic();
}
@Test
assertThat(endlessConsumer.exitStatus())
.describedAs("Consumer should have exited abnormally")
.containsInstanceOf(RecordDeserializationException.class);
+
+ recordGenerator.assertBusinessLogic();
}
@Test
assertThat(endlessConsumer.exitStatus())
.describedAs("Consumer should have exited abnormally")
.containsInstanceOf(RuntimeException.class);
+
+ recordGenerator.assertBusinessLogic();
}
{
return true;
}
+
+ default void assertBusinessLogic()
+ {
+ log.debug("No business-logic to assert");
+ }
}
void sendMessage(ProducerRecord<Bytes, Bytes> record)
newOffsets.put(tp, offset - 1);
});
- Consumer<ConsumerRecord<K, V>> captureOffsetAndExecuteTestHandler =
- record ->
+ TestRecordHandler<K, V> captureOffsetAndExecuteTestHandler =
+ new TestRecordHandler<K, V>(recordHandler)
{
- newOffsets.put(
- new TopicPartition(record.topic(), record.partition()),
- record.offset());
- receivedRecords.add(record);
- consumer.accept(record);
+ @Override
+ public void onNewRecord(ConsumerRecord<K, V> record)
+ {
+ newOffsets.put(
+ new TopicPartition(record.topic(), record.partition()),
+ record.offset());
+ receivedRecords.add(record);
+ }
};
endlessConsumer =
properties.getClientId(),
properties.getTopic(),
kafkaConsumer,
+ rebalanceListener,
captureOffsetAndExecuteTestHandler);
endlessConsumer.start();