MCPcopy Create free account

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

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

↓ 7 callersMethodclear
()
src/main/java/consumer/kafka/DynamicPartitionConnections.java:100
↓ 7 callersMethodgetKey
()
src/main/java/consumer/kafka/MessageAndMetadata.java:33
↓ 7 callersMethodgetOffset
( KafkaConsumer<byte[], byte[]> consumer, String topic, int partition, boolean forceFromStart)
src/main/java/consumer/kafka/KafkaUtils.java:48
↓ 7 callersMethodtoString
()
src/main/java/consumer/kafka/Broker.java:64
↓ 6 callersMethodgetTopic
()
src/main/java/consumer/kafka/MessageAndMetadata.java:73
↓ 5 callersMethodclose
()
src/main/java/consumer/kafka/IBrokerReader.java:29
↓ 5 callersMethodclose
()
src/main/java/consumer/kafka/ZkState.java:136
↓ 4 callersMethodcalculateRate
( KafkaConfig config, long batchDurationMs, int pollSize, int fillFreqMs, long schedulingDelayMs, long
src/main/java/consumer/kafka/PIDController.java:44
↓ 4 callersMethodcreateStream
( JavaStreamingContext jsc, Properties props, int numberOfReceivers, StorageLevel storageLevel,
src/main/java/consumer/kafka/ReceiverLauncher.java:71
↓ 4 callersMethodequals
(Object obj)
src/main/java/consumer/kafka/Broker.java:51
↓ 4 callersMethodgetBrokerInfo
Get all partitions with their current leaders
src/main/java/consumer/kafka/DynamicBrokersReader.java:56
↓ 4 callersMethodreadBytes
(String path)
src/main/java/consumer/kafka/ZkState.java:124
↓ 4 callersMethodsetFetchRate
(KafkaConfig config, Integer rate)
src/main/java/consumer/kafka/Utils.java:50
↓ 4 callersMethodwriteBytes
(String path, byte[] bytes)
src/main/java/consumer/kafka/ZkState.java:98
↓ 3 callersMethoddoPersists
(List<Tuple2<String, Iterable<Long>>> poList, Properties props)
src/main/java/consumer/kafka/ProcessedOffsetManager.java:95
↓ 3 callersMethodgetPartition
()
src/main/java/consumer/kafka/PartitionManager.java:330
↓ 3 callersMethodnext
()
src/main/java/consumer/kafka/PartitionManager.java:133
↓ 3 callersMethodstart
()
src/main/java/consumer/kafka/client/KafkaReceiver.java:69
↓ 2 callersMethodfetchMessages
( KafkaConfig config, KafkaConsumer<byte[], byte[]> consumer, Partition partition, long offset)
src/main/java/consumer/kafka/KafkaUtils.java:66
↓ 2 callersMethodgetClassTag
Scala 2.10 use ClassTag to replace ClassManifest
src/main/java/consumer/kafka/ScalaUtil.java:35
↓ 2 callersMethodgetCurator
()
src/main/java/consumer/kafka/ZkState.java:62
↓ 2 callersMethodhasNext
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:74
↓ 2 callersMethodinit
()
src/main/java/consumer/kafka/KafkaMessageHandler.java:26
↓ 2 callersMethoditerator
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:70
↓ 2 callersMethodopen
(int partitionId)
src/main/java/consumer/kafka/KafkaSparkConsumer.java:59
↓ 2 callersMethodpartitionPath
()
src/main/java/consumer/kafka/DynamicBrokersReader.java:90
↓ 2 callersMethodrefresh
()
src/main/java/consumer/kafka/PartitionCoordinator.java:33
↓ 2 callersMethodremove
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:83
↓ 2 callersMethodstart
()
src/main/java/consumer/kafka/client/SampleConsumer.java:45
↓ 2 callersMethodtriggerBlockManagerWrite
()
src/main/java/consumer/kafka/PartitionManager.java:166
↓ 2 callersMethodzkPath
(String type)
src/main/java/consumer/kafka/PartitionManager.java:309
↓ 1 callersMethodaddPartition
(int partitionId, Broker broker)
src/main/java/consumer/kafka/GlobalPartitionInformation.java:45
↓ 1 callersMethodassignReceiversToPartitions
(int numberOfReceivers, int numberOfPartition, List<JavaDStream<MessageAndMetadata<E>>> streamsLi
src/main/java/consumer/kafka/ReceiverLauncher.java:125
↓ 1 callersMethodbrokerPath
()
src/main/java/consumer/kafka/DynamicBrokersReader.java:94
↓ 1 callersMethodcalculateOffsetLag
()
src/main/java/consumer/kafka/PartitionManager.java:154
↓ 1 callersMethodclone
()
src/main/java/consumer/kafka/KafkaMessageHandler.java:43
↓ 1 callersMethodclose
()
src/main/java/consumer/kafka/PartitionManager.java:334
↓ 1 callersMethodclose
()
src/main/java/consumer/kafka/KafkaSparkConsumer.java:79
↓ 1 callersMethodcompareTo
(Broker o)
src/main/java/consumer/kafka/Broker.java:82
↓ 1 callersMethodcreateStream
()
src/main/java/consumer/kafka/KafkaSparkConsumer.java:86
↓ 1 callersMethodfetchMessages
(String topic)
src/main/java/consumer/kafka/PartitionManager.java:255
↓ 1 callersMethodfill
()
src/main/java/consumer/kafka/PartitionManager.java:186
↓ 1 callersMethodgetBrokerFor
(Integer partitionId)
src/main/java/consumer/kafka/GlobalPartitionInformation.java:58
↓ 1 callersMethodgetBrokerHost
[zk: localhost:2181(CONNECTED) 56] get /brokers/ids/0 { "host":"localhost", "jmx_port":9999, "port":9092, "version":1 } @param contents @return
src/main/java/consumer/kafka/DynamicBrokersReader.java:135
↓ 1 callersMethodgetCurrentBrokers
()
src/main/java/consumer/kafka/IBrokerReader.java:28
↓ 1 callersMethodgetFetchSize
()
src/main/java/consumer/kafka/PartitionManager.java:287
↓ 1 callersMethodgetFetchSize
()
src/main/java/consumer/kafka/ReceiverStreamListener.java:147
↓ 1 callersMethodgetId
()
src/main/java/consumer/kafka/Partition.java:65
↓ 1 callersMethodgetKafkaOffset
()
src/main/java/consumer/kafka/PartitionManager.java:137
↓ 1 callersMethodgetLeaderFor
get /brokers/topics/distributedTopic/partitions/1/state { "controller_epoch":4, "isr":[ 1, 0 ], "leader":1, "leader_epoch":1, "version":1 } @param pa
src/main/java/consumer/kafka/DynamicBrokersReader.java:106
↓ 1 callersMethodgetMaximum
(Iterable<T> values)
src/main/java/consumer/kafka/ProcessedOffsetManager.java:104
↓ 1 callersMethodgetMessageAndMetadataClassTag
()
src/main/java/consumer/kafka/ScalaUtil.java:39
↓ 1 callersMethodgetMyManagedPartitions
()
src/main/java/consumer/kafka/PartitionCoordinator.java:31
↓ 1 callersMethodgetNumPartitions
()
src/main/java/consumer/kafka/DynamicBrokersReader.java:80
↓ 1 callersMethodgetNumPartitions
(ZkState zkState, String topic)
src/main/java/consumer/kafka/ReceiverLauncher.java:157
↓ 1 callersMethodgetProperties
()
src/main/java/consumer/kafka/KafkaConfig.java:162
↓ 1 callersMethodgetTuple2ClassTag
()
src/main/java/consumer/kafka/ScalaUtil.java:44
↓ 1 callersMethodgetZKPath
(Properties props)
src/main/java/consumer/kafka/ProcessedOffsetManager.java:132
↓ 1 callersMethodhandle
(long offset, Partition partition, String topic, String consumer, byte[] payload)
src/main/java/consumer/kafka/KafkaMessageHandler.java:30
↓ 1 callersMethodlaunch
( StreamingContext ssc, Properties pros, int numberOfReceivers, StorageLevel storageLe
src/main/java/consumer/kafka/ReceiverLauncher.java:47
↓ 1 callersMethodnewCurator
(Map<String, String> stateConf)
src/main/java/consumer/kafka/ZkState.java:47
↓ 1 callersMethodpartitionPath
(String topic)
src/main/java/consumer/kafka/ReceiverLauncher.java:168
↓ 1 callersMethodpersistProcessedOffsets
(Properties props, Map<String, Long> partitionOffsetMap)
src/main/java/consumer/kafka/ProcessedOffsetManager.java:114
↓ 1 callersMethodpersistsPartition
(JavaRDD<MessageAndMetadata<T>> rdd, Properties props)
src/main/java/consumer/kafka/ProcessedOffsetManager.java:61
↓ 1 callersMethodprocess
(byte[] payload)
src/main/java/consumer/kafka/KafkaMessageHandler.java:41
↓ 1 callersMethodprocessedPath
(String topic, int partition, Properties props)
src/main/java/consumer/kafka/ProcessedOffsetManager.java:151
↓ 1 callersMethodratePath
()
src/main/java/consumer/kafka/PartitionManager.java:304
↓ 1 callersMethodratePath
()
src/main/java/consumer/kafka/ReceiverStreamListener.java:168
↓ 1 callersMethodrefresh
()
src/main/java/consumer/kafka/ZkCoordinator.java:88
↓ 1 callersMethodregister
(Partition partition, String topic)
src/main/java/consumer/kafka/DynamicPartitionConnections.java:68
↓ 1 callersMethodrun
()
src/main/java/consumer/kafka/client/SampleConsumer.java:50
↓ 1 callersMethodsetConsumer
(String consumer)
src/main/java/consumer/kafka/MessageAndMetadata.java:53
↓ 1 callersMethodsetKey
(byte[] key)
src/main/java/consumer/kafka/MessageAndMetadata.java:37
↓ 1 callersMethodsetOffset
(long offset)
src/main/java/consumer/kafka/MessageAndMetadata.java:69
↓ 1 callersMethodsetPartition
(Partition partition)
src/main/java/consumer/kafka/MessageAndMetadata.java:45
↓ 1 callersMethodsetPayload
(E msg)
src/main/java/consumer/kafka/MessageAndMetadata.java:61
↓ 1 callersMethodsetTopic
(String topic)
src/main/java/consumer/kafka/MessageAndMetadata.java:77
↓ 1 callersMethodsetZkCoordinator
()
src/main/java/consumer/kafka/PartitionManager.java:114
↓ 1 callersMethodstart
()
src/main/java/consumer/kafka/client/KafkaRangeReceiver.java:76
↓ 1 callersMethodtoString
()
src/main/java/consumer/kafka/Partition.java:60
↓ 1 callersMethodunregister
(Broker host, int partition)
src/main/java/consumer/kafka/DynamicPartitionConnections.java:90
↓ 1 callersMethodwriteJSON
(String path, Map<Object, Object> data)
src/main/java/consumer/kafka/ZkState.java:93
↓ 1 callersMethodzkCordPath
(String type)
src/main/java/consumer/kafka/PartitionManager.java:321
MethodBroker
(String host, int port)
src/main/java/consumer/kafka/Broker.java:37
MethodConnectionInfo
(KafkaConsumer<byte[], byte[]> consumer, String topic, int partition)
src/main/java/consumer/kafka/DynamicPartitionConnections.java:48
MethodDynamicBrokersReader
(KafkaConfig config, ZkState zkState)
src/main/java/consumer/kafka/DynamicBrokersReader.java:47
MethodDynamicPartitionConnections
( KafkaConfig config, IBrokerReader brokerReader)
src/main/java/consumer/kafka/DynamicPartitionConnections.java:61
MethodFailedFetchException
(String message)
src/main/java/consumer/kafka/FailedFetchException.java:30
MethodGlobalPartitionInformation
()
src/main/java/consumer/kafka/GlobalPartitionInformation.java:41
MethodKafkaConfig
(Properties props)
src/main/java/consumer/kafka/KafkaConfig.java:66
MethodKafkaRangeReceiver
(KafkaConfig config, Set<Integer> partitionSet, Ka
src/main/java/consumer/kafka/client/KafkaRangeReceiver.java:50
MethodKafkaReceiver
(KafkaConfig config, int partitionId, KafkaMessageHandler me
src/main/java/consumer/kafka/client/KafkaReceiver.java:47
MethodKafkaSparkConsumer
( KafkaConfig config, ZkState zkState, Receiver<MessageAndMetadata<E>> rec
src/main/java/consumer/kafka/KafkaSparkConsumer.java:48
MethodOutOfRangeException
(String message)
src/main/java/consumer/kafka/OutOfRangeException.java:30
MethodPIDController
(double proportional, double integral, double derivative)
src/main/java/consumer/kafka/PIDController.java:35
MethodPartition
(Broker host, int partition)
src/main/java/consumer/kafka/Partition.java:37
MethodPartitionManager
( DynamicPartitionConnections connections, ZkState state, KafkaConfig kafk
src/main/java/consumer/kafka/PartitionManager.java:63
MethodReceiverStreamListener
(KafkaConfig config,long batchDuration)
src/main/java/consumer/kafka/ReceiverStreamListener.java:48
MethodZkBrokerReader
(KafkaConfig config, ZkState zkState)
src/main/java/consumer/kafka/ZkBrokerReader.java:43
MethodZkCoordinator
( DynamicPartitionConnections connections, KafkaConfig config, ZkState state, int partitionId,
src/main/java/consumer/kafka/ZkCoordinator.java:58
next →1–100 of 161, ranked by callers