Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.
Comment: Member ID and client instance ID alignment clarification

...

Proposed Changes

This KIP proposes two the following changes to the Kafka protocol.

...

This KIP proposes adding a UUID called the client instance ID into the request header of all Kafka protocol requests. Those familiar with KIP-714 will be aware that it already introduced the client instance ID, so this KIP actually proposes extending its scope to elevating that concept into a universal unique identifier for client instances in all RPCs.

The client instance ID is now calculated by the client during the constructor of the client before it makes its initial connection to the cluster. This differs slightly from how KIP-714 initialized the client instance ID, but the change is compatible with the existing behavior of the broker. The client uses the same client instance ID for its connections to every broker throughout its lifetime, even when it rebootstraps and makes new connections. The consistency and uniqueness of the identifiers are both important characteristics.The

In KIP-714, the client instance ID was created by the broker which responded to a client's first GetTelemetrySubscriptions RPC. In KIP-848, the member ID was created by the group coordinator in response to a heartbeat, but subsequently KIP-1082 changed this so that the client created its own member ID. As a result, this KIP calculates the client instance ID on the client.

The client instance ID is added as a tagged field in the RPC request header, so it is is added as a tagged field in the RPC request header, so it is present on all RPCs that use the v2 request header, which is almost all of the current versions of the RPCs (only SaslHandshake and OffsetDelete use the v1 request header for their latest versions). The client is required to be consistent in its use of client instance ID in request headers. If it specifies the value in its request headers, every request on a connection must specify the same value. If it does not specify the value in its request headers, every request on a connection must not specified a value.

If a connection specifies Note that omitting a client instance ID in from the request header of its first request which uses the v2 request header, it must specify the same and explicitly sending a zero UUID as the client instance ID in the request header are indistinguishable in the protocol and are considered semantically equivalent. This means that if a client explicitly sets a zero UUID, the broker will treat it as if the client had not set a client instance ID. When the following text says "does not specify a client instance ID", this includes specifying a zero UUID as the client instance ID.

If a connection specifies for all subsequent requests which use the v2 request header. The initial client instance ID for each connection will be cached by the broker for checking (this is an implementation detail, but caching it in the ChannelMetadataRegistry is an option). Once a client has specified a client instance ID in the request header of its first request, any subsequent requests which are missing the client instance ID (with the exception of requests using the v1 request header) or which specify a different value for the client instance ID will be rejected with error code INVALID_REQUEST.If a client does not specify a client instance ID in the request header of its first request which uses the v2 request header, it must not specify a the same client instance ID in the request header of any for all subsequent requests . If it does so, the request will be rejected with error code INVALID_REQUEST.

In KIP-714, the client instance ID was created by the broker which responded to a client's first GetTelemetrySubscriptions RPC. In KIP-848, the member ID was created by the group coordinator in response to a heartbeat, but subsequently KIP-1082 changed this so that the client created its own member ID. As a result, this KIP also changes the definition of the client instance ID so that the client creates it. After this KIP, when the client makes its first GetTelemetrySubscriptions request, it supplies the client instance ID which it has already created rather than specifying a zero client instance ID; the Apache Kafka Java client no longer supplies a zero client instance ID on GetTelemetrySubscriptions requests. The original behavior of KIP-714 is still supported for clients which do not yet send ClientInstanceId in the request header; they just specify a zero client instance ID in their first GetTelemetrySubscriptions request (v0) and then receive a client instance ID to use for client telemetry calculated by the broker in its response.

Member ID in group protocols

Consumer groups with the consumer group protocol (KIP-848 and KIP-1082), share groups (KIP-932) and streams groups (KIP-1071) all make use of a client-generated UUID as the member ID. This KIP proposes using the same client-generated UUID as the client instance ID and the member ID. This makes no change to the protocol at all. It simply creates a single UUID when a client is constructed, and when a client-generated member ID is required for the group protocol, it uses the client instance ID.

Public Interfaces

Client API Changes

The client instance ID is calculated during the constructor of the Producer , Consumer , ShareConsumer and Admin  implementations so there is no need to have a timeout parameter on the accessor method. The following method is added to these interfaces:

public Uuid clientInstanceId()

