Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.

...

  • There will be new addition in the org/apache/kafka/clients/consumer/AcknowledgeType.java  enum where we will add a new entry RENEW with id 4. 

    Code Block
    public enum AcknowledgeType {
        /** The record was consumed successfully. */
        ACCEPT((byte) 1),
        /** The record was not consumed successfully. Release it for another delivery attempt. */
        RELEASE((byte) 2),
        /** The record was not consumed successfully. Reject it and do not release it for another delivery attempt. */
        REJECT((byte) 3),
        /** Consumer needs more time to process the record. Renew the lease. */
        RENEW((byte) 4);	// New entry per KIP-1222
        ...
    }


  • The acknowledgement type will be used in the existing ShareFetch and ShareAcknowledge RPCs. Hence, we must add new versions for them.
    • clients/src/main/resources/common/message/ShareFetchRequest.json 

      Code 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": "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.json  

      Code 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.java interface 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.
    */
    public Optional<Integer> acquisitionLockTimeoutMs();
    New exception class will be added clients/src/main/java/org/apache/kafka/common/errors/UnsupportedAcknowledgeTypeException.java  Code Blockpublic class UnsupportedAcknowledgeTypeException extends InvalidConfigurationException { private static final long serialVersionUID = 1L; public UnsupportedAcknowledgeTypeException(String message, Throwable cause) { super(message, cause); } public UnsupportedAcknowledgeTypeException(String message) { super(message); } }

     

Proposed Changes

KIP-932 added two new RPCs for record fetching and acknowledging: ShareFetch and ShareAcknowledge. From the share consumer, these RPCs are called when an application invokes poll(Duration)and commitSync(Duration)/commitAsync().

...

There is a case where an application tries to send a RENEW acknowledgement to an old broker which does not support the new type. To remedy this, we will also update the ShareConsumer.acknowledge(ConsumerRecord, AcknowledgeType) to throw a new exception UnsupportedAcknowledgeTypeExceptionan IllegalArgumentException. This will be handled completely on the client side where the share consumer will use the RPC API version to determine acknowledgement type validity.

...