Fetch
This page encodes the smallest legal instance of the request and the response: numeric fields are zero, strings and byte arrays are empty, every array carries exactly one sample element, and any records field holds one empty 61-byte RecordBatch v2. Version 17 is a flexible version, so every struct is terminated by a uvarint tagged-field count and strings and arrays carry compact length-plus-one prefixes. Sizes below include the leading int32 size prefix.
- API key
- 1
- Encoded at
- v17
- Flexible versions
- 12+
- Headers
- req v2, resp v1
- Request versions
- 4-17
- Response versions
- 4-17
- Request size
- 113 bytes
- Response size
- 154 bytes
framerpc headerrequest bodyRecordBatchRecordresponse bodytagged_fields
Request
FetchRequest v17, request header v2, 113 bytes on the wire
byte layout (113 bytes, 16 bytes per row)
0
1
2
3
4
5
6
7
8
9
A
B
C
D
E
F
0000
0010
0020
0030
0040
0050
0060
0070
object tree
FetchRequest message v17 [0x0000, 113B] +-- Frame [0x0000, 4B] length-delimited framing | +-- size int32 = 109 [0x0000, 4B] number of bytes that follow, patched after encoding +-- RequestHeader v2 [0x0004, 11B] common request header | +-- request_api_key int16 = 1 (Fetch) [0x0004, 2B] numeric id of the API being invoked | +-- request_api_version int16 = 17 [0x0006, 2B] version of the API being invoked | +-- correlation_id int32 = 0 [0x0008, 4B] echoed back by the broker in the response | +-- client_id nullable_string = "" (int16 len=0) [0x000c, 2B] always a non-flexible int16-prefixed string | +-- tagged_fields uvarint = 0 [0x000e, 1B] number of tagged fields in the header +-- FetchRequest struct [0x000f, 98B] message body, version 17 +-- MaxWaitMs int32 = 0 [0x000f, 4B] The maximum time in milliseconds to wait for the response. +-- MinBytes int32 = 0 [0x0013, 4B] The minimum bytes to accumulate in the response. +-- MaxBytes int32 = 0 [0x0017, 4B] The maximum bytes to fetch. See KIP-74 for cases where this limit may not b... +-- IsolationLevel int8 = 0 [0x001b, 1B] This setting controls the visibility of transactional records. Using READ_U... +-- SessionId int32 = 0 [0x001c, 4B] The fetch session ID. +-- SessionEpoch int32 = 0 [0x0020, 4B] The fetch session epoch, which is used for ordering requests in a session. +-- Topics []FetchTopic = 1 element [0x0024, 52B] The topics to fetch. | +-- length uvarint = 2 (compact, n+1) [0x0024, 1B] one sample element follows | +-- FetchTopic[0] FetchTopic = struct [0x0025, 51B] | +-- TopicId uuid = 00000000-0000-0000-0000-000000000000 [0x0025, 16B] The unique topic ID. | +-- Partitions []FetchPartition = 1 element [0x0035, 34B] The partitions to fetch. | | +-- length uvarint = 2 (compact, n+1) [0x0035, 1B] one sample element follows | | +-- FetchPartition[0] FetchPartition = struct [0x0036, 33B] | | +-- Partition int32 = 0 [0x0036, 4B] The partition index. | | +-- CurrentLeaderEpoch int32 = 0 [0x003a, 4B] The current leader epoch of the partition. | | +-- FetchOffset int64 = 0 [0x003e, 8B] The message offset. | | +-- LastFetchedEpoch int32 = 0 [0x0046, 4B] The epoch of the last fetched record or -1 if there is none. | | +-- LogStartOffset int64 = 0 [0x004a, 8B] The earliest available offset of the follower replica. The field is only us... | | +-- PartitionMaxBytes int32 = 0 [0x0052, 4B] The maximum bytes to fetch from this partition. See KIP-74 for cases where ... | | +-- tagged_fields uvarint = 0 [0x0056, 1B] number of tagged fields in this struct | +-- tagged_fields uvarint = 0 [0x0057, 1B] number of tagged fields in this struct +-- ForgottenTopicsData []ForgottenTopic = 1 element [0x0058, 23B] In an incremental fetch request, the partitions to remove. | +-- length uvarint = 2 (compact, n+1) [0x0058, 1B] one sample element follows | +-- ForgottenTopic[0] ForgottenTopic = struct [0x0059, 22B] | +-- TopicId uuid = 00000000-0000-0000-0000-000000000000 [0x0059, 16B] The unique topic ID. | +-- Partitions []int32 = 1 element [0x0069, 5B] The partitions indexes to forget. | | +-- length uvarint = 2 (compact, n+1) [0x0069, 1B] one sample element follows | | +-- int32[0] int32 = 0 [0x006a, 4B] | +-- tagged_fields uvarint = 0 [0x006e, 1B] number of tagged fields in this struct +-- RackId string = "" (compact, len+1=1) [0x006f, 1B] Rack ID of the consumer making this request. +-- tagged_fields uvarint = 0 [0x0070, 1B] number of tagged fields in this struct
kafka message schema (.json)
{ "apiKey": 1, "type": "request", "listeners": ["broker", "controller"], "name": "FetchRequest", // Versions 0-3 were removed in Apache Kafka 4.0, Version 4 is the new baseline. // // Version 1 is the same as version 0. // Starting in Version 2, the requester must be able to handle Kafka Log // Message format version 1. // Version 3 adds MaxBytes. Starting in version 3, the partition ordering in // the request is now relevant. Partitions will be processed in the order // they appear in the request. // // Version 4 adds IsolationLevel. Starting in version 4, the reqestor must be // able to handle Kafka log message format version 2. // // Version 5 adds LogStartOffset to indicate the earliest available offset of // partition data that can be consumed. // // Version 6 is the same as version 5. // // Version 7 adds incremental fetch request support. // // Version 8 is the same as version 7. // // Version 9 adds CurrentLeaderEpoch, as described in KIP-320. // // Version 10 indicates that we can use the ZStd compression algorithm, as // described in KIP-110. // Version 12 adds flexible versions support as well as epoch validation through // the `LastFetchedEpoch` field // // Version 13 replaces topic names with topic IDs (KIP-516). May return UNKNOWN_TOPIC_ID error code. // // Version 14 is the same as version 13 but it also receives a new error called OffsetMovedToTieredStorageException(KIP-405) // // Version 15 adds the ReplicaState which includes new field ReplicaEpoch and the ReplicaId. Also, // deprecate the old ReplicaId field and set its default value to -1. (KIP-903) // // Version 16 is the same as version 15 (KIP-951). // // Version 17 adds directory id support from KIP-853 "validVersions": "4-17", "flexibleVersions": "12+", "fields": [ { "name": "ClusterId", "type": "string", "versions": "12+", "nullableVersions": "12+", "default": "null", "taggedVersions": "12+", "tag": 0, "ignorable": true, "about": "The clusterId if known. This is used to validate metadata fetches prior to broker registration." }, { "name": "ReplicaId", "type": "int32", "versions": "0-14", "default": "-1", "entityType": "brokerId", "about": "The broker ID of the follower, of -1 if this request is from a consumer." }, { "name": "ReplicaState", "type": "ReplicaState", "versions": "15+", "taggedVersions": "15+", "tag": 1, "about": "The state of the replica in the follower.", "fields": [ { "name": "ReplicaId", "type": "int32", "versions": "15+", "default": "-1", "entityType": "brokerId", "about": "The replica ID of the follower, or -1 if this request is from a consumer." }, { "name": "ReplicaEpoch", "type": "int64", "versions": "15+", "default": "-1", "about": "The epoch of this follower, or -1 if not available." } ]}, { "name": "MaxWaitMs", "type": "int32", "versions": "0+", "about": "The maximum time in milliseconds to wait for the response." }, { "name": "MinBytes", "type": "int32", "versions": "0+", "about": "The minimum bytes to accumulate in the response." }, { "name": "MaxBytes", "type": "int32", "versions": "3+", "default": "0x7fffffff", "ignorable": true, "about": "The maximum bytes to fetch. See KIP-74 for cases where this limit may not be honored." }, { "name": "IsolationLevel", "type": "int8", "versions": "4+", "default": "0", "ignorable": true, "about": "This setting controls the visibility of transactional records. Using READ_UNCOMMITTED (isolation_level = 0) makes all records visible. With READ_COMMITTED (isolation_level = 1), non-transactional and COMMITTED transactional records are visible. To be more concrete, READ_COMMITTED returns all data from offsets smaller than the current LSO (last stable offset), and enables the inclusion of the list of aborted transactions in the result, which allows consumers to discard ABORTED transactional records." }, { "name": "SessionId", "type": "int32", "versions": "7+", "default": "0", "ignorable": true, "about": "The fetch session ID." }, { "name": "SessionEpoch", "type": "int32", "versions": "7+", "default": "-1", "ignorable": true, "about": "The fetch session epoch, which is used for ordering requests in a session." }, { "name": "Topics", "type": "[]FetchTopic", "versions": "0+", "about": "The topics to fetch.", "fields": [ { "name": "Topic", "type": "string", "versions": "0-12", "entityType": "topicName", "ignorable": true, "about": "The name of the topic to fetch." }, { "name": "TopicId", "type": "uuid", "versions": "13+", "ignorable": true, "about": "The unique topic ID."}, { "name": "Partitions", "type": "[]FetchPartition", "versions": "0+", "about": "The partitions to fetch.", "fields": [ { "name": "Partition", "type": "int32", "versions": "0+", "about": "The partition index." }, { "name": "CurrentLeaderEpoch", "type": "int32", "versions": "9+", "default": "-1", "ignorable": true, "about": "The current leader epoch of the partition." }, { "name": "FetchOffset", "type": "int64", "versions": "0+", "about": "The message offset." }, { "name": "LastFetchedEpoch", "type": "int32", "versions": "12+", "default": "-1", "ignorable": false, "about": "The epoch of the last fetched record or -1 if there is none."}, { "name": "LogStartOffset", "type": "int64", "versions": "5+", "default": "-1", "ignorable": true, "about": "The earliest available offset of the follower replica. The field is only used when the request is sent by the follower."}, { "name": "PartitionMaxBytes", "type": "int32", "versions": "0+", "about": "The maximum bytes to fetch from this partition. See KIP-74 for cases where this limit may not be honored." }, { "name": "ReplicaDirectoryId", "type": "uuid", "versions": "17+", "taggedVersions": "17+", "tag": 0, "ignorable": true, "about": "The directory id of the follower fetching." } ]} ]}, { "name": "ForgottenTopicsData", "type": "[]ForgottenTopic", "versions": "7+", "ignorable": false, "about": "In an incremental fetch request, the partitions to remove.", "fields": [ { "name": "Topic", "type": "string", "versions": "7-12", "entityType": "topicName", "ignorable": true, "about": "The topic name." }, { "name": "TopicId", "type": "uuid", "versions": "13+", "ignorable": true, "about": "The unique topic ID."}, { "name": "Partitions", "type": "[]int32", "versions": "7+", "about": "The partitions indexes to forget." } ]}, { "name": "RackId", "type": "string", "versions": "11+", "default": "", "ignorable": true, "about": "Rack ID of the consumer making this request."} ] }
Response
FetchResponse v17, response header v1, 154 bytes on the wire
byte layout (154 bytes, 16 bytes per row)
0
1
2
3
4
5
6
7
8
9
A
B
C
D
E
F
0000
0010
0020
0030
0040
0050
0060
0070
0080
0090
object tree
FetchResponse message v17 [0x0000, 154B] +-- Frame [0x0000, 4B] length-delimited framing | +-- size int32 = 150 [0x0000, 4B] number of bytes that follow, patched after encoding +-- ResponseHeader v1 [0x0004, 5B] common response header | +-- correlation_id int32 = 0 [0x0004, 4B] matches the correlation_id of the request | +-- tagged_fields uvarint = 0 [0x0008, 1B] number of tagged fields in the header +-- FetchResponse struct [0x0009, 145B] message body, version 17 +-- ThrottleTimeMs int32 = 0 [0x0009, 4B] The duration in milliseconds for which the request was throttled due to a q... +-- ErrorCode int16 = 0 [0x000d, 2B] The top level response error code. +-- SessionId int32 = 0 [0x000f, 4B] The fetch session ID, or 0 if this is not part of a fetch session. +-- Responses []FetchableTopicResponse = 1 element [0x0013, 134B] The response topics. | +-- length uvarint = 2 (compact, n+1) [0x0013, 1B] one sample element follows | +-- FetchableTopicResponse[0] FetchableTopicResponse = struct [0x0014, 133B] | +-- TopicId uuid = 00000000-0000-0000-0000-000000000000 [0x0014, 16B] The unique topic ID. | +-- Partitions []PartitionData = 1 element [0x0024, 116B] The topic partitions. | | +-- length uvarint = 2 (compact, n+1) [0x0024, 1B] one sample element follows | | +-- PartitionData[0] PartitionData = struct [0x0025, 115B] | | +-- PartitionIndex int32 = 0 [0x0025, 4B] The partition index. | | +-- ErrorCode int16 = 0 [0x0029, 2B] The error code, or 0 if there was no fetch error. | | +-- HighWatermark int64 = 0 [0x002b, 8B] The current high water mark. | | +-- LastStableOffset int64 = 0 [0x0033, 8B] The last stable offset (or LSO) of the partition. This is the last offset s... | | +-- LogStartOffset int64 = 0 [0x003b, 8B] The current log start offset. | | +-- AbortedTransactions []AbortedTransaction = 1 element [0x0043, 18B] The aborted transactions. | | | +-- length uvarint = 2 (compact, n+1) [0x0043, 1B] one sample element follows | | | +-- AbortedTransaction[0] AbortedTransaction = struct [0x0044, 17B] | | | +-- ProducerId int64 = 0 [0x0044, 8B] The producer id associated with the aborted transaction. | | | +-- FirstOffset int64 = 0 [0x004c, 8B] The first offset in the aborted transaction. | | | +-- tagged_fields uvarint = 0 [0x0054, 1B] number of tagged fields in this struct | | +-- PreferredReadReplica int32 = 0 [0x0055, 4B] The preferred read replica for the consumer to use on its next fetch request. | | +-- Records records = 1 RecordBatch [0x0059, 62B] The record data. | | | +-- length uvarint = 62 (compact, n+1) [0x0059, 1B] size of the record set in bytes | | | +-- RecordBatch v2 = empty [0x005a, 61B] fixed 61-byte RecordBatch v2 header, zero records | | | +-- baseOffset int64 = 0 [0x005a, 8B] offset of the first record in the batch | | | +-- batchLength int32 = 49 [0x0062, 4B] bytes after this field to the end of the batch | | | +-- partitionLeaderEpoch int32 = -1 [0x0066, 4B] leader epoch, -1 when produced by a client | | | +-- magic int8 = 2 [0x006a, 1B] record format version, 2 | | | +-- crc uint32 = crc32c of the bytes that follow [0x006b, 4B] CRC-32C over attributes .. end of batch | | | +-- attributes int16 = 0 [0x006f, 2B] compression, timestamp type, txn and control flags | | | +-- lastOffsetDelta int32 = -1 [0x0071, 4B] offset delta of the last record, -1 when empty | | | +-- baseTimestamp int64 = -1 [0x0075, 8B] timestamp of the first record | | | +-- maxTimestamp int64 = -1 [0x007d, 8B] largest timestamp in the batch | | | +-- producerId int64 = -1 [0x0085, 8B] producer id, -1 when non-idempotent | | | +-- producerEpoch int16 = -1 [0x008d, 2B] producer epoch, -1 when non-idempotent | | | +-- baseSequence int32 = -1 [0x008f, 4B] sequence of the first record, -1 when non-idempotent | | | +-- recordsCount int32 = 0 [0x0093, 4B] number of records that follow, 0 here | | +-- tagged_fields uvarint = 0 [0x0097, 1B] number of tagged fields in this struct | +-- tagged_fields uvarint = 0 [0x0098, 1B] number of tagged fields in this struct +-- tagged_fields uvarint = 0 [0x0099, 1B] number of tagged fields in this struct
kafka message schema (.json)
{ "apiKey": 1, "type": "response", "name": "FetchResponse", // Versions 0-3 were removed in Apache Kafka 4.0, Version 4 is the new baseline. // // Version 1 adds throttle time. Version 2 and 3 are the same as version 1. // // Version 4 adds features for transactional consumption. // // Version 5 adds LogStartOffset to indicate the earliest available offset of // partition data that can be consumed. // // Starting in version 6, we may return KAFKA_STORAGE_ERROR as an error code. // // Version 7 adds incremental fetch request support. // // Starting in version 8, on quota violation, brokers send out responses before throttling. // // Version 9 is the same as version 8. // // Version 10 indicates that the response data can use the ZStd compression // algorithm, as described in KIP-110. // Version 12 adds support for flexible versions, epoch detection through the `TruncationOffset` field, // and leader discovery through the `CurrentLeader` field // // Version 13 replaces the topic name field with topic ID (KIP-516). // // Version 14 is the same as version 13 but it also receives a new error called OffsetMovedToTieredStorageException (KIP-405) // // Version 15 is the same as version 14 (KIP-903). // // Version 16 adds the 'NodeEndpoints' field (KIP-951). // // Version 17 no changes to the response (KIP-853). "validVersions": "4-17", "flexibleVersions": "12+", "fields": [ { "name": "ThrottleTimeMs", "type": "int32", "versions": "1+", "ignorable": true, "about": "The duration in milliseconds for which the request was throttled due to a quota violation, or zero if the request did not violate any quota." }, { "name": "ErrorCode", "type": "int16", "versions": "7+", "ignorable": true, "about": "The top level response error code." }, { "name": "SessionId", "type": "int32", "versions": "7+", "default": "0", "ignorable": false, "about": "The fetch session ID, or 0 if this is not part of a fetch session." }, { "name": "Responses", "type": "[]FetchableTopicResponse", "versions": "0+", "about": "The response topics.", "fields": [ { "name": "Topic", "type": "string", "versions": "0-12", "ignorable": true, "entityType": "topicName", "about": "The topic name." }, { "name": "TopicId", "type": "uuid", "versions": "13+", "ignorable": true, "about": "The unique topic ID."}, { "name": "Partitions", "type": "[]PartitionData", "versions": "0+", "about": "The topic partitions.", "fields": [ { "name": "PartitionIndex", "type": "int32", "versions": "0+", "about": "The partition index." }, { "name": "ErrorCode", "type": "int16", "versions": "0+", "about": "The error code, or 0 if there was no fetch error." }, { "name": "HighWatermark", "type": "int64", "versions": "0+", "about": "The current high water mark." }, { "name": "LastStableOffset", "type": "int64", "versions": "4+", "default": "-1", "ignorable": true, "about": "The last stable offset (or LSO) of the partition. This is the last offset such that the state of all transactional records prior to this offset have been decided (ABORTED or COMMITTED)." }, { "name": "LogStartOffset", "type": "int64", "versions": "5+", "default": "-1", "ignorable": true, "about": "The current log start offset." }, { "name": "DivergingEpoch", "type": "EpochEndOffset", "versions": "12+", "taggedVersions": "12+", "tag": 0, "about": "In case divergence is detected based on the `LastFetchedEpoch` and `FetchOffset` in the request, this field indicates the largest epoch and its end offset such that subsequent records are known to diverge.", "fields": [ { "name": "Epoch", "type": "int32", "versions": "12+", "default": "-1", "about": "The largest epoch." }, { "name": "EndOffset", "type": "int64", "versions": "12+", "default": "-1", "about": "The end offset of the epoch." } ]}, { "name": "CurrentLeader", "type": "LeaderIdAndEpoch", "versions": "12+", "taggedVersions": "12+", "tag": 1, "about": "The current leader of the partition.", "fields": [ { "name": "LeaderId", "type": "int32", "versions": "12+", "default": "-1", "entityType": "brokerId", "about": "The ID of the current leader or -1 if the leader is unknown."}, { "name": "LeaderEpoch", "type": "int32", "versions": "12+", "default": "-1", "about": "The latest known leader epoch." } ]}, { "name": "SnapshotId", "type": "SnapshotId", "versions": "12+", "taggedVersions": "12+", "tag": 2, "about": "In the case of fetching an offset less than the LogStartOffset, this is the end offset and epoch that should be used in the FetchSnapshot request.", "fields": [ { "name": "EndOffset", "type": "int64", "versions": "0+", "default": "-1", "about": "The end offset of the epoch." }, { "name": "Epoch", "type": "int32", "versions": "0+", "default": "-1", "about": "The largest epoch." } ]}, { "name": "AbortedTransactions", "type": "[]AbortedTransaction", "versions": "4+", "nullableVersions": "4+", "ignorable": true, "about": "The aborted transactions.", "fields": [ { "name": "ProducerId", "type": "int64", "versions": "4+", "entityType": "producerId", "about": "The producer id associated with the aborted transaction." }, { "name": "FirstOffset", "type": "int64", "versions": "4+", "about": "The first offset in the aborted transaction." } ]}, { "name": "PreferredReadReplica", "type": "int32", "versions": "11+", "default": "-1", "ignorable": false, "entityType": "brokerId", "about": "The preferred read replica for the consumer to use on its next fetch request."}, { "name": "Records", "type": "records", "versions": "0+", "nullableVersions": "0+", "about": "The record data."} ]} ]}, { "name": "NodeEndpoints", "type": "[]NodeEndpoint", "versions": "16+", "taggedVersions": "16+", "tag": 0, "about": "Endpoints for all current-leaders enumerated in PartitionData, with errors NOT_LEADER_OR_FOLLOWER & FENCED_LEADER_EPOCH.", "fields": [ { "name": "NodeId", "type": "int32", "versions": "16+", "mapKey": true, "entityType": "brokerId", "about": "The ID of the associated node."}, { "name": "Host", "type": "string", "versions": "16+", "about": "The node's hostname." }, { "name": "Port", "type": "int32", "versions": "16+", "about": "The node's port." }, { "name": "Rack", "type": "string", "versions": "16+", "nullableVersions": "16+", "default": "null", "about": "The rack of the node, or null if it has not been assigned to a rack." } ]} ] }