Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/dibbhatt/kafka-spark-consumer
/ functions
Functions
161 in github.com/dibbhatt/kafka-spark-consumer
⨍
Functions
161
◇
Types & classes
34
Method
ZkState
(KafkaConfig config)
src/main/java/consumer/kafka/ZkState.java:67
Method
call
(JavaPairRDD<String, Iterable<Long>> po)
src/main/java/consumer/kafka/ProcessedOffsetManager.java:52
Method
call
(Iterator<MessageAndMetadata<E>> it)
src/main/java/consumer/kafka/PartitionOffsetPair.java:41
Method
call
(JavaRDD<MessageAndMetadata<byte[]>> rdd)
src/main/java/consumer/kafka/client/SampleConsumer.java:83
Method
close
()
src/main/java/consumer/kafka/DynamicBrokersReader.java:124
Method
close
()
src/main/java/consumer/kafka/ZkBrokerReader.java:62
Method
equals
(Object obj)
src/main/java/consumer/kafka/Partition.java:47
Method
equals
(Object obj)
src/main/java/consumer/kafka/GlobalPartitionInformation.java:94
Method
fromString
(String host)
src/main/java/consumer/kafka/Broker.java:69
Method
getConnection
(Partition partition)
src/main/java/consumer/kafka/DynamicPartitionConnections.java:82
Method
getConsumer
()
src/main/java/consumer/kafka/MessageAndMetadata.java:49
Method
getCurrentBrokers
()
src/main/java/consumer/kafka/ZkBrokerReader.java:51
Method
getError
(int errorCode)
src/main/java/consumer/kafka/KafkaError.java:43
Method
getInt
(Object o)
src/main/java/consumer/kafka/Utils.java:31
Method
getManager
(Partition partition)
src/main/java/consumer/kafka/PartitionCoordinator.java:32
Method
getManager
(Partition partition)
src/main/java/consumer/kafka/ZkCoordinator.java:140
Method
getMyManagedPartitions
()
src/main/java/consumer/kafka/ZkCoordinator.java:78
Method
getOffset
()
src/main/java/consumer/kafka/MessageAndMetadata.java:65
Method
getOrderedPartitions
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:62
Method
getPartition
()
src/main/java/consumer/kafka/MessageAndMetadata.java:41
Method
getPartitionMap
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:49
Method
getPartitionOffset
( JavaDStream<MessageAndMetadata<T>> unionStreams, Properties props)
src/main/java/consumer/kafka/ProcessedOffsetManager.java:42
Method
getPayload
()
src/main/java/consumer/kafka/MessageAndMetadata.java:57
Method
getRate
()
src/test/java/consumer/kafka/PIDControllerTest.java:100
Method
hashCode
()
src/main/java/consumer/kafka/Partition.java:42
Method
hashCode
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:89
Method
hashCode
()
src/main/java/consumer/kafka/Broker.java:46
Method
lastCommittedOffset
()
src/main/java/consumer/kafka/PartitionManager.java:326
Method
main
(String[] args)
src/main/java/consumer/kafka/client/SampleConsumer.java:127
Method
next
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:78
Method
onBatchCompleted
()
src/test/java/consumer/kafka/PIDControllerTest.java:118
Method
onBatchCompleted
( StreamingListenerBatchCompleted batchCompleted)
src/main/java/consumer/kafka/ReceiverStreamListener.java:110
Method
onBatchStarted
(StreamingListenerBatchStarted arg0)
src/main/java/consumer/kafka/ReceiverStreamListener.java:106
Method
onBatchSubmitted
()
src/test/java/consumer/kafka/PIDControllerTest.java:111
Method
onBatchSubmitted
(StreamingListenerBatchSubmitted arg0)
src/main/java/consumer/kafka/ReceiverStreamListener.java:93
Method
onOutputOperationCompleted
(StreamingListenerOutputOperationCompleted outPutOpsComplete)
src/main/java/consumer/kafka/ReceiverStreamListener.java:82
Method
onOutputOperationStarted
(StreamingListenerOutputOperationStarted outPutOpsStart)
src/main/java/consumer/kafka/ReceiverStreamListener.java:78
Method
onReceiverError
(StreamingListenerReceiverError error)
src/main/java/consumer/kafka/ReceiverStreamListener.java:74
Method
onReceiverStarted
(StreamingListenerReceiverStarted arg0)
src/main/java/consumer/kafka/ReceiverStreamListener.java:70
Method
onReceiverStopped
(StreamingListenerReceiverStopped receiverStopped)
src/main/java/consumer/kafka/ReceiverStreamListener.java:65
Method
onStart
()
src/main/java/consumer/kafka/client/KafkaRangeReceiver.java:67
Method
onStart
()
src/main/java/consumer/kafka/client/KafkaReceiver.java:64
Method
onStop
()
src/main/java/consumer/kafka/client/KafkaRangeReceiver.java:117
Method
onStop
()
src/main/java/consumer/kafka/client/KafkaReceiver.java:109
Method
onStreamingStarted
(StreamingListenerStreamingStarted arg0)
src/main/java/consumer/kafka/ReceiverStreamListener.java:61
Method
persists
(JavaPairDStream<String, Iterable<Long>> partitonOffset, Properties props)
src/main/java/consumer/kafka/ProcessedOffsetManager.java:49
Method
process
(byte[] payload)
src/main/java/consumer/kafka/IdentityMessageHandler.java:24
Method
readJSON
(String path)
src/main/java/consumer/kafka/ZkState.java:111
Method
reportOffsetLag
()
src/main/java/consumer/kafka/PartitionManager.java:145
Method
run
()
src/main/java/consumer/kafka/KafkaSparkConsumer.java:108
Method
setUp
()
src/test/java/consumer/kafka/PIDControllerTest.java:39
Method
tearDown
()
src/test/java/consumer/kafka/PIDControllerTest.java:53
Method
testPIControllerWithNoDelay
()
src/test/java/consumer/kafka/PIDControllerTest.java:86
Method
testPIControllerWithProcessingDelay
()
src/test/java/consumer/kafka/PIDControllerTest.java:58
Method
testPIControllerWithSchedulingDelay
()
src/test/java/consumer/kafka/PIDControllerTest.java:72
Method
toByteArray
(ByteBuffer buffer)
src/main/java/consumer/kafka/Utils.java:44
Method
toScalaSeq
(List<T> list)
src/main/java/consumer/kafka/ScalaUtil.java:49
Method
toString
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:53
Method
uncaughtException
(Thread th, Throwable ex)
src/main/java/consumer/kafka/client/KafkaRangeReceiver.java:89
Method
uncaughtException
(Thread th, Throwable ex)
src/main/java/consumer/kafka/client/KafkaReceiver.java:80
Method
zkIdsPath
(String type)
src/main/java/consumer/kafka/PartitionManager.java:315
← previous
101–161 of 161, ranked by callers