MCPcopy Create free account

hub / github.com/dibbhatt/kafka-spark-consumer / functions

Functions161 in github.com/dibbhatt/kafka-spark-consumer

MethodZkState
(KafkaConfig config)
src/main/java/consumer/kafka/ZkState.java:67
Methodcall
(JavaPairRDD<String, Iterable<Long>> po)
src/main/java/consumer/kafka/ProcessedOffsetManager.java:52
Methodcall
(Iterator<MessageAndMetadata<E>> it)
src/main/java/consumer/kafka/PartitionOffsetPair.java:41
Methodcall
(JavaRDD<MessageAndMetadata<byte[]>> rdd)
src/main/java/consumer/kafka/client/SampleConsumer.java:83
Methodclose
()
src/main/java/consumer/kafka/DynamicBrokersReader.java:124
Methodclose
()
src/main/java/consumer/kafka/ZkBrokerReader.java:62
Methodequals
(Object obj)
src/main/java/consumer/kafka/Partition.java:47
Methodequals
(Object obj)
src/main/java/consumer/kafka/GlobalPartitionInformation.java:94
MethodfromString
(String host)
src/main/java/consumer/kafka/Broker.java:69
MethodgetConnection
(Partition partition)
src/main/java/consumer/kafka/DynamicPartitionConnections.java:82
MethodgetConsumer
()
src/main/java/consumer/kafka/MessageAndMetadata.java:49
MethodgetCurrentBrokers
()
src/main/java/consumer/kafka/ZkBrokerReader.java:51
MethodgetError
(int errorCode)
src/main/java/consumer/kafka/KafkaError.java:43
MethodgetInt
(Object o)
src/main/java/consumer/kafka/Utils.java:31
MethodgetManager
(Partition partition)
src/main/java/consumer/kafka/PartitionCoordinator.java:32
MethodgetManager
(Partition partition)
src/main/java/consumer/kafka/ZkCoordinator.java:140
MethodgetMyManagedPartitions
()
src/main/java/consumer/kafka/ZkCoordinator.java:78
MethodgetOffset
()
src/main/java/consumer/kafka/MessageAndMetadata.java:65
MethodgetOrderedPartitions
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:62
MethodgetPartition
()
src/main/java/consumer/kafka/MessageAndMetadata.java:41
MethodgetPartitionMap
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:49
MethodgetPartitionOffset
( JavaDStream<MessageAndMetadata<T>> unionStreams, Properties props)
src/main/java/consumer/kafka/ProcessedOffsetManager.java:42
MethodgetPayload
()
src/main/java/consumer/kafka/MessageAndMetadata.java:57
MethodgetRate
()
src/test/java/consumer/kafka/PIDControllerTest.java:100
MethodhashCode
()
src/main/java/consumer/kafka/Partition.java:42
MethodhashCode
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:89
MethodhashCode
()
src/main/java/consumer/kafka/Broker.java:46
MethodlastCommittedOffset
()
src/main/java/consumer/kafka/PartitionManager.java:326
Methodmain
(String[] args)
src/main/java/consumer/kafka/client/SampleConsumer.java:127
Methodnext
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:78
MethodonBatchCompleted
()
src/test/java/consumer/kafka/PIDControllerTest.java:118
MethodonBatchCompleted
( StreamingListenerBatchCompleted batchCompleted)
src/main/java/consumer/kafka/ReceiverStreamListener.java:110
MethodonBatchStarted
(StreamingListenerBatchStarted arg0)
src/main/java/consumer/kafka/ReceiverStreamListener.java:106
MethodonBatchSubmitted
()
src/test/java/consumer/kafka/PIDControllerTest.java:111
MethodonBatchSubmitted
(StreamingListenerBatchSubmitted arg0)
src/main/java/consumer/kafka/ReceiverStreamListener.java:93
MethodonOutputOperationCompleted
(StreamingListenerOutputOperationCompleted outPutOpsComplete)
src/main/java/consumer/kafka/ReceiverStreamListener.java:82
MethodonOutputOperationStarted
(StreamingListenerOutputOperationStarted outPutOpsStart)
src/main/java/consumer/kafka/ReceiverStreamListener.java:78
MethodonReceiverError
(StreamingListenerReceiverError error)
src/main/java/consumer/kafka/ReceiverStreamListener.java:74
MethodonReceiverStarted
(StreamingListenerReceiverStarted arg0)
src/main/java/consumer/kafka/ReceiverStreamListener.java:70
MethodonReceiverStopped
(StreamingListenerReceiverStopped receiverStopped)
src/main/java/consumer/kafka/ReceiverStreamListener.java:65
MethodonStart
()
src/main/java/consumer/kafka/client/KafkaRangeReceiver.java:67
MethodonStart
()
src/main/java/consumer/kafka/client/KafkaReceiver.java:64
MethodonStop
()
src/main/java/consumer/kafka/client/KafkaRangeReceiver.java:117
MethodonStop
()
src/main/java/consumer/kafka/client/KafkaReceiver.java:109
MethodonStreamingStarted
(StreamingListenerStreamingStarted arg0)
src/main/java/consumer/kafka/ReceiverStreamListener.java:61
Methodpersists
(JavaPairDStream<String, Iterable<Long>> partitonOffset, Properties props)
src/main/java/consumer/kafka/ProcessedOffsetManager.java:49
Methodprocess
(byte[] payload)
src/main/java/consumer/kafka/IdentityMessageHandler.java:24
MethodreadJSON
(String path)
src/main/java/consumer/kafka/ZkState.java:111
MethodreportOffsetLag
()
src/main/java/consumer/kafka/PartitionManager.java:145
Methodrun
()
src/main/java/consumer/kafka/KafkaSparkConsumer.java:108
MethodsetUp
()
src/test/java/consumer/kafka/PIDControllerTest.java:39
MethodtearDown
()
src/test/java/consumer/kafka/PIDControllerTest.java:53
MethodtestPIControllerWithNoDelay
()
src/test/java/consumer/kafka/PIDControllerTest.java:86
MethodtestPIControllerWithProcessingDelay
()
src/test/java/consumer/kafka/PIDControllerTest.java:58
MethodtestPIControllerWithSchedulingDelay
()
src/test/java/consumer/kafka/PIDControllerTest.java:72
MethodtoByteArray
(ByteBuffer buffer)
src/main/java/consumer/kafka/Utils.java:44
MethodtoScalaSeq
(List<T> list)
src/main/java/consumer/kafka/ScalaUtil.java:49
MethodtoString
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:53
MethoduncaughtException
(Thread th, Throwable ex)
src/main/java/consumer/kafka/client/KafkaRangeReceiver.java:89
MethoduncaughtException
(Thread th, Throwable ex)
src/main/java/consumer/kafka/client/KafkaReceiver.java:80
MethodzkIdsPath
(String type)
src/main/java/consumer/kafka/PartitionManager.java:315
← previous101–161 of 161, ranked by callers