Produce
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 12 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
- 0
- Encoded at
- v12
- Flexible versions
- 9+
- Headers
- req v2, resp v1
- Request versions
- 3-12
- Response versions
- 3-12
- Request size
- 94 bytes
- Response size
- 57 bytes
framerpc headerrequest bodyRecordBatchRecordresponse bodytagged_fields
Request
ProduceRequest v12, request header v2, 94 bytes on the wire
byte layout (94 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
object tree
ProduceRequest message v12 [0x0000, 94B] +-- Frame [0x0000, 4B] length-delimited framing | +-- size int32 = 90 [0x0000, 4B] number of bytes that follow, patched after encoding +-- RequestHeader v2 [0x0004, 11B] common request header | +-- request_api_key int16 = 0 (Produce) [0x0004, 2B] numeric id of the API being invoked | +-- request_api_version int16 = 12 [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 +-- ProduceRequest struct [0x000f, 79B] message body, version 12 +-- TransactionalId string = "" (compact, len+1=1) [0x000f, 1B] The transactional ID, or null if the producer is not transactional. +-- Acks int16 = 0 [0x0010, 2B] The number of acknowledgments the producer requires the leader to have rece... +-- TimeoutMs int32 = 0 [0x0012, 4B] The timeout to await a response in milliseconds. +-- TopicData []TopicProduceData = 1 element [0x0016, 71B] Each topic to produce to. | +-- length uvarint = 2 (compact, n+1) [0x0016, 1B] one sample element follows | +-- TopicProduceData[0] TopicProduceData = struct [0x0017, 70B] | +-- Name string = "" (compact, len+1=1) [0x0017, 1B] The topic name. | +-- PartitionData []PartitionProduceData = 1 element [0x0018, 68B] Each partition to produce to. | | +-- length uvarint = 2 (compact, n+1) [0x0018, 1B] one sample element follows | | +-- PartitionProduceData[0] PartitionProduceData = struct [0x0019, 67B] | | +-- Index int32 = 0 [0x0019, 4B] The partition index. | | +-- Records records = 1 RecordBatch [0x001d, 62B] The record data to be produced. | | | +-- length uvarint = 62 (compact, n+1) [0x001d, 1B] size of the record set in bytes | | | +-- RecordBatch v2 = empty [0x001e, 61B] fixed 61-byte RecordBatch v2 header, zero records | | | +-- baseOffset int64 = 0 [0x001e, 8B] offset of the first record in the batch | | | +-- batchLength int32 = 49 [0x0026, 4B] bytes after this field to the end of the batch | | | +-- partitionLeaderEpoch int32 = -1 [0x002a, 4B] leader epoch, -1 when produced by a client | | | +-- magic int8 = 2 [0x002e, 1B] record format version, 2 | | | +-- crc uint32 = crc32c of the bytes that follow [0x002f, 4B] CRC-32C over attributes .. end of batch | | | +-- attributes int16 = 0 [0x0033, 2B] compression, timestamp type, txn and control flags | | | +-- lastOffsetDelta int32 = -1 [0x0035, 4B] offset delta of the last record, -1 when empty | | | +-- baseTimestamp int64 = -1 [0x0039, 8B] timestamp of the first record | | | +-- maxTimestamp int64 = -1 [0x0041, 8B] largest timestamp in the batch | | | +-- producerId int64 = -1 [0x0049, 8B] producer id, -1 when non-idempotent | | | +-- producerEpoch int16 = -1 [0x0051, 2B] producer epoch, -1 when non-idempotent | | | +-- baseSequence int32 = -1 [0x0053, 4B] sequence of the first record, -1 when non-idempotent | | | +-- recordsCount int32 = 0 [0x0057, 4B] number of records that follow, 0 here | | +-- tagged_fields uvarint = 0 [0x005b, 1B] number of tagged fields in this struct | +-- tagged_fields uvarint = 0 [0x005c, 1B] number of tagged fields in this struct +-- tagged_fields uvarint = 0 [0x005d, 1B] number of tagged fields in this struct
kafka message schema (.json)
{ "apiKey": 0, "type": "request", "listeners": ["broker"], "name": "ProduceRequest", // Versions 0-2 were removed in Apache Kafka 4.0, version 3 is the new baseline. Due to a bug in librdkafka, // these versions have to be included in the api versions response (see KAFKA-18659), but are rejected otherwise. // See `ApiKeys.PRODUCE_API_VERSIONS_RESPONSE_MIN_VERSION` for more details. // // Version 1 and 2 are the same as version 0. // // Version 3 adds the transactional ID, which is used for authorization when attempting to write // transactional data. Version 3 also adds support for Kafka Message Format v2. // // Version 4 is the same as version 3, but the requester must be prepared to handle a // KAFKA_STORAGE_ERROR. // // Version 5 and 6 are the same as version 3. // // Starting in version 7, records can be produced using ZStandard compression. See KIP-110. // // Starting in Version 8, response has RecordErrors and ErrorMessage. See KIP-467. // // Version 9 enables flexible versions. // // Version 10 is the same as version 9 (KIP-951). // // Version 11 adds support for new error code TRANSACTION_ABORTABLE (KIP-890). // // Version 12 is the same as version 11 (KIP-890). Note when produce requests are used in transaction, if // transaction V2 (KIP_890 part 2) is enabled, the produce request will also include the function for a // AddPartitionsToTxn call. If V2 is disabled, the client can't use produce request version higher than 11 within // a transaction. "validVersions": "3-12", "flexibleVersions": "9+", "fields": [ { "name": "TransactionalId", "type": "string", "versions": "3+", "nullableVersions": "3+", "default": "null", "entityType": "transactionalId", "about": "The transactional ID, or null if the producer is not transactional." }, { "name": "Acks", "type": "int16", "versions": "0+", "about": "The number of acknowledgments the producer requires the leader to have received before considering a request complete. Allowed values: 0 for no acknowledgments, 1 for only the leader and -1 for the full ISR." }, { "name": "TimeoutMs", "type": "int32", "versions": "0+", "about": "The timeout to await a response in milliseconds." }, { "name": "TopicData", "type": "[]TopicProduceData", "versions": "0+", "about": "Each topic to produce to.", "fields": [ { "name": "Name", "type": "string", "versions": "0+", "entityType": "topicName", "mapKey": true, "about": "The topic name." }, { "name": "PartitionData", "type": "[]PartitionProduceData", "versions": "0+", "about": "Each partition to produce to.", "fields": [ { "name": "Index", "type": "int32", "versions": "0+", "about": "The partition index." }, { "name": "Records", "type": "records", "versions": "0+", "nullableVersions": "0+", "about": "The record data to be produced." } ]} ]} ] }
Response
ProduceResponse v12, response header v1, 57 bytes on the wire
byte layout (57 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
object tree
ProduceResponse message v12 [0x0000, 57B] +-- Frame [0x0000, 4B] length-delimited framing | +-- size int32 = 53 [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 +-- ProduceResponse struct [0x0009, 48B] message body, version 12 +-- Responses []TopicProduceResponse = 1 element [0x0009, 43B] Each produce response. | +-- length uvarint = 2 (compact, n+1) [0x0009, 1B] one sample element follows | +-- TopicProduceResponse[0] TopicProduceResponse = struct [0x000a, 42B] | +-- Name string = "" (compact, len+1=1) [0x000a, 1B] The topic name. | +-- PartitionResponses []PartitionProduceResponse = 1 element [0x000b, 40B] Each partition that we produced to within the topic. | | +-- length uvarint = 2 (compact, n+1) [0x000b, 1B] one sample element follows | | +-- PartitionProduceResponse[0] PartitionProduceResponse = struct [0x000c, 39B] | | +-- Index int32 = 0 [0x000c, 4B] The partition index. | | +-- ErrorCode int16 = 0 [0x0010, 2B] The error code, or 0 if there was no error. | | +-- BaseOffset int64 = 0 [0x0012, 8B] The base offset. | | +-- LogAppendTimeMs int64 = 0 [0x001a, 8B] The timestamp returned by broker after appending the messages. If CreateTim... | | +-- LogStartOffset int64 = 0 [0x0022, 8B] The log start offset. | | +-- RecordErrors []BatchIndexAndErrorMessage = 1 element [0x002a, 7B] The batch indices of records that caused the batch to be dropped. | | | +-- length uvarint = 2 (compact, n+1) [0x002a, 1B] one sample element follows | | | +-- BatchIndexAndErrorMessage[0] BatchIndexAndErrorMessage = struct [0x002b, 6B] | | | +-- BatchIndex int32 = 0 [0x002b, 4B] The batch index of the record that caused the batch to be dropped. | | | +-- BatchIndexErrorMessage string = "" (compact, len+1=1) [0x002f, 1B] The error message of the record that caused the batch to be dropped. | | | +-- tagged_fields uvarint = 0 [0x0030, 1B] number of tagged fields in this struct | | +-- ErrorMessage string = "" (compact, len+1=1) [0x0031, 1B] The global error message summarizing the common root cause of the records t... | | +-- tagged_fields uvarint = 0 [0x0032, 1B] number of tagged fields in this struct | +-- tagged_fields uvarint = 0 [0x0033, 1B] number of tagged fields in this struct +-- ThrottleTimeMs int32 = 0 [0x0034, 4B] The duration in milliseconds for which the request was throttled due to a q... +-- tagged_fields uvarint = 0 [0x0038, 1B] number of tagged fields in this struct
kafka message schema (.json)
{ "apiKey": 0, "type": "response", "name": "ProduceResponse", // Versions 0-2 were removed in Apache Kafka 4.0, version 3 is the new baseline. Due to a bug in librdkafka, // these versions have to be included in the api versions response (see KAFKA-18659), but are rejected otherwise. // See `ApiKeys.PRODUCE_API_VERSIONS_RESPONSE_MIN_VERSION` for more details. // // Version 1 added the throttle time. // Version 2 added the log append time. // // Version 3 is the same as version 2. // // Version 4 added KAFKA_STORAGE_ERROR as a possible error code. // // Version 5 added LogStartOffset to filter out spurious OutOfOrderSequenceExceptions on the client. // // Version 8 added RecordErrors and ErrorMessage to include information about // records that cause the whole batch to be dropped. See KIP-467 for details. // // Version 9 enables flexible versions. // // Version 10 adds 'CurrentLeader' and 'NodeEndpoints' as tagged fields (KIP-951) // // Version 11 adds support for new error code TRANSACTION_ABORTABLE (KIP-890). // // Version 12 is the same as version 10 (KIP-890). "validVersions": "3-12", "flexibleVersions": "9+", "fields": [ { "name": "Responses", "type": "[]TopicProduceResponse", "versions": "0+", "about": "Each produce response.", "fields": [ { "name": "Name", "type": "string", "versions": "0+", "entityType": "topicName", "mapKey": true, "about": "The topic name." }, { "name": "PartitionResponses", "type": "[]PartitionProduceResponse", "versions": "0+", "about": "Each partition that we produced to within the topic.", "fields": [ { "name": "Index", "type": "int32", "versions": "0+", "about": "The partition index." }, { "name": "ErrorCode", "type": "int16", "versions": "0+", "about": "The error code, or 0 if there was no error." }, { "name": "BaseOffset", "type": "int64", "versions": "0+", "about": "The base offset." }, { "name": "LogAppendTimeMs", "type": "int64", "versions": "2+", "default": "-1", "ignorable": true, "about": "The timestamp returned by broker after appending the messages. If CreateTime is used for the topic, the timestamp will be -1. If LogAppendTime is used for the topic, the timestamp will be the broker local time when the messages are appended." }, { "name": "LogStartOffset", "type": "int64", "versions": "5+", "default": "-1", "ignorable": true, "about": "The log start offset." }, { "name": "RecordErrors", "type": "[]BatchIndexAndErrorMessage", "versions": "8+", "ignorable": true, "about": "The batch indices of records that caused the batch to be dropped.", "fields": [ { "name": "BatchIndex", "type": "int32", "versions": "8+", "about": "The batch index of the record that caused the batch to be dropped." }, { "name": "BatchIndexErrorMessage", "type": "string", "default": "null", "versions": "8+", "nullableVersions": "8+", "about": "The error message of the record that caused the batch to be dropped."} ]}, { "name": "ErrorMessage", "type": "string", "default": "null", "versions": "8+", "nullableVersions": "8+", "ignorable": true, "about": "The global error message summarizing the common root cause of the records that caused the batch to be dropped."}, { "name": "CurrentLeader", "type": "LeaderIdAndEpoch", "versions": "10+", "taggedVersions": "10+", "tag": 0, "about": "The leader broker that the producer should use for future requests.", "fields": [ { "name": "LeaderId", "type": "int32", "versions": "10+", "default": "-1", "entityType": "brokerId", "about": "The ID of the current leader or -1 if the leader is unknown."}, { "name": "LeaderEpoch", "type": "int32", "versions": "10+", "default": "-1", "about": "The latest known leader epoch."} ]} ]} ]}, { "name": "ThrottleTimeMs", "type": "int32", "versions": "1+", "ignorable": true, "default": "0", "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": "NodeEndpoints", "type": "[]NodeEndpoint", "versions": "10+", "taggedVersions": "10+", "tag": 0, "about": "Endpoints for all current-leaders enumerated in PartitionProduceResponses, with errors NOT_LEADER_OR_FOLLOWER.", "fields": [ { "name": "NodeId", "type": "int32", "versions": "10+", "mapKey": true, "entityType": "brokerId", "about": "The ID of the associated node."}, { "name": "Host", "type": "string", "versions": "10+", "about": "The node's hostname." }, { "name": "Port", "type": "int32", "versions": "10+", "about": "The node's port." }, { "name": "Rack", "type": "string", "versions": "10+", "nullableVersions": "10+", "default": "null", "about": "The rack of the node, or null if it has not been assigned to a rack." } ]} ] }