X-Git-Url: https://juplo.de/gitweb/?a=blobdiff_plain;f=src%2Ftest%2Fjava%2Fde%2Fjuplo%2Fkafka%2FApplicationTest.java;fp=src%2Ftest%2Fjava%2Fde%2Fjuplo%2Fkafka%2FApplicationTest.java;h=d3ff3b1d3927c5ea54719f177939d9e30d683f22;hb=2bf77d19d90e7356e1a7c6e13202971fd1b9897b;hp=0000000000000000000000000000000000000000;hpb=a126175808eeec80355deb8eb3f4ef7e85e84780;p=demos%2Fkafka%2Ftraining diff --git a/src/test/java/de/juplo/kafka/ApplicationTest.java b/src/test/java/de/juplo/kafka/ApplicationTest.java new file mode 100644 index 0000000..d3ff3b1 --- /dev/null +++ b/src/test/java/de/juplo/kafka/ApplicationTest.java @@ -0,0 +1,57 @@ +package de.juplo.kafka; + +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.serialization.LongSerializer; +import org.apache.kafka.common.serialization.StringSerializer; +import org.apache.kafka.common.utils.Bytes; + +import java.util.Set; +import java.util.function.Consumer; + + +public class ApplicationTest extends GenericApplicationTest +{ + public ApplicationTest() + { + super( + new StringSerializer(), + new RecordGenerator<> () + { + final StringSerializer stringSerializer = new StringSerializer(); + final LongSerializer longSerializer = new LongSerializer(); + + + @Override + public void generate( + int numberOfMessagesToGenerate, + Set poisonPills, + Consumer> messageSender) + { + int i = 0; + + for (int partition = 0; partition < 10; partition++) + { + for (int key = 0; key < 10; key++) + { + if (++i > numberOfMessagesToGenerate) + return; + + Bytes value = + poisonPills.contains(i) + ? new Bytes(stringSerializer.serialize(TOPIC, "BOOM!")) + : new Bytes(longSerializer.serialize(TOPIC, (long)i)); + + ProducerRecord record = + new ProducerRecord<>( + TOPIC, + partition, + Integer.toString(partition*10+key%2), + value); + + messageSender.accept(record); + } + } + } + }); + } +}