Versions Compared

Key

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

Table of Contents

Status

Current state: VotingAccepted

Discussion thread: here

Voting thread: here

...

  • 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 field IsRenewAck defaulting to false to both the RPCs. This field will help in optimising the broker side implementation for processing renew acknowledgements.
    • 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": "IsRenewAck", "type": "bool", "versions": "2+", "default": false,
            "about": "Whether renew type acknowledgement present in AcknowledgementBatches." }, // Version 2 supports RENEW ack type and this field serves as an indicator (KIP-1222)
          { "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": "IsRenewAck", "type": "bool", "versions": "2+", "default": false,
            "about": "Whether renew type acknowledgement present in AcknowledgementBatches." }, // Version 2 supports RENEW ack type and this field serves as an indicator (KIP-1222)
          { "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.
     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();

     

  • An additional tag for the metric RecordAcknowledgementsPerSec, defined in KIP-1103, will get introduced as a side effect of introducing the RENEW ack type. The new definition of the metric is: 

    RationaleMetric NameTypeGroupTagsDescriptionJMX Bean
    Tracks the rate of records acknowledged per acknowledgement type.RecordAcknowledgementsPerSecMeterShareGroupMetricsackType:{Accept|Release|Reject|Renew}The rate per second of records acknowledged per acknowledgement type.

    kafka.server:type=ShareGroupMetrics,name=RecordAcknowledgementsPerSec,ackType={Accept|Release|Reject|Renew}



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().

...