and then the following method is deprecated for removal in Apache Kafka 5.0:

public Uuid clientInstanceId(Duration timeout)

In a similar vein, the following method in KafkaStreams is added:

public ClientInstanceIds clientInstanceIds()

and then the following method is deprecated for removal in Apache Kafka 5.0:

public ClientInstanceIds clientInstanceIds(Duration timeout)

Kafka Protocol Changes

Request Header

This KIP introduces a tagged field ClientInstanceId into version 2 of the request header. This means it can be introduced without any other RPC changes.

Code Block
{
  "type": "header",
  "name": "RequestHeader",
  // Version 0 was removed in Apache Kafka 4.0, Version 1 is the new baseline.
  //
  // Version 0 of the RequestHeader is only used by v0 of ControlledShutdownRequest.
  //
  // Version 1 is the first version with ClientId.
  //
  // Version 2 is the first flexible version.
  // Tagged field 0 introduces ClientInstanceId. (KIP-1313)
  "validVersions": "1-2",
  "flexibleVersions": "2+",
  "fields": [
    { "name": "RequestApiKey", "type": "int16", "versions": "0+",
      "about": "The API key of this request." },
    { "name": "RequestApiVersion", "type": "int16", "versions": "0+",
      "about": "The API version of this request." },
    { "name": "CorrelationId", "type": "int32", "versions": "0+",
      "about": "The correlation ID of this request." },

    // The ClientId string must be serialized with the old-style two-byte length prefix.
    // The reason is that older brokers must be able to read the request header for any
    // ApiVersionsRequest, even if it is from a newer version.
    // Since the client is sending the ApiVersionsRequest in order to discover what
    // versions are supported, the client does not know the best version to use.
    { "name": "ClientId", "type": "string", "versions": "1+", "nullableVersions": "1+", "flexibleVersions": "none",
      "about": "The client ID, for identifying an application." },
    { "name": "ClientInstanceId", "type": "uuid", "versions": "2+", "taggedVersions": "2+", "tag": 0, "ignorable": "true",
      "about": "The client instance ID, for identifying an instance of an application." }
  ]
}

GetTelemetrySubscriptions API

This KIP introduces GetTelemetrySubscriptions v1.

Request

If the client uses v0, the original KIP-714 behavior is used. The client sends the value 0 for the ClientInstanceId in order to 

After this KIP, the Apache Kafka Java client will send the ClientInstanceId  in the request header and also in the request body. If present in the request header, it must match the value in the request body, and that value will not be zero, or else the request will be rejected with error code INVALID_REQUEST.

It is still permitted to send the value zero for the ClientInstanceId  in the request, which will cause the broker to create a ClientInstanceId  and send it back in the response. This is the original KIP-714 behavior and it is still supported for clients which do not use the ClientInstanceId  in the request header.

Code Block
{
  "apiKey": 71,
  "type": "request",
  "listeners": ["broker"],
  "name": "GetTelemetrySubscriptionsRequest",
  // Version 0 is the initial version. (KIP-714)
  //
  // Version 1 removes ClientInstanceId from the request and response. (KIP-1313)
  "validVersions": "0-1",
  "flexibleVersions": "0+",
  "fields": [
    {
      "name": "ClientInstanceId", "type": "uuid", "versions": "0", << Only v0 supports this
      "about": "Unique id for this client instance, must be set to 0 on the first request (v0 only)."
    }
  ]
}

Response

The ClientInstanceId  is removed in the v1 schema.

