MCPcopy Create free account

hub / github.com/ajmalbabu/kafka-clients / functions

Functions42 in github.com/ajmalbabu/kafka-clients

↓ 3 callersMethodgetSchema
()
src/main/java/poc/AvroSupport.java:23
↓ 3 callersMethodsaveOffsetInExternalStore
Overwrite the offset for the topic in an external storage. @param topic - Topic name. @param partition - Partition of the topic. @param offset
src/main/java/poc/OffsetManager.java:30
↓ 2 callersMethodreadOffsetFromExternalStore
@return he last offset + 1 for the provided topic and partition.
src/main/java/poc/OffsetManager.java:50
↓ 2 callersMethodstorageName
(String topic, int partition)
src/main/java/poc/OffsetManager.java:65
↓ 1 callersMethodbyteArrayToData
(Schema schema, byte[] byteData)
src/main/java/poc/AvroSupport.java:48
↓ 1 callersMethodcreateConsumer
()
src/main/java/poc/AtLeastOnceConsumer.java:47
↓ 1 callersMethodcreateConsumer
()
src/main/java/poc/AtMostOnceConsumer.java:65
↓ 1 callersMethodcreateConsumer
()
src/main/java/poc/ExactlyOnceStaticConsumer.java:77
↓ 1 callersMethodcreateConsumer
()
src/main/java/poc/AvroConsumerExample.java:60
↓ 1 callersMethodcreateConsumer
()
src/main/java/poc/ExactlyOnceDynamicConsumer.java:64
↓ 1 callersMethodcreateProducer
()
src/main/java/poc/ProducerExample.java:34
↓ 1 callersMethodcreateProducer
()
src/main/java/poc/AvroProducerExample.java:33
↓ 1 callersMethoddataToByteArray
(Schema schema, GenericRecord datum)
src/main/java/poc/AvroSupport.java:34
↓ 1 callersMethodexecute
()
src/main/java/poc/AtLeastOnceConsumer.java:35
↓ 1 callersMethodexecute
()
src/main/java/poc/AtMostOnceConsumer.java:54
↓ 1 callersMethodgetValue
(GenericRecord genericRecord, String name, Class<T> clazz)
src/main/java/poc/AvroSupport.java:67
↓ 1 callersMethodprocess
()
src/main/java/poc/AtLeastOnceConsumer.java:97
↓ 1 callersMethodprocess
()
src/main/java/poc/AtMostOnceConsumer.java:109
↓ 1 callersMethodprocessRecords
(KafkaConsumer<String, String> consumer)
src/main/java/poc/AtLeastOnceConsumer.java:72
↓ 1 callersMethodprocessRecords
(KafkaConsumer<String, String> consumer)
src/main/java/poc/AtMostOnceConsumer.java:91
↓ 1 callersMethodprocessRecords
Process data and store offset in external store. Best practice is to do these operations atomically. Read class level comments.
src/main/java/poc/ExactlyOnceStaticConsumer.java:117
↓ 1 callersMethodprocessRecords
(KafkaConsumer<String, byte[]> consumer)
src/main/java/poc/AvroConsumerExample.java:37
↓ 1 callersMethodprocessRecords
(KafkaConsumer<String, String> consumer)
src/main/java/poc/ExactlyOnceDynamicConsumer.java:84
↓ 1 callersMethodreadMessages
()
src/main/java/poc/ExactlyOnceStaticConsumer.java:57
↓ 1 callersMethodreadMessages
()
src/main/java/poc/AvroConsumerExample.java:27
↓ 1 callersMethodreadMessages
()
src/main/java/poc/ExactlyOnceDynamicConsumer.java:52
↓ 1 callersMethodrecord
(String name)
src/main/java/poc/AvroProducerExample.java:60
↓ 1 callersMethodregisterConsumerToSpecificPartition
Manually listens for specific topic partition. But, if you are looking for example of how to dynamically listens to partition and want to manually con
src/main/java/poc/ExactlyOnceStaticConsumer.java:104
↓ 1 callersMethodsendMessages
()
src/main/java/poc/ProducerExample.java:23
↓ 1 callersMethodsendMessages
()
src/main/java/poc/AvroProducerExample.java:26
↓ 1 callersMethodsendRecords
(Producer<String, byte[]> producer)
src/main/java/poc/AvroProducerExample.java:48
MethodMyConsumerRebalancerListener
(Consumer<String, String> consumer)
src/main/java/poc/MyConsumerRebalancerListener.java:16
MethodOffsetManager
(String storagePrefix)
src/main/java/poc/OffsetManager.java:19
Methodmain
(String[] str)
src/main/java/poc/ProducerExample.java:15
Methodmain
(String[] str)
src/main/java/poc/AtLeastOnceConsumer.java:26
Methodmain
(String[] str)
src/main/java/poc/AtMostOnceConsumer.java:46
Methodmain
(String[] str)
src/main/java/poc/ExactlyOnceStaticConsumer.java:47
Methodmain
(String[] str)
src/main/java/poc/AvroProducerExample.java:17
Methodmain
(String[] str)
src/main/java/poc/AvroConsumerExample.java:18
Methodmain
(String[] str)
src/main/java/poc/ExactlyOnceDynamicConsumer.java:42
MethodonPartitionsAssigned
(Collection<TopicPartition> partitions)
src/main/java/poc/MyConsumerRebalancerListener.java:28
MethodonPartitionsRevoked
(Collection<TopicPartition> partitions)
src/main/java/poc/MyConsumerRebalancerListener.java:20