Current State: Under Discussion
Item | Value |
|---|---|
Discussion Thread | |
JIRA | |
Release |
Kafka brokers authenticate every client connection and store the authenticated principal (KafkaPrincipal) in memory on each KafkaChannel. However, there is no admin API, CLI command, JMX MBean, or log output that allows an operator to answer:
Which user principals currently have active connections to this broker?
This is a fundamental observability gap. Every comparable system provides this capability:
System | Command / API |
|---|---|
MySQL |
|
PostgreSQL |
|
RabbitMQ | Management API |
MongoDB |
|
Apache Kafka | Nothing |
The broker already holds all the data in memory:
SocketServer
└── DataPlaneAcceptor (one per listener endpoint)
└── Processor (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.
Workaround | Limitation |
|---|---|
Set | Only logs principals when they make requests that trigger authorization. Truly idle connections are invisible. |
Set | Extremely verbose. Still misses connections that send zero requests. |
JMX quota metrics ( | Sensors expire after 600s of inactivity. Requires quotas to be enabled. |
Heap dump ( | Causes GC pause. Requires post-processing. Not suitable for real-time use. |
PushTelemetry RPC. Client-side metrics, not server-side connection listing. Requires client opt-in.ListClientMetricsResources API — lists telemetry configuration resources, not connections.ClientInstanceId in all request headers — enriches request tracing but provides no API to query active connections.None of these address the core gap of listing currently authenticated connections.
A new RPC pair: ListClientConnections.
{
"apiKey": 93,
"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." }
]
}
|
{
"apiKey": 93,
"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." }
]
}
]
}
|
New method on Admin / KafkaAdminClient:
public interface Admin {
/**
* List active client connections on the specified broker(s).
*/
ListClientConnectionsResult listClientConnections(ListClientConnectionsOptions options);
}
|
public class ListClientConnectionsOptions extends AbstractOptions<ListClientConnectionsOptions> {
private String principalFilter;
private String clientAddressFilter;
private String listenerFilter;
public ListClientConnectionsOptions principalFilter(String principal) { ... }
public ListClientConnectionsOptions clientAddressFilter(String address) { ... }
public ListClientConnectionsOptions listenerFilter(String listener) { ... }
}
|
public class ListClientConnectionsResult {
public KafkaFuture<Collection<ClientConnectionInfo>> all() { ... }
}
|
public class ClientConnectionInfo {
public String connectionId() { ... }
public KafkaPrincipal principal() { ... }
public InetAddress clientAddress() { ... }
public int clientPort() { ... }
public String listenerName() { ... }
public SecurityProtocol securityProtocol() { ... }
public String softwareName() { ... }
public String softwareVersion() { ... }
public long connectedSince() { ... }
}
|
New CLI command:
# List all connections on all brokers kafka-client-connections.sh --bootstrap-server <broker>:9092 --list # Filter by principal kafka-client-connections.sh --bootstrap-server <broker>:9092 --list --principal User:alice # Filter by client address kafka-client-connections.sh --bootstrap-server <broker>:9092 --list --client-address 10.0.0.5 # Filter by listener kafka-client-connections.sh --bootstrap-server <broker>:9092 --list --listener SASL_SSL |
Example output:
CONNECTION-ID PRINCIPAL CLIENT-ADDRESS PORT LISTENER PROTOCOL SOFTWARE CONNECTED-SINCE 10.0.0.5:54321-10.0.0.1:9092-0 User:alice 10.0.0.5 54321 SASL_SSL SASL_SSL apache-kafka-java/3.9.0 2026-04-27T10:15:30Z 10.0.0.6:43210-10.0.0.1:9092-1 User:bob 10.0.0.6 43210 SASL_SSL SASL_SSL confluent-kafka-go/2.3.0 2026-04-27T11:22:45Z |
A new method SocketServer.collectClientConnections() encapsulates the connection enumeration:
DataPlaneAcceptor instances (one per listener endpoint)Processor poolselector.channels() to get a snapshot of open KafkaChannel objectsid, principal(), socketAddress(), socketPort(), channelMetadataRegistry().clientInformation()ConnectionInfo objects with all extracted dataThe handler in KafkaApis receives a reference to SocketServer (added as a constructor parameter), calls collectClientConnections(), applies request filters, and builds the response.
This design keeps internal state (Selector, Processor) encapsulated within the kafka.network package while exposing only the aggregated result.
Each NetworkProcessor owns its Selector and runs on its own I/O thread. Selector.channels() returns new ArrayList<>(channels.values()) — a snapshot copy that is safe for cross-thread reads.
POC validation: A working prototype confirmed that calling selector.channels() from the request handler thread (which is different from the I/O threads) works correctly without synchronization. The snapshot semantics of channels() are sufficient — the result may be slightly stale (a connection established or closed during iteration may or may not appear) but this is acceptable for a point-in-time listing.
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).
KafkaChannel does not currently track connection establishment time. This field is included in the response schema for version 0 but will initially return 0 (epoch) until the channel tracking is added. The change is minimal — adding a connectedSince field set during Selector.register() — but is deferred from the initial implementation to keep the patch small. A follow-up PR will populate this field.
The SoftwareName and SoftwareVersion fields are populated from data sent in the client's ApiVersionsRequest (per KIP-511, shipped in Kafka 2.4). Connections that have not yet completed the ApiVersions handshake (e.g., raw TCP connections or very early in the connection lifecycle) will report empty strings for these fields.
File | Change |
|---|---|
| New request schema |
| New response schema |
| Register new API key (93) |
| Request wrapper class |
| Response wrapper class |
| Add parse case for API key 93 |
| Add parse case for API key 93 |
| Add |
| Add |
| Pass |
| New admin client method (future) |
| New CLI tool (future) |
A working POC patch is available at: KIP-1329 POC branch. The POC validates:
ListClientConnections(93): 0)Selector.channels() snapshot semanticsThis API exposes information about authenticated connections, including principals and client addresses.
Operation | Resource | Permission |
|---|---|---|
ListClientConnections | CLUSTER | DESCRIBE |
This is consistent with other cluster-level describe operations (e.g., DescribeCluster, ListClientMetricsResources).
Operators without DESCRIBE on CLUSTER will receive an CLUSTER_AUTHORIZATION_FAILED error.
CLUSTER_AUTHORIZATION_FAILED when caller lacks DESCRIBE on CLUSTER.ListClientConnections to the protocol documentation (protocol.html).kafka-client-connections.sh in the operations section.A JMX MBean could expose connection info, but:
DescribeCluster returns broker metadata, not connection metadata. Overloading it would conflate two different concerns and break the single-responsibility principle of the API.
All log-based approaches (authorizer log, request log) are fundamentally incomplete because they only capture connections that actively send requests. Truly idle connections — which are the most security-relevant during incident response — remain invisible.