Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/bjpublic/apache-kafka-with-java
/ functions
Functions
138 in github.com/bjpublic/apache-kafka-with-java
⨍
Functions
138
◇
Types & classes
65
↓ 140 callers
Method
put
(Collection<SinkRecord> records)
Chapter3/3.6 kafka-connector/simple-sink-connector/src/main/java/com/example/SingleFileSinkTask.java:37
↓ 14 callers
Method
poll
()
Chapter3/3.6 kafka-connector/simple-source-connector/src/main/java/com/example/SingleFileSourceTask.java:65
↓ 12 callers
Method
close
()
Chapter3/3.5 kafka-streams/simple-kafka-processor/src/main/java/com/example/FilterProcessor.java:24
↓ 11 callers
Method
put
(Collection<SinkRecord> records)
Chapter5/5.1 web-page-event-pipeline/elasticsearch-kafka-connector/src/main/java/com/pipeline/ElasticSearchSinkTask.java:52
↓ 10 callers
Method
flush
(Map<TopicPartition, OffsetAndMetadata> offsets)
Chapter3/3.6 kafka-connector/simple-sink-connector/src/main/java/com/example/SingleFileSinkTask.java:48
↓ 8 callers
Method
close
()
Chapter3/3.4.1 kafka-producer/kafka-producer-custom-partitioner/src/main/java/com/example/CustomPartitioner.java:33
↓ 7 callers
Method
start
(Map<String, String> props)
Chapter3/3.6 kafka-connector/simple-sink-connector/src/main/java/com/example/SingleFileSinkTask.java:25
↓ 5 callers
Method
partition
(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster
Chapter3/3.4.1 kafka-producer/kafka-producer-custom-partitioner/src/main/java/com/example/CustomPartitioner.java:14
↓ 5 callers
Method
run
(String... args)
Chapter4/4.4 spring-kafka/spring-kafka-producer/src/main/java/com/example/SpringProducerApplication.java:22
↓ 2 callers
Method
getMetricName
(String value)
Chapter5/5.2 metric-log-pipeline-streams/metric-kafka-streams/src/main/java/com/pipeline/MetricJsonUtils.java:13
↓ 2 callers
Method
save
(int partitionNo)
Chapter5/5.1 web-page-event-pipeline/kafka-multi-consumer-thread-hdfs-save/src/main/java/com/pipeline/consumer/ConsumerWorker.java:81
↓ 1 callers
Method
addHdfsFileBuffer
(ConsumerRecord<String, String> record)
Chapter5/5.1 web-page-event-pipeline/kafka-multi-consumer-thread-hdfs-save/src/main/java/com/pipeline/consumer/ConsumerWorker.java:60
↓ 1 callers
Method
checkFlushCount
(int partitionNo)
Chapter5/5.1 web-page-event-pipeline/kafka-multi-consumer-thread-hdfs-save/src/main/java/com/pipeline/consumer/ConsumerWorker.java:73
↓ 1 callers
Method
getDefaultPartitionSize
()
Chapter4/4.3 kafka-consumer/kafka-multi-consumer-thread-by-partition/src/main/java/com/example/MultiConsumerThreadByPartition.java:66
↓ 1 callers
Method
getHostTimestamp
(String value)
Chapter5/5.2 metric-log-pipeline-streams/metric-kafka-streams/src/main/java/com/pipeline/MetricJsonUtils.java:18
↓ 1 callers
Method
getLines
(long readLine)
Chapter3/3.6 kafka-connector/simple-source-connector/src/main/java/com/example/SingleFileSourceTask.java:87
↓ 1 callers
Method
getPartitionSize
(String topic)
Chapter4/4.3 kafka-consumer/kafka-multi-consumer-thread-by-partition/src/main/java/com/example/MultiConsumerThreadByPartition.java:47
↓ 1 callers
Method
getTotalCpuPercent
(String value)
Chapter5/5.2 metric-log-pipeline-streams/metric-kafka-streams/src/main/java/com/pipeline/MetricJsonUtils.java:8
↓ 1 callers
Method
printKeyValueStoreData
()
Chapter3/3.5 kafka-streams/queryable-store/src/main/java/com/example/QueryableStore.java:59
↓ 1 callers
Method
run
()
Chapter5/5.1 web-page-event-pipeline/kafka-multi-consumer-thread-hdfs-save/src/main/java/com/pipeline/consumer/ConsumerWorker.java:37
↓ 1 callers
Method
run
(String... args)
Chapter4/4.4 spring-kafka/spring-kafka-template-producer/src/main/java/com/example/SpringProducerApplication.java:27
↓ 1 callers
Method
saveBufferToHdfsFile
(Set<TopicPartition> partitions)
Chapter5/5.1 web-page-event-pipeline/kafka-multi-consumer-thread-hdfs-save/src/main/java/com/pipeline/consumer/ConsumerWorker.java:69
↓ 1 callers
Method
saveRemainBufferToHdfsFile
()
Chapter5/5.1 web-page-event-pipeline/kafka-multi-consumer-thread-hdfs-save/src/main/java/com/pipeline/consumer/ConsumerWorker.java:98
↓ 1 callers
Method
start
(Map<String, String> props)
Chapter5/5.1 web-page-event-pipeline/elasticsearch-kafka-connector/src/main/java/com/pipeline/ElasticSearchSinkTask.java:39
Method
ConsumerWorker
(Properties prop, String topic, int number)
Chapter5/5.1 web-page-event-pipeline/kafka-multi-consumer-thread-hdfs-save/src/main/java/com/pipeline/consumer/ConsumerWorker.java:30
Method
ConsumerWorker
(Properties prop, String topic, int number)
Chapter4/4.3 kafka-consumer/kafka-multi-consumer-thread-by-partition/src/main/java/com/example/ConsumerWorker.java:23
Method
ConsumerWorker
(String recordValue)
Chapter4/4.3 kafka-consumer/kafka-consumer-with-multi-worker-thread/src/main/java/com/example/ConsumerWorker.java:11
Method
ConsumerWorker
(Properties prop, String topic, int number)
Chapter4/4.3 kafka-consumer/kafka-multi-consumer-thread/src/main/java/com/example/ConsumerWorker.java:20
Method
ElasticSearchSinkConnectorConfig
(Map<String, String> props)
Chapter5/5.1 web-page-event-pipeline/elasticsearch-kafka-connector/src/main/java/com/pipeline/config/ElasticSearchSinkConnectorConfig.java:27
Method
ProduceController
(KafkaTemplate<String, String> kafkaTemplate)
Chapter5/5.1 web-page-event-pipeline/kafka-spring-producer-with-rest-controller/src/main/java/com/pipeline/ProduceController.java:23
Method
SingleFileSinkConnectorConfig
(Map<String, String> props)
Chapter3/3.6 kafka-connector/simple-sink-connector/src/main/java/com/example/SingleFileSinkConnectorConfig.java:21
Method
SingleFileSourceConnectorConfig
(Map<String, String> props)
Chapter3/3.6 kafka-connector/simple-source-connector/src/main/java/com/example/SingleFileSourceConnectorConfig.java:31
Method
UserEventVO
(String timestamp, String userAgent, String colorName, String userName)
Chapter5/5.1 web-page-event-pipeline/kafka-spring-producer-with-rest-controller/src/main/java/com/pipeline/UserEventVO.java:5
Method
batchListener
(ConsumerRecords<String, String> records)
Chapter4/4.4 spring-kafka/spring-kafka-batch-listener/src/main/java/com/example/SpringConsumerApplication.java:21
Method
commitListener
(ConsumerRecords<String, String> records, Acknowledgment ack)
Chapter4/4.4 spring-kafka/spring-kafka-commit-listener/src/main/java/com/example/SpringConsumerApplication.java:22
Method
concurrentBatchListener
(ConsumerRecords<String, String> records)
Chapter4/4.4 spring-kafka/spring-kafka-batch-listener/src/main/java/com/example/SpringConsumerApplication.java:33
Method
concurrentTopicListener
(String messageValue)
Chapter4/4.4 spring-kafka/spring-kafka-record-listener/src/main/java/com/example/SpringConsumerApplication.java:43
Method
config
()
Chapter3/3.6 kafka-connector/simple-source-connector/src/main/java/com/example/SingleFileSourceConnector.java:53
Method
config
()
Chapter3/3.6 kafka-connector/simple-sink-connector/src/main/java/com/example/SingleFileSinkConnector.java:49
Method
config
()
Chapter5/5.1 web-page-event-pipeline/elasticsearch-kafka-connector/src/main/java/com/pipeline/ElasticSearchSinkConnector.java:54
Method
configure
(Map<String, ?> configs)
Chapter3/3.4.1 kafka-producer/kafka-producer-custom-partitioner/src/main/java/com/example/CustomPartitioner.java:30
Method
consumerCommitListener
(ConsumerRecords<String, String> records, Consumer<String, String> consumer)
Chapter4/4.4 spring-kafka/spring-kafka-commit-listener/src/main/java/com/example/SpringConsumerApplication.java:28
Method
customContainerFactory
()
Chapter4/4.4 spring-kafka/spring-kafka-listener-container/src/main/java/com/example/ListenerContainerConfiguration.java:21
Method
customKafkaTemplate
()
Chapter4/4.4 spring-kafka/spring-kafka-template-producer/src/main/java/com/example/KafkaTemplateConfiguration.java:14
Method
customListener
(String data)
Chapter4/4.4 spring-kafka/spring-kafka-listener-container/src/main/java/com/example/SpringConsumerApplication.java:18
Method
flush
(Map<TopicPartition, OffsetAndMetadata> offsets)
Chapter5/5.1 web-page-event-pipeline/elasticsearch-kafka-connector/src/main/java/com/pipeline/ElasticSearchSinkTask.java:82
Method
init
(ProcessorContext context)
Chapter3/3.5 kafka-streams/simple-kafka-processor/src/main/java/com/example/FilterProcessor.java:11
Method
listenSpecificPartition
(ConsumerRecord<String, String> record)
Chapter4/4.4 spring-kafka/spring-kafka-record-listener/src/main/java/com/example/SpringConsumerApplication.java:50
Method
main
(String[] args)
Chapter6/6.1 confluent-kafka/confluent-kafka-producer/src/main/java/com/example/SimpleProducer.java:23
Method
main
(String[] args)
Chapter6/6.1 confluent-kafka/confluent-kafka-consumer/src/main/java/com/example/SimpleConsumer.java:28
Method
main
(String[] args)
Chapter3/3.5 kafka-streams/kstream-globalktable-join/src/main/java/com/example/KStreamJoinGlobalKTable.java:20
Method
main
(String[] args)
Chapter3/3.5 kafka-streams/queryable-store/src/main/java/com/example/QueryableStore.java:29
Method
main
(String[] args)
Chapter3/3.5 kafka-streams/simple-kafka-processor/src/main/java/com/example/SimpleKafkaProcessor.java:17
Method
main
(String[] args)
Chapter3/3.5 kafka-streams/kstream-count/src/main/java/com/example/KStreamCountApplication.java:22
Method
main
(String[] args)
Chapter3/3.5 kafka-streams/kafka-streams-filter/src/main/java/com/example/StreamsFilter.java:18
Method
main
(String[] args)
Chapter3/3.5 kafka-streams/kstream-ktable-join/src/main/java/com/example/KStreamJoinKTable.java:20
Method
main
(String[] args)
Chapter3/3.5 kafka-streams/simple-kafka-streams/src/main/java/com/example/SimpleStreamApplication.java:18
Method
main
(String[] args)
Chapter3/3.4.1 kafka-producer/kafka-producer-key-value/src/main/java/com/example/ProducerWithKeyValue.java:13
Method
main
(String[] args)
Chapter3/3.4.1 kafka-producer/kafka-producer-sync-callback/src/main/java/com/example/ProducerWithSyncCallback.java:18
Method
main
(String[] args)
Chapter3/3.4.1 kafka-producer/kafka-producer-exact-partition/src/main/java/com/example/ProducerExactPartition.java:13
Method
main
(String[] args)
Chapter3/3.4.1 kafka-producer/simple-kafka-producer/src/main/java/com/example/SimpleProducer.java:17
Method
main
(String[] args)
Chapter3/3.4.1 kafka-producer/kafka-producer-async-callback/src/main/java/com/example/ProducerWithAsyncCallback.java:13
Method
main
(String[] args)
Chapter3/3.4.1 kafka-producer/kafka-producer-custom-partitioner/src/main/java/com/example/ProducerWithCustomPartitioner.java:13
Method
main
(String[] args)
Chapter3/3.4.3 kafka-admin/kafka-admin-client/src/main/java/com/example/KafkaAdminClient.java:19
Method
main
(String[] args)
Chapter3/3.4.2 kafka-consumer/kafka-consumer-sync-commit/src/main/java/com/example/ConsumerWithSyncCommit.java:18
Method
main
(String[] args)
Chapter3/3.4.2 kafka-consumer/simple-kafka-consumer/src/main/java/com/example/SimpleConsumer.java:21
Method
main
(String[] args)
Chapter3/3.4.2 kafka-consumer/kafka-consumer-auto-commit/src/main/java/com/example/ConsumerWithAutoCommit.java:21
Method
main
(String[] args)
Chapter3/3.4.2 kafka-consumer/kafka-consumer-exact-partition/src/main/java/com/example/ConsumerWithExactPartition.java:19
Method
main
(String[] args)
Chapter3/3.4.2 kafka-consumer/kafka-consumer-sync-offset-commit/src/main/java/com/example/ConsumerWithSyncOffsetCommit.java:21
Method
main
(String[] args)
Chapter3/3.4.2 kafka-consumer/kafka-consumer-sync-offset-commit-shutdown-hook/src/main/java/com/example/ConsumerWithSyncOffsetCommit.java:21
Method
main
(String[] args)
Chapter3/3.4.2 kafka-consumer/kafka-consumer-async-commit/src/main/java/com/example/ConsumerWithASyncCommit.java:20
Method
main
(String[] args)
Chapter3/3.4.2 kafka-consumer/kafka-consumer-rebalance-listener/src/main/java/com/example/ConsumerWithRebalanceListener.java:20
Method
main
(String[] args)
Chapter5/5.1 web-page-event-pipeline/kafka-multi-consumer-thread-hdfs-save/src/main/java/com/pipeline/HdfsSinkApplication.java:23
Method
main
(String[] args)
Chapter5/5.1 web-page-event-pipeline/kafka-spring-producer-with-rest-controller/src/main/java/com/pipeline/RestApiProducer.java:8
Method
main
(final String[] args)
Chapter5/5.2 metric-log-pipeline-streams/metric-kafka-streams/src/main/java/com/pipeline/MetricStreams.java:14
Method
main
(String[] args)
Chapter4/4.4 spring-kafka/spring-kafka-commit-listener/src/main/java/com/example/SpringConsumerApplication.java:17
Method
main
(String[] args)
Chapter4/4.4 spring-kafka/spring-kafka-template-producer/src/main/java/com/example/SpringProducerApplication.java:22
Method
main
(String[] args)
Chapter4/4.4 spring-kafka/spring-kafka-listener-container/src/main/java/com/example/SpringConsumerApplication.java:13
Method
main
(String[] args)
Chapter4/4.4 spring-kafka/spring-kafka-producer/src/main/java/com/example/SpringProducerApplication.java:17
Method
main
(String[] args)
Chapter4/4.4 spring-kafka/spring-kafka-record-listener/src/main/java/com/example/SpringConsumerApplication.java:17
Method
main
(String[] args)
Chapter4/4.4 spring-kafka/spring-kafka-batch-listener/src/main/java/com/example/SpringConsumerApplication.java:16
Method
main
(String[] args)
Chapter4/4.2.2 idempotence-producer/src/main/java/com/example/IdempotenceProducer.java:17
Method
main
(String[] args)
Chapter4/4.2.3 transaction-producer/src/main/java/com/example/TransactionProducer.java:17
Method
main
(String[] args)
Chapter4/4.3 kafka-consumer/kafka-multi-consumer-thread-by-partition/src/main/java/com/example/MultiConsumerThreadByPartition.java:27
Method
main
(String[] args)
Chapter4/4.3 kafka-consumer/kafka-consumer-with-multi-worker-thread/src/main/java/com/example/ConsumerWithMultiWorkerThread.java:23
Method
main
(String[] args)
Chapter4/4.3 kafka-consumer/kafka-multi-consumer-thread/src/main/java/com/example/MultiConsumerThread.java:17
Method
main
(String[] args)
Chapter4/4.2.3 transaction-consumer/src/main/java/com/example/TransactionConsumer.java:21
Method
onComplete
(Map<TopicPartition, OffsetAndMetadata> offsets, Exception e)
Chapter3/3.4.2 kafka-consumer/kafka-consumer-async-commit/src/main/java/com/example/ConsumerWithASyncCommit.java:37
Method
onCompletion
(RecordMetadata recordMetadata, Exception e)
Chapter3/3.4.1 kafka-producer/kafka-producer-async-callback/src/main/java/com/example/ProducerCallback.java:11
Method
onFailure
(Exception e)
Chapter5/5.1 web-page-event-pipeline/elasticsearch-kafka-connector/src/main/java/com/pipeline/ElasticSearchSinkTask.java:74
Method
onFailure
(Throwable ex)
Chapter5/5.1 web-page-event-pipeline/kafka-spring-producer-with-rest-controller/src/main/java/com/pipeline/ProduceController.java:43
Method
onFailure
(KafkaProducerException ex)
Chapter4/4.4 spring-kafka/spring-kafka-template-producer/src/main/java/com/example/SpringProducerApplication.java:36
Method
onPartitionsAssigned
(Collection<TopicPartition> partitions)
Chapter3/3.4.2 kafka-consumer/kafka-consumer-rebalance-listener/src/main/java/com/example/RebalanceListener.java:13
Method
onPartitionsAssigned
(Collection<TopicPartition> partitions)
Chapter4/4.4 spring-kafka/spring-kafka-listener-container/src/main/java/com/example/ListenerContainerConfiguration.java:43
Method
onPartitionsLost
(Collection<TopicPartition> partitions)
Chapter4/4.4 spring-kafka/spring-kafka-listener-container/src/main/java/com/example/ListenerContainerConfiguration.java:48
Method
onPartitionsRevoked
(Collection<TopicPartition> partitions)
Chapter3/3.4.2 kafka-consumer/kafka-consumer-rebalance-listener/src/main/java/com/example/RebalanceListener.java:18
Method
onPartitionsRevokedAfterCommit
(Consumer<?, ?> consumer, Collection<TopicPartition> partitions)
Chapter4/4.4 spring-kafka/spring-kafka-listener-container/src/main/java/com/example/ListenerContainerConfiguration.java:38
Method
onPartitionsRevokedBeforeCommit
(Consumer<?, ?> consumer, Collection<TopicPartition> partitions)
Chapter4/4.4 spring-kafka/spring-kafka-listener-container/src/main/java/com/example/ListenerContainerConfiguration.java:33
Method
onResponse
(BulkResponse bulkResponse)
Chapter5/5.1 web-page-event-pipeline/elasticsearch-kafka-connector/src/main/java/com/pipeline/ElasticSearchSinkTask.java:65
Method
onSuccess
(SendResult<String, String> result)
Chapter5/5.1 web-page-event-pipeline/kafka-spring-producer-with-rest-controller/src/main/java/com/pipeline/ProduceController.java:38
next →
1–100 of 138, ranked by callers