Versions Compared

Key

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

...

Item

Value

Discussion Thread

TBD — post to dev@kafka.apache.org and link heredev@ thread

JIRA

KAFKA-20526

Release


Motivation

...

The broker already holds all the data in memory:

Code Block
SocketServer
  └── DataPlaneAcceptor (one per listener endpoint)
        └── NetworkProcessorProcessor (one per network thread)
              └── Selector
                    └── channels: Map[String, KafkaChannel]
                          └── KafkaChannel
                                ├── principal(): KafkaPrincipal    ← authenticated user
                                ├── socketAddress(): InetAddress      ← client IP
                                ├── channelMetadataRegistry
                                │     └── clientInformation         ← software name/version (KIP-511)
                                └── id: String                      ← connection ID

The data is simply not surfaced through any external interface.

...

Code Block
javascript
javascript
{
  "apiKey": TBD93,
  "type": "request",
  "name": "ListClientConnectionsRequest",
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "PrincipalFilter", "type": "string", "versions": "0+",
      "nullableVersions": "0+",
      "about": "If set, only return connections for this principal. Null returns all." },
    { "name": "ClientAddressFilter", "type": "string", "versions": "0+",
      "nullableVersions": "0+",
      "about": "If set, only return connections from this client address. Null returns all." },
    { "name": "ListenerFilter", "type": "string", "versions": "0+",
      "nullableVersions": "0+",
      "about": "If set, only return connections on this listener. Null returns all." }
  ]
}

...

Code Block
javascript
javascript
{
  "apiKey": TBD93,
  "type": "response",
  "name": "ListClientConnectionsResponse",
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+",
      "about": "Duration in milliseconds for which the request was throttled." },
    { "name": "ErrorCode", "type": "int16", "versions": "0+" },
    { "name": "Connections", "type": "[]ConnectionInfo", "versions": "0+",
      "about": "List of active connections on this broker.",
      "fields": [
        { "name": "ConnectionId", "type": "string", "versions": "0+",
          "about": "The broker-assigned connection ID." },
        { "name": "Principal", "type": "string", "versions": "0+",
          "about": "The authenticated principal, e.g. User:alice." },
        { "name": "PrincipalType", "type": "string", "versions": "0+",
          "about": "The principal type, e.g. User." },
        { "name": "ClientAddress", "type": "string", "versions": "0+",
          "about": "The remote client IP address." },
        { "name": "ClientPort", "type": "int32", "versions": "0+",
          "about": "The remote client port." },
        { "name": "ListenerName", "type": "string", "versions": "0+",
          "about": "The listener the client connected to." },
        { "name": "SecurityProtocol", "type": "string", "versions": "0+",
          "about": "The security protocol: PLAINTEXT, SSL, SASL_PLAINTEXT, or SASL_SSL." },
        { "name": "SoftwareName", "type": "string", "versions": "0+",
          "about": "The client software name (from ApiVersionsRequest, per KIP-511)." },
        { "name": "SoftwareVersion", "type": "string", "versions": "0+",
          "about": "The client software version (from ApiVersionsRequest, per KIP-511)." },
        { "name": "ConnectedSince", "type": "int64", "versions": "0+",
          "about": "Timestamp (epoch ms) when the connection was established." }
      ]
    }
  ]
}

...

  1. Iterates over all DataPlaneAcceptor instances (one per listener endpoint)
  2. For each acceptor, iterates its Processor pool
  3. For each processor, calls selector.channels() to get a snapshot of open KafkaChannel objects
  4. For each channel, extracts: id, principal(), socketAddress(), socketPort(), channelMetadataRegistry().clientInformation()
  5. Returns a Seqcollection of ConnectionInfo objects with all extracted data

The handler in KafkaApis receives a reference to SocketServer (added as a constructor parameter), calls collectClientConnections(), applies request filters, and builds the response.

...

The implementation adds a collectClientConnections() method on SocketServer that iterates all processors and their selectors, collecting channel metadata into a response-ready structure. This keeps the access pattern encapsulated within the kafka.network package where Selector is accessible (it is package-private to kafka.network).

ConnectedSince Tracking

...