Code Block
{
  "apiKey": 71,
  "type": "response",
  "name": "GetTelemetrySubscriptionsResponse",
  // Version 0 is the initial version. (KIP-714)
  //
  // Version 1 removes ClientInstanceId from the request and response. (KIP-1313)
  "validVersions": "0-1",
  "flexibleVersions": "0+",
  "fields": [
    {
      "name": "ThrottleTimeMs", "type": "int32", "versions": "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": "ErrorCode", "type": "int16", "versions": "0+",
      "about": "The error code, or 0 if there was no error."
    },
    {
      "name": "ClientInstanceId", "type": "uuid", "versions": "0", << Only v0 supports this
      "about": "Assigned client instance id if ClientInstanceId was 0 in the request, else 0 (v0 only)."
    },
    {
      "name": "SubscriptionId", "type": "int32", "versions": "0+",
      "about": "Unique identifier for the current subscription set for this client instance."
    },
    {
      "name": "AcceptedCompressionTypes", "type": "[]int8", "versions": "0+",
      "about": "Compression types that broker accepts for the PushTelemetryRequest."
    },
    {
      "name": "PushIntervalMs", "type": "int32", "versions": "0+",
      "about": "Configured push interval, which is the lowest configured interval in the current subscription set."
    },
    {
      "name": "TelemetryMaxBytes", "type": "int32", "versions": "0+",
      "about": "The maximum bytes of binary data the broker accepts in PushTelemetryRequest."
    },
    {
      "name": "DeltaTemporality", "type": "bool", "versions": "0+",
      "about": "Flag to indicate monotonic/counter metrics are to be emitted as deltas or cumulative values."
    },
    {
      "name": "RequestedMetrics", "type": "[]string", "versions": "0+",
      "about": "Requested metrics prefix string match. Empty array: No metrics subscribed, Array[0] empty string: All metrics subscribed."
    }
  ]
}

PushTelemetry API

This KIP introduces PushTelemetry v1.

Request

v1 removes ClientInstanceId  which is now sent by the client in the request header. The ClientInstanceId  in the request header must not be zero, or else the request will be rejected with error code INVALID_REQUEST .

which use the v2 request header. The initial client instance ID for each connection will be cached by the broker for checking (this is an implementation detail, but caching it in the ChannelMetadataRegistry is an option). Once a client has specified a client instance ID in the request header of its first request, any subsequent requests which are missing the client instance ID (with the exception of requests using the v1 request header) or which specify a different value for the client instance ID will be rejected with error code INVALID_REQUEST.

If a client does not specify a client instance ID in the request header of its first request which uses the v2 request header, it must not specify a client instance ID in the request header of any subsequent requests. If it does so, the request will be rejected with error code INVALID_REQUEST.

Alignment of identifiers

By adding client instance ID to the request headers, we now had a unique application instance identifier which we can use for other purposes such as the member ID in the group protocols and the client telemetry client instance ID. The client MUST use the same UUID in request headers as it uses for the client telemetry client instance ID. The alignment of other identifiers is by convention (and the Java client will follow the convention) rather than mandate. In the language of standards, the client SHOULD use the same UUID in request headers as it uses for the member ID in the group protocols. This alignment just makes traceability and problem determination more straightforward.

In the modern group protocol RPCs such as ConsumerGroupHeartbeat  and ShareGroupHeartbeat , the member ID is a string. In practice, it is a UUID which is encoded into a string, but the nature of this conversion is not specified in the protocol and it really is treated as a string in the broker. The Java client uses the org.apache.kafka.common.Uuid  class to generate the member ID and convert it into a string. The client instance ID really is a UUID in the protocol. When the Java code in the broker converts this into a string, it uses the same Java code as the client does for the member ID. As a result, the member ID and client instance ID can trivially be the same when represented as strings, even though they have different data types in the protocol. For traceability and problem determination, this "conventional" alignment works well.

However, at least one other Kafka client implementation uses a slightly different string encoding of the member ID which does not match that generated by org.apache.kafka.common.Uuid . This is valid with regards to KIP-848 and KIP-932, but does unfortunately mean that the member ID does not look the same as the client instance ID, even if they have the same original UUID value in the client.

It would be possible to go around the existing RPCs such as ConsumerGroupHeartbeat  and GetTelemetrySubscriptions , and remove the fields containing the existing identifiers which are intended to be aligned. Doing so would be a bad idea though, because we would then have RPC versions which essentially depend upon the presence of a tagged field in the request header. This is a protocol-compatibility nightmare because tagged fields are optional and having a hard dependency on an optional field seems unwise.

This KIP makes one change to the GetTelemetrySubscriptions  behavior. The client may only request a new client instance ID on its initial GetTelemetrySubscriptions  request (request.clientInstanceId  is 0), if it does not also send a client instance ID in the request header. After this KIP, the broker is only expecting to generate telemetry client instance ID for older clients which do not use the request header. This will automatically align the UUID in the request headers and client telemetry.

