X-Git-Url: https://juplo.de/gitweb/?a=blobdiff_plain;f=src%2Ftest%2Fjava%2Fde%2Fjuplo%2Fkafka%2FTestRecordHandler.java;h=37d3f659f9c78b24a0a62da4e26b8a1169651f68;hb=595eab489c638b07072f6ec7e3c6f52000295931;hp=4047093b03de41932c10d4e6cd7acc876ab066e4;hpb=2d84eda74475aaffff11ddfebe56d309b9aff2e9;p=demos%2Fkafka%2Ftraining diff --git a/src/test/java/de/juplo/kafka/TestRecordHandler.java b/src/test/java/de/juplo/kafka/TestRecordHandler.java index 4047093..37d3f65 100644 --- a/src/test/java/de/juplo/kafka/TestRecordHandler.java +++ b/src/test/java/de/juplo/kafka/TestRecordHandler.java @@ -4,15 +4,26 @@ import lombok.RequiredArgsConstructor; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.TopicPartition; +import java.util.Map; +import java.util.Set; + @RequiredArgsConstructor -public abstract class TestRecordHandler implements RecordHandler +public class TestRecordHandler implements RecordHandler { private final RecordHandler handler; + Map seenOffsets; + Set> receivedRecords; - public abstract void onNewRecord(ConsumerRecord record); + public void onNewRecord(ConsumerRecord record) + { + seenOffsets.put( + new TopicPartition(record.topic(), record.partition()), + record.offset()); + receivedRecords.add(record); + } @Override public void accept(ConsumerRecord record) @@ -20,22 +31,4 @@ public abstract class TestRecordHandler implements RecordHandler this.onNewRecord(record); handler.accept(record); } - @Override - - public void beforeNextPoll() - { - handler.beforeNextPoll(); - } - - @Override - public void onPartitionAssigned(TopicPartition tp) - { - handler.onPartitionAssigned(tp); - } - - @Override - public void onPartitionRevoked(TopicPartition tp) - { - handler.onPartitionRevoked(tp); - } }