Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/apache/kafka
/ functions
Functions
1,332 in github.com/apache/kafka
⨍
Functions
1,332
◇
Types & classes
265
↓ 1 callers
Method
fmtNodeIds
(Node[] nodes)
clients/src/main/java/org/apache/kafka/common/PartitionInfo.java:81
↓ 1 callers
Method
forId
(int id)
clients/src/main/java/org/apache/kafka/common/record/CompressionType.java:35
↓ 1 callers
Method
forId
(int id)
clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java:63
↓ 1 callers
Method
freeUp
Attempt to ensure we have at least the requested number of bytes of memory for allocation by deallocating pooled buffers (if needed)
clients/src/main/java/org/apache/kafka/clients/producer/internals/BufferPool.java:187
↓ 1 callers
Method
fromBin
(int b)
clients/src/main/java/org/apache/kafka/common/metrics/stats/Histogram.java:102
↓ 1 callers
Method
fromByte
(byte flg)
clients/src/main/java/org/apache/kafka/common/message/KafkaLZ4BlockOutputStream.java:263
↓ 1 callers
Method
fromByte
(byte bd)
clients/src/main/java/org/apache/kafka/common/message/KafkaLZ4BlockOutputStream.java:354
↓ 1 callers
Method
generateData
()
examples/src/main/java/kafka/examples/SimpleConsumerDemo.java:44
↓ 1 callers
Method
generateOffsets
()
contrib/hadoop-consumer/src/main/java/kafka/etl/impl/DataGenerator.java:103
↓ 1 callers
Method
getAttribute
(String name)
clients/src/main/java/org/apache/kafka/common/metrics/JmxReporter.java:160
↓ 1 callers
Method
getBoolean
get boolean value with default value @param key @param defaultValue @return boolean value @throws Exception if value is not of type boolean or string
contrib/hadoop-consumer/src/main/java/kafka/etl/Props.java:260
↓ 1 callers
Method
getBoolean
(String key)
clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java:77
↓ 1 callers
Function
getCSVFileNameFromMetricsMbeanName
(mbeanName)
system_test/utils/metrics.py:65
↓ 1 callers
Method
getChecksum
()
contrib/hadoop-consumer/src/main/java/kafka/etl/KafkaETLKey.java:66
↓ 1 callers
Method
getClientBufferSize
(Props props)
contrib/hadoop-consumer/src/main/java/kafka/etl/KafkaETLContext.java:262
↓ 1 callers
Method
getClientTimeout
(Props props)
contrib/hadoop-consumer/src/main/java/kafka/etl/KafkaETLContext.java:266
↓ 1 callers
Method
getData
(Message message)
contrib/hadoop-consumer/src/main/java/kafka/etl/impl/SimpleKafkaETLMapper.java:43
↓ 1 callers
Method
getFieldOrDefault
Return the value of the given pre-validated field, or if the value is missing return the default value. @param field The field for which to get the d
clients/src/main/java/org/apache/kafka/common/protocol/types/Struct.java:52
↓ 1 callers
Method
getJobConf
Helper function to initialize a job configuration
contrib/hadoop-consumer/src/main/java/kafka/etl/KafkaETLJob.java:61
↓ 1 callers
Method
getMBeanName
@param metricName @return standard JMX MBean name in the following format domainName:type=metricType,key1=val1,key2=val2
clients/src/main/java/org/apache/kafka/common/metrics/JmxReporter.java:101
↓ 1 callers
Method
getNext
(KafkaETLKey key, BytesWritable value)
contrib/hadoop-consumer/src/main/java/kafka/etl/KafkaETLContext.java:134
↓ 1 callers
Method
getOffset
()
contrib/hadoop-consumer/src/main/java/kafka/etl/KafkaETLKey.java:62
↓ 1 callers
Method
getOffsetRange
Get offset ranges
contrib/hadoop-consumer/src/main/java/kafka/etl/KafkaETLContext.java:218
↓ 1 callers
Method
getOutputPath
(JobContext job)
contrib/hadoop-producer/src/main/java/kafka/bridge/hadoop/KafkaOutputFormat.java:75
↓ 1 callers
Method
getPropsFromJob
(Configuration conf)
contrib/hadoop-consumer/src/main/java/kafka/etl/KafkaETLUtils.java:111
↓ 1 callers
Method
getStringSerDeser
(String encoder)
clients/src/test/java/org/apache/kafka/common/serialization/SerializationTest.java:57
↓ 1 callers
Method
getTags
(String... keyValue)
clients/src/main/java/org/apache/kafka/common/MetricName.java:82
↓ 1 callers
Method
getTotalBytes
()
contrib/hadoop-consumer/src/main/java/kafka/etl/KafkaETLContext.java:78
↓ 1 callers
Method
getURI
()
contrib/hadoop-consumer/src/main/java/kafka/etl/KafkaETLRequest.java:88
↓ 1 callers
Method
getValue
()
core/src/main/scala/kafka/utils/Crc32.java:46
↓ 1 callers
Function
get_broker_shutdown_log_line
(systemTestEnv, testcaseEnv, leaderAttributesDict)
system_test/utils/kafka_system_test_utils.py:568
↓ 1 callers
Function
get_data_from_list_of_dicts
(listOfDicts, lookupKey, lookupVal, fieldToRetrieve)
system_test/utils/system_test_utils.py:131
↓ 1 callers
Function
get_jira
()
kafka-patch-review.py:19
↓ 1 callers
Function
get_leader_elected_log_line
(systemTestEnv, testcaseEnv, leaderAttributesDict)
system_test/utils/kafka_system_test_utils.py:626
↓ 1 callers
Function
get_mbeans_for_role
(dashboardsForRole)
system_test/utils/metrics.py:296
↓ 1 callers
Function
get_message_checksum
(logPathName)
system_test/utils/kafka_system_test_utils.py:1338
↓ 1 callers
Method
handleCompletedReceives
Handle any completed receives and update the response list with the responses received. @param responses The list of responses to update @param now Th
clients/src/main/java/org/apache/kafka/clients/NetworkClient.java:289
↓ 1 callers
Method
handleCompletedSends
Handle any completed request send. In particular if no response is expected consider the request complete. @param responses The list of responses to u
clients/src/main/java/org/apache/kafka/clients/NetworkClient.java:273
↓ 1 callers
Method
handleConnections
Record any newly completed connections
clients/src/main/java/org/apache/kafka/clients/NetworkClient.java:344
↓ 1 callers
Method
handleDisconnect
(ClientResponse response, long now)
clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java:209
↓ 1 callers
Method
handleDisconnections
Handle any disconnected connections @param responses The list of responses that completed with the disconnection @param now The current time
clients/src/main/java/org/apache/kafka/clients/NetworkClient.java:323
↓ 1 callers
Method
handleMetadataResponse
(RequestHeader header, Struct body, long now)
clients/src/main/java/org/apache/kafka/clients/NetworkClient.java:306
↓ 1 callers
Method
handleResponse
Handle a produce response
clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java:221
↓ 1 callers
Method
hasDefault
()
clients/src/main/java/org/apache/kafka/common/config/ConfigDef.java:333
↓ 1 callers
Method
hasField
Check if the struct contains a field. @param name @return Whether a field exists.
clients/src/main/java/org/apache/kafka/common/protocol/types/Struct.java:91
↓ 1 callers
Method
hasKey
Does the record have a key?
clients/src/main/java/org/apache/kafka/common/record/Record.java:249
↓ 1 callers
Method
hasReceive
()
clients/src/main/java/org/apache/kafka/common/network/Selector.java:398
↓ 1 callers
Method
hasRoomFor
Check if we have room for a new record containing the given key/value pair Note that the return value is based on the estimate of the bytes written t
clients/src/main/java/org/apache/kafka/common/record/MemoryRecords.java:103
↓ 1 callers
Method
hasSend
()
clients/src/main/java/org/apache/kafka/common/network/Selector.java:390
↓ 1 callers
Method
hasUnsent
@return Whether there is any unsent record in the accumulator.
clients/src/main/java/org/apache/kafka/clients/producer/internals/RecordAccumulator.java:249
↓ 1 callers
Method
hashCode
()
clients/src/main/java/org/apache/kafka/common/Node.java:52
↓ 1 callers
Method
hashCode
()
clients/src/main/java/org/apache/kafka/common/requests/AbstractRequestResponse.java:50
↓ 1 callers
Method
history
Get the list of sent records since the last call to {@link #clear()}
clients/src/main/java/org/apache/kafka/clients/producer/MockProducer.java:147
↓ 1 callers
Method
in
(List<String> validStrings)
clients/src/main/java/org/apache/kafka/common/config/ConfigDef.java:278
↓ 1 callers
Method
inSyncReplicas
The subset of the replicas that are in sync, that is caught-up to the leader and ready to take over as leader if the leader should fail
clients/src/main/java/org/apache/kafka/common/PartitionInfo.java:66
↓ 1 callers
Method
initCommonFields
(String groupId, Map<TopicPartition, PartitionData> offsetData, int versionId)
clients/src/main/java/org/apache/kafka/common/requests/OffsetCommitRequest.java:108
↓ 1 callers
Method
initiateClose
Start closing the sender (won't actually complete until all data is sent out)
clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java:203
↓ 1 callers
Method
innerDone
()
clients/src/main/java/org/apache/kafka/common/record/MemoryRecords.java:233
↓ 1 callers
Method
isBlackedOut
Return true if we are disconnected from the given node and can't re-establish a connection yet @param node The node to check @param now The current ti
clients/src/main/java/org/apache/kafka/clients/ClusterConnectionStates.java:51
↓ 1 callers
Method
isComplete
(long timeMs, MetricConfig config)
clients/src/main/java/org/apache/kafka/common/metrics/stats/SampledStat.java:125
↓ 1 callers
Method
isDone
()
clients/src/main/java/org/apache/kafka/clients/producer/internals/FutureRecordMetadata.java:70
↓ 1 callers
Method
isReady
(Node node, long now)
clients/src/test/java/org/apache/kafka/clients/MockClient.java:32
↓ 1 callers
Method
isReady
Check if the node with the given id is ready to send more requests. @param node The given node id @param now The current time in ms @return true if th
clients/src/main/java/org/apache/kafka/clients/NetworkClient.java:140
↓ 1 callers
Method
isValidOffset
()
contrib/hadoop-consumer/src/main/java/kafka/etl/KafkaETLRequest.java:93
↓ 1 callers
Method
keyForId
Get the selection key associated with this numeric id
clients/src/main/java/org/apache/kafka/common/network/Selector.java:357
↓ 1 callers
Method
lastSent
Get the last request we sent to the given node (but don't remove it from the queue) @param node The node id
clients/src/main/java/org/apache/kafka/clients/InFlightRequests.java:66
↓ 1 callers
Method
lastUpdate
The last time metadata was updated.
clients/src/main/java/org/apache/kafka/clients/producer/internals/Metadata.java:139
↓ 1 callers
Method
leaderFor
Get the current leader for the given topic-partition @param topicPartition The topic and partition we want to know the leader for @return The node tha
clients/src/main/java/org/apache/kafka/common/Cluster.java:106
↓ 1 callers
Method
leastLoadedNode
Choose the node with the fewest outstanding requests which is at least eligible for connection. This method will prefer a node with an existing connec
clients/src/main/java/org/apache/kafka/clients/NetworkClient.java:247
↓ 1 callers
Method
lessThan
(double upperBound)
clients/src/main/java/org/apache/kafka/common/metrics/Quota.java:32
↓ 1 callers
Method
logAll
()
clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java:101
↓ 1 callers
Function
main
main(), shut up, pylint
kafka-patch-review.py:32
↓ 1 callers
Function
main
()
system_test/system_test_runner.py:55
↓ 1 callers
Method
makeNext
()
clients/src/main/java/org/apache/kafka/common/utils/AbstractIterator.java:75
↓ 1 callers
Method
maybeComputeNext
()
clients/src/main/java/org/apache/kafka/common/utils/AbstractIterator.java:77
↓ 1 callers
Method
maybeRegisterNodeMetrics
(int node)
clients/src/main/java/org/apache/kafka/common/network/Selector.java:471
↓ 1 callers
Method
maybeRegisterTopicMetrics
(String topic)
clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java:394
↓ 1 callers
Method
maybeUpdateMetadata
Add a metadata request to the list of sends if we can make one
clients/src/main/java/org/apache/kafka/clients/NetworkClient.java:374
↓ 1 callers
Method
metadataRequest
Create a metadata request for the given topics
clients/src/main/java/org/apache/kafka/clients/NetworkClient.java:365
↓ 1 callers
Method
metricChange
This is called whenever a metric is updated or added @param metric
clients/src/main/java/org/apache/kafka/common/metrics/MetricsReporter.java:34
↓ 1 callers
Method
metrics
()
clients/src/main/java/org/apache/kafka/common/metrics/Sensor.java:170
↓ 1 callers
Method
metrics
Return a map of metrics maintained by the consumer
clients/src/main/java/org/apache/kafka/clients/consumer/Consumer.java:119
↓ 1 callers
Method
moreThan
(double lowerBound)
clients/src/main/java/org/apache/kafka/common/metrics/Quota.java:36
↓ 1 callers
Method
murmur2
Generates 32 bit murmur2 hash from byte array @param data byte array to hash @return 32 bit hash of the given array
clients/src/main/java/org/apache/kafka/common/utils/Utils.java:244
↓ 1 callers
Method
newWindow
()
clients/src/main/java/org/apache/kafka/clients/tools/ProducerPerformance.java:156
↓ 1 callers
Method
nextCompletion
(long start, int bytes, Stats stats)
clients/src/main/java/org/apache/kafka/clients/tools/ProducerPerformance.java:138
↓ 1 callers
Method
nextRequestHeader
Generate a request header for the given API key @param key The api key @return A request header with the appropriate client id and correlation id
clients/src/main/java/org/apache/kafka/clients/NetworkClient.java:219
↓ 1 callers
Method
numFields
The number of fields in this schema
clients/src/main/java/org/apache/kafka/common/protocol/types/Schema.java:88
↓ 1 callers
Method
offset
The position of this record in the corresponding Kafka partition. @throws Exception The exception thrown while fetching this record.
clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerRecord.java:118
↓ 1 callers
Method
oldest
(long now)
clients/src/main/java/org/apache/kafka/common/metrics/stats/SampledStat.java:80
↓ 1 callers
Method
originals
()
clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java:95
↓ 1 callers
Method
output
(String fileprefix)
contrib/hadoop-consumer/src/main/java/kafka/etl/KafkaETLContext.java:170
↓ 1 callers
Method
parse
(ByteBuffer buffer)
clients/src/main/java/org/apache/kafka/common/requests/ResponseHeader.java:52
↓ 1 callers
Method
parseAcks
(String acksString)
clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java:234
↓ 1 callers
Method
partition
The partition id
clients/src/main/java/org/apache/kafka/common/PartitionInfo.java:44
↓ 1 callers
Method
partition
Compute the partition for the given record. @param record The record being sent @param cluster The current cluster metadata
clients/src/main/java/org/apache/kafka/clients/producer/internals/Partitioner.java:46
↓ 1 callers
Method
partitionsForNode
Get the list of partitions whose leader is this node @param nodeId The node id @return A list of partitions
clients/src/main/java/org/apache/kafka/common/Cluster.java:137
↓ 1 callers
Method
percentile
()
clients/src/main/java/org/apache/kafka/common/metrics/stats/Percentile.java:36
↓ 1 callers
Method
percentiles
(int[] latencies, int count, double... percentiles)
clients/src/main/java/org/apache/kafka/clients/tools/ProducerPerformance.java:181
↓ 1 callers
Function
plot_graphs
(inputCsvFiles, labels, title, xLabel, yLabel, attribute, outputGraphFile)
system_test/utils/metrics.py:108
← previous
next →
401–500 of 1,332, ranked by callers