Key:

  • UUID-H - client instance ID sent by client in header, generated by client
  • UUID-R - client instance ID sent by client in request, which is not equal to UUID-H
  • UUID-B - client instance ID sent by broker in response, generated by broker

Pre-KIP-1313 - broker does not expect ClientInstanceId  in request header and ignores it

Client sends

Broker responds

Notes

GetTelemetrySubscriptions v0

request.ClientInstanceId = 0

response.ClientInstanceId = UUID-B

response.ErrorCode = NONE

This is KIP-714 initial GetTelemetrySubscriptions .

Client is requesting a new client instance ID from the broker.

The client will henceforth use UUID-B for client telemetry.

GetTelemetrySubscriptions v0

request.ClientInstanceId = UUID-R

response.ClientInstanceId = 0

response.ErrorCode = NONE

This is KIP-714 non-initial GetTelemetrySubscriptions .

The client is using UUID-R for client telemetry.

Post-KIP-1313 - broker is aware of optional ClientInstanceId  in request header

Client sends

Broker responds

Notes

GetTelemetrySubscriptions v0

header.clientInstanceId not present or 0

request.ClientInstanceId = 0

response.ClientInstanceId = UUID-B

response.ErrorCode = NONE

This is KIP-714 initial GetTelemetrySubscriptions  request from a pre-KIP-1313 client.

The client is requesting a new client instance ID from the broker.

The client will henceforth use UUID-B for client telemetry.

GetTelemetrySubscriptions v0

header.ClientInstanceId not present or 0

request.ClientInstanceId = UUID-R

response.ClientInstanceId = 0

response.ErrorCode = NONE

This is KIP-714 non-initial GetTelemetrySubscriptions  request from a pre-KIP-1313 client.

The client is using UUID-R for client telemetry.

GetTelemetrySubscriptions v0

header.ClientInstanceId = UUID-H

request.ClientInstanceId = UUID-H

response.ClientInstanceId = 0

response.ErrorCode = NONE

This is KIP-714 GetTelemetrySubscriptions  request from a post-KIP-1313 client.

The client is using UUID-H for request headers and client telemetry. This is what we expect from the Apache Kafka Java client after this KIP.

GetTelemetrySubscriptions v0

header.ClientInstanceId = UUID-H

request.ClientInstanceId = 0

response.ErrorCode =  INVALID_REQUEST

This is not allowed.

If a client sends UUID-H in the request header, the request client instance ID must also be UUID-H.

After this KIP, the client generates the client instance ID and does not ask the broker to do so.

GetTelemetrySubscriptions v0

header.ClientInstanceId = UUID-H

request.ClientInstanceId = UUID-R≠UUID-H

response.ErrorCode =  INVALID_REQUEST

This is not allowed.

If a client sends UUID-H in the request header, the request client instance ID must also be UUID-H.

In summary, for GetTelemetrySubscriptions v0, here are the combinations:


Old brokerNew broker
Old client

Initial request:

  • request.ClientInstanceID=0
  • response.ClientInstanceId=UUID-B

Subsequent requests:

  • request.ClientInstanceId=UUID-B
  • response.ClientInstanceId=0

Initial request:

  • header.ClientInstanceId=0 (not sent by the client, but its absence is treated as 0)
  • request.ClientInstanceID=0
  • response.ClientInstanceId=UUID-B

Subsequent requests:

  • header.ClientInstanceId=0 (not sent by the client, but its absence is treated as 0)
  • request.ClientInstanceId=UUID-B
  • response.ClientInstanceId=0
New client

Initial request:

  • header.ClientInstanceId=UUID-H (ignored by broker)
  • request.ClientInstanceId=UUID-H
  • response.ClientInstanceId=0

Subsequent requests:

  • header.ClientInstanceId=UUID-H (ignored by broker)
  • request.ClientInstanceId=UUID-H
  • response.ClientInstanceId=0 

Initial request:

  • header.ClientInstanceId=UUID-H
  • request.ClientInstanceID=UUID-H
  • response.ClientInstanceId=0

