DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
...
There will be an addition in the
org.apache.kafka.clients.consumer.ShareAcquireModeenum defined in KIP-1206.Code Block public enum ShareAcquireMode { BATCH_OPTIMIZED("batch_optimized"), RECORD_LIMIT("record_limit"), RENEW_ACK_MODE("renew_ack_mode"); // Added as part for KIP-1222 public final String name; ShareAcquireMode(final String name) { this.name = name; } /** * Case-insensitive acquire mode lookup by string name. */ public static ShareAcquireMode of(final String name) { return ShareAcquireMode.valueOf(name.toUpperCase(Locale.ROOT)); } }- The acknowledgement type will be used in the existing ShareFetch and ShareAcknowledge RPCs. Hence, we must add new versions for them. Additionally, we will add a new value for the field shareAcquireMode called
renew-ack-modewith value id value 255. This field will help in optimising the broker side implementation for processing renew acknowledgements.clients/src/main/resources/common/message/ShareFetchRequest.jsonCode Block // ShareFetch { "apiKey": 78, "type": "request", "listeners": ["broker"], "name": "ShareFetchRequest", // Version 0 was used for early access of KIP-932 in Apache Kafka 4.0 but removed in Apacke Kafka 4.1. // // Version 1 is the initial stable version (KIP-932). // // Version 2 supports RENEW ack type "validVersions": "1-2", "flexibleVersions": "0+", "fields": [ { "name": "GroupId", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null", "entityType": "groupId", "about": "The group identifier." }, { "name": "MemberId", "type": "string", "versions": "0+", "nullableVersions": "0+", "about": "The member ID." }, { "name": "ShareSessionEpoch", "type": "int32", "versions": "0+", "about": "The current share session epoch: 0 to open a share session; -1 to close it; otherwise increments for consecutive requests." }, { "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": "0+", "default": "0x7fffffff", "about": "The maximum bytes to fetch. See KIP-74 for cases where this limit may not be honored." }, { "name": "MaxRecords", "type": "int32", "versions": "1+", "about": "The maximum number of records to fetch. This limit can be exceeded for alignment of batch boundaries." }, { "name": "BatchSize", "type": "int32", "versions": "1+", "about": "The optimal number of records for batches of acquired records and acknowledgements." }, { "name": "ShareAcquireMode", "type": "int8", "versions": "2+", "about": "The acquire mode to control the fetch behavior: 0 - batch-optimized, 1 - record-limit, 255 - renew-ack-mode" }, // From KIP-1206 { "name": "Topics", "type": "[]FetchTopic", "versions": "0+", "about": "The topics to fetch.", "fields": [ { "name": "TopicId", "type": "uuid", "versions": "0+", "about": "The unique topic ID.", "mapKey": true }, { "name": "Partitions", "type": "[]FetchPartition", "versions": "0+", "about": "The partitions to fetch.", "fields": [ { "name": "PartitionIndex", "type": "int32", "versions": "0+", "mapKey": true, "about": "The partition index." }, { "name": "PartitionMaxBytes", "type": "int32", "versions": "0", "about": "The maximum bytes to fetch from this partition. 0 when only acknowledgement with no fetching is required. See KIP-74 for cases where this limit may not be honored." }, { "name": "AcknowledgementBatches", "type": "[]AcknowledgementBatch", "versions": "0+", "about": "Record batches to acknowledge.", "fields": [ { "name": "FirstOffset", "type": "int64", "versions": "0+", "about": "First offset of batch of records to acknowledge."}, { "name": "LastOffset", "type": "int64", "versions": "0+", "about": "Last offset (inclusive) of batch of records to acknowledge."}, { "name": "AcknowledgeTypes", "type": "[]int8", "versions": "0+", "about": "Array of acknowledge types - 0:Gap,1:Accept,2:Release,3:Reject,4:Renew."} // Version 2 supports RENEW ack type (KIP-1222) ]} ]} ]}, { "name": "ForgottenTopicsData", "type": "[]ForgottenTopic", "versions": "0+", "about": "The partitions to remove from this share session.", "fields": [ { "name": "TopicId", "type": "uuid", "versions": "0+", "about": "The unique topic ID."}, { "name": "Partitions", "type": "[]int32", "versions": "0+", "about": "The partitions indexes to forget." } ]} ] }clients/src/main/resources/common/message/ShareAcknowledgeRequest.jsonCode Block // ShareAcknowledge { "apiKey": 79, "type": "request", "listeners": ["broker"], "name": "ShareAcknowledgeRequest", // Version 0 was used for early access of KIP-932 in Apache Kafka 4.0 but removed in Apacke Kafka 4.1. // // Version 1 is the initial stable version (KIP-932). // // Version 2 will have RENEW ack type "validVersions": "1-2", "flexibleVersions": "0+", "fields": [ { "name": "GroupId", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null", "entityType": "groupId", "about": "The group identifier." }, { "name": "MemberId", "type": "string", "versions": "0+", "nullableVersions": "0+", "about": "The member ID." }, { "name": "ShareSessionEpoch", "type": "int32", "versions": "0+", "about": "The current share session epoch: 0 to open a share session; -1 to close it; otherwise increments for consecutive requests." }, { "name": "Topics", "type": "[]AcknowledgeTopic", "versions": "0+", "about": "The topics containing records to acknowledge.", "fields": [ { "name": "TopicId", "type": "uuid", "versions": "0+", "about": "The unique topic ID.", "mapKey": true }, { "name": "Partitions", "type": "[]AcknowledgePartition", "versions": "0+", "about": "The partitions containing records to acknowledge.", "fields": [ { "name": "PartitionIndex", "type": "int32", "versions": "0+", "mapKey": true, "about": "The partition index." }, { "name": "AcknowledgementBatches", "type": "[]AcknowledgementBatch", "versions": "0+", "about": "Record batches to acknowledge.", "fields": [ { "name": "FirstOffset", "type": "int64", "versions": "0+", "about": "First offset of batch of records to acknowledge." }, { "name": "LastOffset", "type": "int64", "versions": "0+", "about": "Last offset (inclusive) of batch of records to acknowledge." }, { "name": "AcknowledgeTypes", "type": "[]int8", "versions": "0+", "about": "Array of acknowledge types - 0:Gap,1:Accept,2:Release,3:Reject,4:Renew" } // Version 2 supports RENEW ack type (KIP-1222) ]} ]} ]} ] }
New method exposed in
clients/src/main/java/org/apache/kafka/clients/consumer/ShareConsumer.javainterface to help the application determine RENEW interval.Code Block /** Returns the acquisition lock timeout value in milliseconds for the last set of records fetched from the brokers. This is an optional field as unless an application poll() results in a ShareFetch to a broker, this value cannot be determined. */ public Optional<Integer> acquisitionLockTimeoutMs();
...