Current state: Under discussion
Discussion thread: here
JIRA: TBD
Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).
Kafka Streams applications are opaque to cluster operators. When an application hits performance degradation or a protracted rebalance, operators cannot inspect the processing graph — the individual source, processor, and sink nodes, their predecessor/successor relationships, and the state stores they touch — without access to the application source code or logs. KIP-1071 (Streams Rebalance Protocol) already has clients send a topology to the broker, but only the subset the coordinator needs for assignment: subtopologies, source topics, changelog topics, and copartition groups. The full processing graph is not part of that assignment topology and is not currently available on the broker.
This KIP proposes a mechanism to send the full topology description from the client to the broker, where a pluggable backend can store and expose it for operational tooling and topology visualization in management UIs. The design follows the pattern established by KIP-714 (Client Metrics and Observability): the broker acts as a conduit, receiving data from clients and delegating storage and presentation to a plugin implementation. This keeps the broker itself simple and the feature extensible.
| Configuration name | Description | Values |
|---|---|---|
group.streams.topology.description.plugin.class | The fully qualified class name of a StreamsGroupTopologyDescriptionPlugin implementation. When not set, the feature is disabled. | Type: class, Default: empty string |
| Configuration name | Description | Values |
|---|---|---|
topology.description.push.enabled | Controls whether the Kafka Streams client sends topology descriptions to the broker when requested. When set to | Type: boolean, Default: true |
StreamsGroupHeartbeatResponse is bumped to the next version (N) and gains a new field at that version:
{ "name": "TopologyDescriptionRequired", "type": "bool", "versions": "N+", "default": "false",
"about": "True if the client should send the topology description via UpdateStreamsGroupTopologyDescription." } |
The broker sets this field to true when a topology description plugin is configured and plugin.requiresTopologyPush(requestContext, groupId, groupCreationTimeMs, topologyEpoch) returns true, where requestContext is the context of the heartbeat request.
A new RPC is introduced for setting the topology description for a streams group. Like StreamsGroupHeartbeat, the request is sent to the group coordinator for the group.
Request:
{
"apiKey": "TBD",
"type": "request",
"listeners": ["broker"],
"name": "UpdateStreamsGroupTopologyDescriptionRequest",
"validVersions": "0",
"flexibleVersions": "0+",
"fields": [
{ "name": "GroupId", "type": "string", "versions": "0+", "entityType": "groupId",
"about": "The streams group identifier." },
{ "name": "TopologyEpoch", "type": "int32", "versions": "0+",
"about": "The epoch of the topology being described." },
{ "name": "TopologyDescription", "type": "TopologyDescription", "versions": "0+",
"about": "The topology description." }
],
"commonStructs": [
{ "name": "TopologyDescription", "versions": "0+", "fields": [
{ "name": "Subtopologies", "type": "[]Subtopology", "versions": "0+",
"about": "The subtopologies that make up this topology." },
{ "name": "GlobalStores", "type": "[]GlobalStore", "versions": "0+",
"about": "Global state stores used by this topology." }
]},
{ "name": "Subtopology", "versions": "0+", "fields": [
{ "name": "SubtopologyId", "type": "string", "versions": "0+",
"about": "The subtopology identifier, unique within the topology." },
{ "name": "Nodes", "type": "[]TopologyNode", "versions": "0+",
"about": "The processing nodes in this subtopology." }
]},
{ "name": "TopologyNode", "versions": "0+", "fields": [
{ "name": "Name", "type": "string", "versions": "0+",
"about": "The name of this node (e.g., KSTREAM-SOURCE-0000000000)." },
{ "name": "NodeType", "type": "int8", "versions": "0+",
"about": "The type of this node: 1=SOURCE, 2=PROCESSOR, 3=SINK." },
{ "name": "SourceTopics", "type": "[]string", "versions": "0+", "entityType": "topicName",
"about": "The source topics this node reads from. Defined only for source nodes, may be empty if source topics are dynamically determined." },
{ "name": "SinkTopic", "type": "string", "versions": "0+", "entityType": "topicName",
"nullableVersions": "0+", "default": "null",
"about": "The topic this node writes to. Defined only for sink nodes, may be null if sink topic is dynamically determined." },
{ "name": "Stores", "type": "[]string", "versions": "0+",
"about": "The state store names accessed by this node. Defined only for processor nodes." },
{ "name": "Successors", "type": "[]string", "versions": "0+",
"about": "The names of successor nodes in the processing graph. Predecessor relationships are reconstructed from this field on the read side." }
]},
{ "name": "GlobalStore", "versions": "0+", "fields": [
{ "name": "Source", "type": "TopologyNode", "versions": "0+",
"about": "The source node providing data to the global store." },
{ "name": "Processor", "type": "TopologyNode", "versions": "0+",
"about": "The processor node that populates the global store." }
]}
]
} |
Response:
{
"apiKey": "TBD",
"type": "response",
"name": "UpdateStreamsGroupTopologyDescriptionResponse",
"validVersions": "0",
"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 top-level error code, or 0 if there was no error." },
{ "name": "ErrorMessage", "type": "string", "versions": "0+",
"nullableVersions": "0+", "default": "null",
"about": "The top-level error message, or null if there was no error." }
]
} |
Authorization: Requires READ ACL on the GROUP resource for the given group ID. Like offset commits, we don't consider this a modification of the GROUP. This allows deploying apps with READ ACLs on the group.
Error codes:
GROUP_AUTHORIZATION_FAILED — the client is not authorizedINVALID_REQUEST — the request is malformed, or the plugin semantically rejected the payload by completing its future with InvalidRequestExceptionUNSUPPORTED_VERSION — the coordinator cannot serve this RPC because no topology description plugin is configuredTOPOLOGY_DESCRIPTION_TOO_LARGE — the plugin rejected the description because it exceeds the size the plugin is willing to storeTOPOLOGY_DESCRIPTION_UPDATE_FAILED — the plugin failed to process the request for some other reason; the client logs the underlying error at INFO levelGROUP_ID_NOT_FOUND — the specified group does not existNOT_COORDINATOR — the broker is not the coordinator for this groupCOORDINATOR_NOT_AVAILABLE — the coordinator is not availableCOORDINATOR_LOAD_IN_PROGRESS — the coordinator is loadingClient-side retry behavior for each code is described in Error Handling and Retries.
StreamsGroupDescribeRequest is bumped to the next version (N) and gains a new field at that version:
{"name": "IncludeTopologyDescription", "type": "bool", "versions": "N+", "default": "false",
"about": "Whether to include the full topology description from the topology description plugin in the response." } |
A client that negotiates an older version is handled unchanged by any broker. The flag may only be set when version N or later is negotiated.
StreamsGroupDescribeResponse is bumped to the next version (N). Two new fields are added to each DescribedGroup at that version:
{ "name": "TopologyDescription", "type": "TopologyDescription", "versions": "N+",
"nullableVersions": "N+", "default": "null",
"about": "The topology description for this group. Null if not available — see TopologyDescriptionStatus for the reason." },
{ "name": "TopologyDescriptionStatus", "type": "int8", "versions": "N+", "default": "0",
"about": "The status of the topology description for this group: 0=NOT_REQUESTED (client did not set IncludeTopologyDescription), 1=NOT_STORED (no topology description has been recorded for this group), 2=ERROR (the broker failed to fetch the topology description; check broker logs), 3=AVAILABLE (a topology description is present in the TopologyDescription field)." } |
The TopologyDescription common struct mirrors the struct used by UpdateStreamsGroupTopologyDescriptionRequest (same field names and shape). Because Kafka RPC schemas do not share common structs across message files, the struct is duplicated in StreamsGroupDescribeResponse.json.
Setting these fields does not change the ErrorCode on the DescribedGroup: a group with a successful describe but a missing or failed topology fetch still returns ErrorCode=NONE. The TopologyDescriptionStatus field tells the caller why TopologyDescription is null, so that "waiting for first push" (NOT_STORED) can be distinguished from "broker-side fetch failed" (ERROR) without an error-level change to the describe result.
A new interface is introduced in org.apache.kafka.coordinator.group.api.streams :
package org.apache.kafka.coordinator.group.api.streams;
import org.apache.kafka.common.Configurable;
import org.apache.kafka.server.authorizer.AuthorizableRequestContext;
import java.util.concurrent.CompletableFuture;
/**
* A broker-side plugin that manages topology descriptions for streams groups.
*
* <p>Implementations receive topology descriptions pushed by Kafka Streams clients
* and can store, forward, or expose them however they see fit.
*
* <p>The broker calls {@link #requiresTopologyPush} on every heartbeat to determine
* whether the client should send its topology description. Implementations that need
* to consult an external service should kick off that work asynchronously on the
* first call and return {@code false} until the result is available.
*
* <p>Every method takes an {@link AuthorizableRequestContext} as its first argument.
*
* <p>Implementations must be thread-safe. {@link #setTopology} may be called
* concurrently by multiple group members observing
* {@code TopologyDescriptionRequired=true} in the same heartbeat cycle; concurrent
* calls with the same {@code (groupId, topologyEpoch)} pair carry identical data
* and must be treated as idempotent.
*/
public interface StreamsGroupTopologyDescriptionPlugin extends Configurable, AutoCloseable {
/**
* Returns whether the broker should request a topology push from the client.
*
* <p>Called on every successful heartbeat. Not called for members in
* {@code STALE_TOPOLOGY} status. If this method returns {@code true}, the broker
* sets {@code TopologyDescriptionRequired=true} in the heartbeat response.
*
* <p>This method should not throw. Failure modes should be handled internally
* and converted to a return value of {@code false}. The broker defensively
* catches any exception and treats it as {@code false}, logging at WARN; plugins
* should not rely on this backstop.
*
* <p>This method is on the heartbeat path and must return quickly. Implementations
* that need to consult an external service should return {@code false} until the
* result is available.
*
* <p>See the KIP's <em>Plugin Implementation Guidelines</em> section for how
* implementations should handle in-flight tracking, retries, and topology
* expiration.
*
* @param requestContext the context of the heartbeat request
* @param groupId the streams group ID
* @param groupCreationTimeMs the timestamp when the group was created. A value
* of {@code 0} means "unset". Plugins should treat {@code 0}
* as "incarnation indistinguishable" and should not cross-check
* stored epoch values against another stored description that
* also had {@code groupCreationTimeMs == 0}.
* @param topologyEpoch the topology epoch
* @return true if the broker should request a topology push from the client
*/
boolean requiresTopologyPush(AuthorizableRequestContext requestContext,
String groupId, long groupCreationTimeMs, int topologyEpoch);
/**
* Called when a client sends a topology description for a streams group.
* This method may be called concurrently by multiple members of the same group;
* all calls for the same (groupId, topologyEpoch) carry identical data.
*
* <p>The returned future completes when the topology has been persisted or
* forwarded. All failures must be signalled by completing the returned future
* exceptionally — implementations must not throw synchronously from this method.
* The broker handles the future's completion exception as follows:
* {@link org.apache.kafka.common.errors.InvalidRequestException} maps to
* {@code INVALID_REQUEST}; {@link org.apache.kafka.common.errors.TopologyDescriptionTooLargeException}
* maps to {@code TOPOLOGY_DESCRIPTION_TOO_LARGE}; any other exception maps to
* {@code TOPOLOGY_DESCRIPTION_UPDATE_FAILED} and is logged at WARN.
*
* @param requestContext the context of the UpdateStreamsGroupTopologyDescription request
* @param groupId the streams group ID
* @param groupCreationTimeMs the timestamp when the group was created
* @param topologyEpoch the topology epoch
* @param description the topology description
* @return a future that completes when the operation is done
*/
CompletableFuture<Void> setTopology(AuthorizableRequestContext requestContext,
String groupId, long groupCreationTimeMs, int topologyEpoch,
StreamsGroupTopologyDescription description);
/**
* Called when a group is explicitly deleted via DeleteGroups. Removes any topology
* description stored for this group.
*
* <p>The returned future completes when the deletion has been processed.
* If it completes exceptionally, the broker logs the error; the outcome does not
* affect the DeleteGroups response returned to the caller.
*
* @param requestContext the context of the DeleteGroups request
* @param groupId the streams group ID
* @return a future that completes when the operation is done
*/
CompletableFuture<Void> deleteTopology(AuthorizableRequestContext requestContext, String groupId);
/**
* Called to retrieve the stored topology description for a group. This is invoked
* by the broker when a client calls StreamsGroupDescribe with
* {@code IncludeTopologyDescription=true}.
*
* <p>Returns a future that resolves to the stored topology description for the
* given {@code (groupId, groupCreationTimeMs, topologyEpoch)} tuple, or to
* {@code null} if no topology is stored (e.g. no push has succeeded yet, or the
* stored description is for a different epoch). If the future completes
* exceptionally, the plugin signals a read error for this group.
*
* @param requestContext the context of the StreamsGroupDescribe request
* @param groupId the streams group ID
* @param groupCreationTimeMs the timestamp when the group was created
* @param topologyEpoch the topology epoch the caller is asking about
* @return a future resolving to the stored topology description, or null if none
*/
CompletableFuture<StreamsGroupTopologyDescription>
getTopology(AuthorizableRequestContext requestContext,
String groupId, long groupCreationTimeMs, int topologyEpoch);
} |
The plugin uses a StreamsGroupTopologyDescription POJO that mirrors org.apache.kafka.streams.TopologyDescription but lives in the org.apache.kafka.coordinator.group.api.streams package; plugin implementations only need to depend on group-coordinator-api. The only difference is that there is no predecessor relation.
package org.apache.kafka.coordinator.group.api.streams;
public class StreamsGroupTopologyDescription {
public Collection<Subtopology> subtopologies();
public Collection<GlobalStore> globalStores();
public static class Subtopology {
public String id();
public Collection<Node> nodes();
}
/**
* A processing node in the topology. Predecessor nodes can be inferred from successor relation.
*/
public interface Node {
String name();
Set<String> successors();
}
public static class Source implements Node {
public Set<String> topics();
}
public static class Processor implements Node {
public Set<String> stores();
}
public static class Sink implements Node {
public Optional<String> topic();
}
public static class GlobalStore {
public Source source();
public Processor processor();
}
} |
DescribeStreamsGroupsOptions gains an includeTopologyDescription(boolean) setter. When set to true, the admin client sets IncludeTopologyDescription on the StreamsGroupDescribeRequest (at the new version) and the coordinator consults the plugin.
StreamsGroupDescription gains two accessors:
Optional<StreamsGroupTopologyDescription> topologyDescription() — empty unless the field was requested and the plugin returned a description.StreamsGroupTopologyDescriptionStatus topologyDescriptionStatus() — a new enum { AVAILABLE, NOT_REQUESTED, NOT_STORED, ERROR }. AVAILABLE is reported when TopologyDescription is non-null; the remaining values mirror the wire-level TopologyDescriptionStatus int8.The Admin client exposes a POJO hierarchy in org.apache.kafka.clients.admin that mirrors org.apache.kafka.streams.TopologyDescription but lives in the clients module. The admin client converts the wire-format struct into this hierarchy.
package org.apache.kafka.clients.admin;
public class StreamsGroupTopologyDescription {
public Collection<Subtopology> subtopologies();
public Collection<GlobalStore> globalStores();
public static class Subtopology {
public String id();
public Collection<Node> nodes();
}
public interface Node {
String name();
/** Direct predecessor nodes. */
Set<String> predecessors();
/** Direct successor nodes. */
Set<String> successors();
}
public static class Source implements Node {
public Set<String> topics();
}
public static class Processor implements Node {
public Set<String> stores();
}
public static class Sink implements Node {
public Optional<String> topic();
}
public static class GlobalStore {
public Source source();
public Processor processor();
}
} |
kafka-streams-groups.sh gains a new --topology sub-action under --describe, parallel to --members, --offsets, and --state:
kafka-streams-groups.sh --bootstrap-server <broker> --describe --topology --group <group-id>
The output mirrors the format produced by Topology#describe() in the Kafka Streams API:
Topologies:
Sub-topology: 0
Source: KSTREAM-SOURCE-0000000000 (topics: [input-topic])
--> my-processor
Processor: my-processor (stores: [my-store])
<-- KSTREAM-SOURCE-0000000000
--> KSTREAM-SINK-0000000002
Sink: KSTREAM-SINK-0000000002 (topic: output-topic)
<-- my-processor
Global stores, if present, are printed in a trailing Global Stores: block. The command calls the admin client with includeTopologyDescription(true).
When the coordinator returns no topology description, the command picks its output from the TopologyDescriptionStatus on the response:
NOT_STORED → "No topology description has been recorded for group '<group-id>'."ERROR → "The broker failed to retrieve the topology description for group '<group-id>' (check broker logs)."AVAILABLE → normal pretty-printed topology.Exit code. The command exits 0 for AVAILABLE. It exits 1 for NOT_STORED, for ERROR, and when the DescribedGroup.ErrorCode is non-zero (authorization failure, group not found, coordinator unavailable, etc.). NOT_REQUESTED does not occur for this command because it always sets IncludeTopologyDescription=true.
The plugin is instantiated at broker startup if group.streams.topology.description.plugin.class is configured. A broker without a plugin does not advertise UpdateStreamsGroupTopologyDescription in its ApiVersions response and returns UNSUPPORTED_VERSION for any such request. Because a plugin-less broker never sets TopologyDescriptionRequired in heartbeat responses, in normal operation the RPC is only sent against a broker that has a plugin configured.
After a successful StreamsGroupHeartbeat, the broker calls requiresTopologyPush on the plugin. STALE_TOPOLOGY members are skipped. If the plugin returns true, the broker sets TopologyDescriptionRequired=true on the heartbeat response. Per the plugin contract requiresTopologyPush should not throw; the broker defensively catches any exception, logs at WARN, and treats the call as false. Beyond the STALE_TOPOLOGY skip the broker does no other gating — deduplication, timeouts, and back-off are the plugin's responsibility (see Plugin Implementation Guidelines). The groupCreationTimeMs parameter identifies the specific incarnation of a group; a recreated group with the same groupId gets a fresh creation timestamp. It is persisted on the group record as a tagged field shared with KIP-1282. Groups created before KIP-1282 is implemented carry the sentinel value groupCreationTimeMs = 0.
UpdateStreamsGroupTopologyDescription, the broker checks the READ ACL on the group and that a plugin is configured. The broker does not enforce a size limit on the topology description — it is the plugin's responsibility to decide what size it is willing to store. The broker calls setTopology on the plugin. On success, the response carries NONE; InvalidRequestException from the plugin maps to INVALID_REQUEST, TopologyDescriptionTooLargeException maps to TOPOLOGY_DESCRIPTION_TOO_LARGE, any other exception maps to TOPOLOGY_DESCRIPTION_UPDATE_FAILED, and all three are logged at WARN.DeleteGroups, the broker calls deleteTopology on the plugin only after the group delete has returned no error for that group. deleteTopology failures are logged but do not affect the deletion response. For groups that expire naturally (all members leave), deleteTopology is not called — the plugin expires the topology via a wall-clock TTL (see Plugin Implementation Guidelines).On StreamsGroupDescribe with IncludeTopologyDescription=true, the broker calls getTopology on the plugin for each group after assembling the rest of the response, passing the group's current (groupCreationTimeMs, topologyEpoch) tuple. Calls across groups run in parallel. Authorization is unchanged — the existing DESCRIBE ACL on the GROUP resource covers the topology description. DescribedGroup.ErrorCode is never modified by topology-related outcomes; the TopologyDescriptionStatus field carries the reason when TopologyDescription is null (NOT_REQUESTED, NOT_STORED, or ERROR; the last is also logged at WARN). When no plugin is configured, the broker returns NOT_STORED (the broker does not distinguish "no plugin" from "plugin has nothing stored" on the wire).
The Streams client records the TopologyDescriptionRequired flag from each heartbeat response.
topology.description.push.enabled=true, the Streams client converts the topology returned by Topology#describe() to the wire format and stores it internally. When topology.description.push.enabled=false, no description is stored and the feature is disabled on this client.On each consumer background-thread poll, the client sends UpdateStreamsGroupTopologyDescription to the coordinator when a coordinator is known, the flag is set, a stored topology description is available, and no prior request is in flight. Completion handling is described in Error Handling and Retries.
The push runs on the consumer background thread and never blocks user-facing Kafka Streams APIs. The push is best-effort: a failure to send the topology description does not prevent the Streams application from running.
TopologyDescription from Kafka Streams maps one-to-one to the wire types in the RPC schema. Predecessor edges are not sent on the wire; the read side reconstructs them by inverting each node's successor list.
NOT_COORDINATOR and COORDINATOR_NOT_AVAILABLE trigger coordinator rediscovery; the flag stays set. COORDINATOR_LOAD_IN_PROGRESS and network exceptions leave the flag set for retry on the next poll.
All other errors (TOPOLOGY_DESCRIPTION_TOO_LARGE, TOPOLOGY_DESCRIPTION_UPDATE_FAILED, INVALID_REQUEST, UNSUPPORTED_VERSION, GROUP_ID_NOT_FOUND, GROUP_AUTHORIZATION_FAILED) clear the topologyDescriptionRequired flag and log at WARN. The client does not retry on its own; re-attempts, if any, happen when the broker re-sets the flag via a subsequent heartbeat. See Plugin Implementation Guidelines for when the plugin should re-solicit.
A non-zero ThrottleTimeMs on the response delays the next push attempt by that amount, as with other request managers.
Broker-side gating around the plugin is minimal. The plugin can expect requiresTopologyPush to be called on a regular cadence following each group's heartbeat interval.
A correct plugin implementation should:
requiresTopologyPush as a non-blocking, efficient call, since it is invoked frequently.requiresTopologyPush calls can be used as an implicit keep-alive.setTopology future with TopologyDescriptionTooLargeException, and stop returning requiresTopologyPush=true for the same (groupId, topologyEpoch) pair once that pair has been confirmed too large.(groupId, groupCreationTimeMs, topologyEpoch) tuple the plugin should track (i) whether a push is currently in flight and (ii) the time of the last requiresTopologyPush=true. While a push is in flight requiresTopologyPush returns false. Once the push has completed successfully it returns false permanently for that tuple. On a transient setTopology failure the plugin arranges for requiresTopologyPush to return true again after a self-driven back-off (e.g. exponential, starting at 1s) — this is the client re-solicitation mechanism, and it is the only retry pathway: the wire-level TOPOLOGY_DESCRIPTION_UPDATE_FAILED is terminal from the client's perspective and pairs with this plugin-side re-solicitation. On permanent failure (TOPOLOGY_DESCRIPTION_TOO_LARGE or plugin-semantic INVALID_REQUEST) the plugin returns false permanently for the tuple and logs.The topology description contains user-defined processor names, state-store names, and topic names. This KIP does not treat these as inherently sensitive and does not introduce any client-side redaction. Operators who consider node, store, or topic names sensitive should scope DESCRIBE ACLs on the GROUP resource accordingly: the same ACL that guards StreamsGroupDescribe guards the TopologyDescription field on the response. The UpdateStreamsGroupTopologyDescription path is guarded by READ on the GROUP, identical to the existing heartbeat ACL.
This KIP does not introduce any new broker-side or client-side metrics. It is the responsibility of the plugin to define metrics as necessary.
This KIP bumps StreamsGroupHeartbeatResponse, StreamsGroupDescribeRequest, and StreamsGroupDescribeResponse to the next available version of each RPC and adds the new fields (TopologyDescriptionRequired, IncludeTopologyDescription, TopologyDescription, TopologyDescriptionStatus) at that version. Pre-upgrade clients negotiate an older version and never see the new fields, so all three changes are wire-compatible.
The new RPC uses a new API key and is only sent by clients that understand the feature.
Without a configured plugin, no flags are set and no topology descriptions are sent. There is no behavioral change for existing deployments.
During a rolling upgrade of brokers, some brokers may have the plugin configured and some may not. The TopologyDescriptionRequired flag is only set by plugin-equipped coordinators; a client whose current coordinator lacks the plugin never sees the flag. If the coordinator migrates mid-push to a plugin-less broker, the client receives UNSUPPORTED_VERSION, clears the topologyDescriptionRequired flag, and does not retry. The flag is set again only if a future heartbeat response from a plugin-equipped coordinator includes TopologyDescriptionRequired=true.
During a rolling upgrade of the Streams application (topology epoch change), the broker skips the plugin entirely for STALE_TOPOLOGY members — neither requiresTopologyPush is called nor is the response flag set. Only members running the new topology reach the plugin. If every active member is stale during the rollout, the current-epoch topology remains uncaptured until at least one member heartbeats with the new epoch.
During this transition the assignment topology (advanced synchronously when a new-epoch member heartbeats) and the description topology (advanced asynchronously via the plugin) can briefly disagree: StreamsGroupDescribe may report the new topologyEpoch while TopologyDescription is still null with status NOT_STORED until the first push for the new epoch succeeds. The two reconverge once any member at the new epoch pushes its description.
Hash-based mismatch detection. A future enhancement could introduce a topology hash to detect clients on different topology descriptions reporting the same topology epoch.
Multi-version describe. The describe response surfaces only the topology under the current (groupCreationTimeMs, topologyEpoch) tuple. During a rolling topology upgrade, the previous epoch's description may still stored by the plugin but is no longer reachable via describe. A natural extension is a descriptions[] array on the response, tagged by epoch, allowing operators to view both the previous and the in-flight new topology while the rollout completes.
Node grouping. A nodeGroup field on TopologyNode for UI grouping, requiring an extension to the public TopologyDescription.Node interface in Kafka Streams — best addressed in a follow-up KIP.
setTopology future with TopologyDescriptionTooLargeException causes the broker to return TOPOLOGY_DESCRIPTION_TOO_LARGE, and the client clears the topologyDescriptionRequired flag without retrying.Disabling the feature on the client (topology.description.push.enabled=false) stops it from sending topology descriptions altogether; describe returns NOT_STORED.
NOT_STORED on describe and never asks the client to push (the broker-side half of the rolling-upgrade matrix; the client-side half is intrinsic to wire-version negotiation).Adding the topology description to StreamsGroupHeartbeatRequest would bloat the heartbeat with potentially large payloads and mix observability concerns with the rebalance protocol. A separate RPC keeps the feature self-contained and allows the broker to apply independent size limits and authorization. See the discussion in the KIP for a full comparison of both approaches.
Using a compacted internal topic (like __consumer_offsets) would require the broker to materialize all topology descriptions in memory from the topic's compacted log. For clusters with many streams groups, this could consume significant heap space on the coordinator broker. A plugin-based approach avoids this by delegating storage to an external system that can handle retention and retrieval independently.
A dedicated read RPC was considered but rejected in favor of extending StreamsGroupDescribe at a new version with an IncludeTopologyDescription flag. This keeps the number of new API keys minimal, groups the topology description naturally with other group metadata, and avoids bloating the describe response when the caller does not need it. The plugin remains free to expose topologies directly via its own channels in addition to the describe RPC.
KIP-714 supports compression (ZStd, LZ4, GZip, Snappy) for telemetry payloads because metrics are pushed repeatedly at high frequency. Topology descriptions are pushed infrequently (only on topology epoch changes) and the wire payload is expected to fit comfortably for typical applications. The Kafka protocol's flexible versions already provide efficient serialization. Adding compression would add complexity without meaningful benefit for this use case.
For illustration, consider a Kafka Streams application that reads from orders, filters out null values, and writes the remaining records to valid-orders:
StreamsBuilder builder = new StreamsBuilder();
builder.stream("orders")
.filter((key, value) -> value != null)
.to("valid-orders");
Topology#describe() produces the following text representation:
Topologies:
Sub-topology: 0
Source: KSTREAM-SOURCE-0000000000 (topics: [orders])
--> KSTREAM-FILTER-0000000001
Processor: KSTREAM-FILTER-0000000001 (stores: [])
--> KSTREAM-SINK-0000000002
<-- KSTREAM-SOURCE-0000000000
Sink: KSTREAM-SINK-0000000002 (topic: valid-orders)
<-- KSTREAM-FILTER-0000000001
The corresponding UpdateStreamsGroupTopologyDescription request body, with the topology converted to the wire format, looks like this:
{
"GroupId": "orders-app",
"TopologyEpoch": 0,
"TopologyDescription": {
"Subtopologies": [
{
"SubtopologyId": "0",
"Nodes": [
{
"Name": "KSTREAM-SOURCE-0000000000",
"NodeType": 1,
"SourceTopics": ["orders"],
"SinkTopic": null,
"Stores": [],
"Successors": ["KSTREAM-FILTER-0000000001"]
},
{
"Name": "KSTREAM-FILTER-0000000001",
"NodeType": 2,
"SourceTopics": [],
"SinkTopic": null,
"Stores": [],
"Successors": ["KSTREAM-SINK-0000000002"]
},
{
"Name": "KSTREAM-SINK-0000000002",
"NodeType": 3,
"SourceTopics": [],
"SinkTopic": "valid-orders",
"Stores": [],
"Successors": []
}
]
}
],
"GlobalStores": []
}
}
Predecessor edges (<-- in the text representation) are not sent on the wire; the read side reconstructs them by inverting each node's Successors.