Subsequent requests:

  • header.ClientInstanceId=UUID-H
  • request.ClientInstanceId=UUID-H
  • response.ClientInstanceId=0

Public Interfaces

Client API Changes

The client instance ID is calculated during the constructor of the Producer , Consumer , ShareConsumer and Admin  implementations so there is no need to have a timeout parameter on the accessor method. The following method is added to these interfaces:

public Uuid clientInstanceId()

and then the following method is deprecated for removal in Apache Kafka 5.0:

public Uuid clientInstanceId(Duration timeout)

In a similar vein, the following method in KafkaStreams is added:

public ClientInstanceIds clientInstanceIds()

and then the following method is deprecated for removal in Apache Kafka 5.0:

public ClientInstanceIds clientInstanceIds(Duration timeout)

Kafka Protocol Changes

Request Header

This KIP introduces a tagged field ClientInstanceId into version 2 of the request header. This means it can be introduced without any other RPC changes.

Code Block
{
  "type": "header",
  "name": "RequestHeader",
  // Version 0 was removed in Apache Kafka 4.0, Version 1 is the new baseline.
  //
  // Version 0 of the RequestHeader is only used by v0 of ControlledShutdownRequest.
  //
  // Version 1 is the first version with ClientId.
Code Block
{
  "apiKey": 72,
  "type": "request",
  "listeners": ["broker"],
  "name": "PushTelemetryRequest",
  // Version 0 is the initial version. (KIP-714)
  //
  // Version 1 moves ClientInstanceId from the request to the request header 2 is the first flexible version.
  // Tagged field 0 introduces ClientInstanceId. (KIP-1313)
   "validVersions": "01-12",
  "flexibleVersions": "02+",
  "fields": [
    {
      "name": "ClientInstanceIdRequestApiKey", "type": "uuidint16", "versions": "0+", << Only v0 supports this.
      "about": "UniqueThe API idkey forof this client instancerequest."
    },
    {
      "name": "SubscriptionIdRequestApiVersion", "type": "int32int16", "versions": "0+",
      "about": "UniqueThe identifierAPI forversion theof currentthis subscriptionrequest."
    },
    {
      "name": "TerminatingCorrelationId", "type": "boolint32", "versions": "0+",
      "about": "ClientThe correlation isID terminatingof thethis connectionrequest." },

    // The ClientId string must be serialized with the old-style two-byte length  },prefix.
    {
// The reason is that older "name": "CompressionType", "type": "int8", "versions": "0+",
      "about": "Compression codec used to compress the metrics."
    },
    {
      "name": "Metrics", "type": "bytes", "versions": "0+", "zeroCopy": true,
      "about": "Metrics encoded in OpenTelemetry MetricsData v1 protobuf format."
    }
  ]
}

Response

No change to the schema with v1.

...

brokers must be able to read the request header for any
    // ApiVersionsRequest, even if it is from a newer version.
    // Since the client is sending the ApiVersionsRequest in order to discover what
    // versions are supported, the client does not know the best version to use.
    { "name": "ClientId", "type": "string", "versions": "1+", "nullableVersions": "1+", "flexibleVersions": "none",
      "about": "The client ID, for identifying an application." },
    { "name": "ClientInstanceId", "type": "uuid", "versions": "2+", "taggedVersions": "2+", "tag": 0, "ignorable": "true",
      "about": "The client instance ID, for identifying an instance of an application." }
  ]
}

GetTelemetrySubscriptions API

A very small behavioral change is made in the broker handling of GetTelemetrySubscriptions v0.

A post-KIP-1313 client will not send a zero ClientInstanceId  in the request body to ask the broker to assign an ID.

If the client sends a zero ClientInstanceId  in the request body, it must not send a ClientInstanceId  in the request header.

Similarly, if the client sends a ClientInstanceId  in the request header and also sends a non-zero ClientInstanceId  in the request body, the values must be the same. It is not permitted to use one client instance ID for request headers and a different client instance ID for telemetry.

As a result, the description for the ClientInstanceId  field in the request becomes "Unique id for this client instance. If client sends ClientInstanceId in the header, must equal that value. If not, must be set to 0 on the first request." .

Compatibility, Deprecation, and Migration Plan

...