Versions Compared

Key

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

...

  • GROUP_AUTHORIZATION_FAILED — the client is not authorized
  • INVALID_REQUEST — the request is malformed (including an empty MemberId), or the plugin semantically rejected the payload by completing its future with InvalidRequestException
  • UNSUPPORTED_VERSION — the coordinator cannot serve this RPC because no topology description plugin is configured
  • STREAMS_TOPOLOGY_DESCRIPTION_TOO_LARGE — the plugin rejected the description because it exceeds the size the plugin is willing to storeSTREAMS_TOPOLOGY_DESCRIPTION_UPDATE_FAILED — the plugin failed to process the request for some other reason. The accompanying ErrorMessage carries the plugin's exception message. The broker's response shape is identical for both permanent and transient plugin failures; the distinction is broker-internal state that determines whether subsequent heartbeats at the same topology epoch will re-solicit (see Broker Side).
  • UNKNOWN_MEMBER_ID — the member named in MemberId is no longer in the group; the client should treat itself as fenced and rejoin

  • GROUP_ID_NOT_FOUND — the specified group does not exist
  • NOT_COORDINATOR — the broker is not the coordinator for this group
  • COORDINATOR_NOT_AVAILABLE — the coordinator is not available
  • COORDINATOR_LOAD_IN_PROGRESS — the coordinator is loading

...

DeleteGroupsResponse Change

No schema change. A new error code is added to the per-group ErrorCode slot of the existing DeleteGroupsResponse:

  • STREAMS_TOPOLOGY_DESCRIPTION_DELETE_FAILED — the topology description plugin failed to delete the description for this streams group; the group is not tombstoned. The caller may retry the request once the plugin recovers, or unset group.streams.topology.description.plugin.class to bypass the plugin. The plugin's exception is logged at WARN on the broker; the response itself carries only the error code, matching the existing DeleteGroupsResponse shape.

Plugin Interface

A new interface is introduced in org.apache.kafka.coordinator.group.api.streams DeleteGroupsRequest and DeleteGroupsResponse are both bumped to the next version (3). The request shape is unchanged at the new version; the response adds an ErrorMessage field to each per-group DeletableGroupResult:

Code Block
languagejavajs
linenumberstrue
/**
 * A broker-side plugin that stores, forwards, or exposes topology descriptions pushed
 * by Kafka Streams clients.
 *
 * <p>Implementations must be thread-safe. {@link #setTopology} may be called
 * concurrently by multiple members of the same group; calls with the same
 * {@code (groupId, topologyEpoch)} carry identical data and must be idempotent.
 * {@link #deleteTopology} must also be idempotent — it may be called more than once
 * for the same {@code groupId}, including when nothing is stored.
 */
public interface StreamsGroupTopologyDescriptionPlugin extends Configurable, AutoCloseable {

    /**
     * Store the topology description for a streams group.
     *
     * <p>The returned future completes when the topology has been persisted or forwarded.
     * Failures must be signalled by completing the future exceptionally — implementations
     * must not throw synchronously. The completion exception maps to the client-visible
     * error code:
     *
     * <ul>
     *   <li>{@link org.apache.kafka.common.errors.InvalidRequestException} — payloads the
     *       plugin will not accept on semantic grounds; reported as {@code INVALID_REQUEST}.</li>
     *   <li>{@link org.apache.kafka.common.errors.StreamsTopologyDescriptionTooLargeException} —
     *       descriptions larger than the plugin is willing to store; reported as
     *       {@code STREAMS_TOPOLOGY_DESCRIPTION_TOO_LARGE}.</li>
     *   <li>Any other exception — transient backend failure; reported as
     *       {@code STREAMS_TOPOLOGY_DESCRIPTION_UPDATE_FAILED}.</li>
     * </ul>
     *
     * The first two are treated as permanent at this topology epoch and no further push
     * will be solicited until the epoch advances. The third is treated as transient and
     * may be retried.
     */
    CompletableFuture<Void> setTopology(String groupId, int topologyEpoch,
                                        StreamsGroupTopologyDescription description);

	/**
	 * Remove any topology description stored for this group. Called when the group is
	 * deleted or expires. A failure (future completed exceptionally) is reported to the
	 * caller of {@code DeleteGroups} as {@code STREAMS_TOPOLOGY_DESCRIPTION_DELETE_FAILED} on that
	 * group's per-group result and the broker does not tombstone the group; a retry of
	 * {@code DeleteGroups} re-invokes this method idempotently. The periodic-cleanup path
	 * treats a failure identically — the group's tombstone is deferred to a future cycle.
	 */
	CompletableFuture<Void> deleteTopology(String groupId);

    /**
     * Return the stored topology description for {@code (groupId, topologyEpoch)}, or
     * {@code null} if the plugin no longer has the data (e.g. backend wipe). If the future
     * completes exceptionally, the broker reports a read error for the group.
     */
    CompletableFuture<StreamsGroupTopologyDescription> getTopology(String groupId, 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.

{ "name": "ErrorMessage", "type": "string", "versions": "3+", "nullableVersions": "3+", "default": "null",
  "about": "A description of the deletion error, or null if the deletion succeeded or no description is available." }

A new generic error code is added to the per-group ErrorCode slot:

  • DELETE_FAILED — the delete operation could not complete; the accompanying ErrorMessage describes the underlying cause. The group is not tombstoned, and the caller may retry once the underlying condition is resolved. For streams groups configured with a topology description plugin this is returned when plugin.deleteTopology fails; other group types may adopt the same code in the future.

This is the only deletion-blocking failure mode introduced by this KIP. Consumer and share groups are unaffected. See Broker Side for the full ordering rule and the recovery path.

Plugin Interface

A new interface is introduced in org.apache.kafka.coordinator.group.api.streams :

Code Block
languagejava
linenumberstrue
/**
 * A broker-side plugin that stores, forwards, or exposes topology descriptions pushed
 * by Kafka Streams clients.
 *
 * <p>Implementations must be thread-safe. {@link #setTopology} may be called
 * concurrently by multiple members of the same group; calls with the same
 * {@code (groupId, topologyEpoch)} carry identical data and must be idempotent.
 * {@link #deleteTopology} must also be idempotent — it may be called more than once
 * for the same {@code groupId}, including when nothing is stored.
 */
public interface StreamsGroupTopologyDescriptionPlugin extends Configurable, AutoCloseable {

	/**
	 * Store the topology description for a streams group.
	 *
	 * <p>The returned future completes when the topology has been persisted or forwarded.
	 * Failures must be signalled by completing the future exceptionally — implementations
	 * must not throw synchronously. The completion exception drives broker-side behaviour:
	 *
	 * <ul>
	 *   <li>{@link PluginPermanentFailureException} — the description will never be accepted
	 *       at this topology epoch (e.g. too large, semantically rejected). The broker
	 *       ratchets {@code LastFailedTopologyEpoch} and stops re-soliciting until the
	 *       epoch advances.</li>
	 *   <li>{@link PluginTransientFailureException} or any other exception — treated as
	 *       transient. The broker arms or extends the per-group back-off (30 s → 1 h,
	 *       exponential) and re-solicits on a later heartbeat.</li>
	 * </ul>
	 *
	 * In both cases the caller receives error code
	 * {@code STREAMS_TOPOLOGY_DESCRIPTION_UPDATE_FAILED} with the exception's message in
	 * {@code ErrorMessage}; the permanent-vs-transient split is broker-internal state.
	 */
	CompletableFuture<Void> setTopology(String groupId, int topologyEpoch,
                                    	StreamsGroupTopologyDescription description);

	/**
	 * Remove any topology description stored for this group. Called when the group is
	 * deleted or expires. A failure (future completed exceptionally) is reported to the
	 * caller of {@code DeleteGroups} as {@code DELETE_FAILED} with the exception message
	 * in the per-group {@code ErrorMessage}, and the broker does not tombstone the group;
 	* a retry of {@code DeleteGroups} re-invokes this method idempotently. The
	 * periodic-cleanup path treats a failure identically — the group's tombstone is
	 * deferred to a future cycle.
	 */
	CompletableFuture<Void> deleteTopology(String groupId);

    /**
     * Return the stored topology description for {@code (groupId, topologyEpoch)}, or
     * {@code null} if the plugin no longer has the data (e.g. backend wipe). If the future
     * completes exceptionally, the broker reports a read error for the group.
     */
    CompletableFuture<StreamsGroupTopologyDescription> getTopology(String groupId, 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.

Code Block
languagejava
linenumberstrue
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();
    }
Code Block
languagejava
linenumberstrue
package org.apache.kafka.coordinator.group.api.streams;

public class StreamsGroupTopologyDescription {
    public Collection<Subtopology> subtopologies();
    public Collection<GlobalStore> globalStores();

    public static class SubtopologyProcessor implements Node {
        public StringSet<String> idstores();
    }

    public Collection<Node> nodesstatic class Sink implements Node {
        public Optional<String> topic();
      }

	/**
     * 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 static class GlobalStore {
        public Source source();
        public Processor processor();
    }
}

Two new exception classes in the same package let the plugin signal the permanent-vs-transient distinction. Plugins that throw any other exception are treated as transient:

Code Block
languagejava
linenumberstrue
package org.apache.kafka.coordinator.group.api.streams;

import org.apache.kafka.common.errors.ApiException;

/** Signals that the topology description for the current epoch will never be accepted (e.g. too large, semantically rejected). */
public class PluginPermanentFailureException extends ApiException {
    public PluginPermanentFailureException(String message) { public Set<String> topicssuper(message);
    }

    public static class Processor implements Node {
        public Set<String> stores();
    }

    public static class Sink implements Node {
        public Optional<String> topic();
    }

    public static class GlobalStorePluginPermanentFailureException(String message, Throwable cause) { super(message, cause); }
}

/** Signals a transient backend failure; the broker re-solicits on a later heartbeat. Plugins that throw any other exception are treated identically. */
public class PluginTransientFailureException extends ApiException {
    public PluginTransientFailureException(String message) { public Source sourcesuper(message); }
    public PluginTransientFailureException(String message, Throwable cause) public Processor processor({ super(message, cause);
    }
}

Broker-Side Persistence

...

  1. The plugin is instantiated at broker startup if group.streams.topology.description.plugin.class is configured. A broker without a plugin returns UNSUPPORTED_VERSION for StreamsGroupTopologyDescriptionUpdate and never sets TopologyDescriptionRequired, so the RPC is only sent against a plugin-configured broker.

  2. After a successful StreamsGroupHeartbeat, the broker decides whether to set TopologyDescriptionRequired=true purely from the group's persisted state — no plugin RPC is involved. Members with STALE_TOPOLOGY status are skipped. For all other members the broker sets the flag iff StoredTopologyEpoch != currentTopologyEpoch AND LastFailedTopologyEpoch != currentTopologyEpoch AND no per-group back-off is in its window. The back-off is in-memory state (keyed by groupId, carrying topologyEpoch + nextAttemptMs) that arms or extends every time the flag is set and additionally on a transient setTopology failure; consecutive arms double the window from 30 s up to 1 h. It clears on a successful push, on a permanent failure (where LastFailedTopologyEpoch ratchets), and implicitly on any topology-epoch advance. The same mechanism covers unresponsive plugins and clients that never push (for example, with topology.description.push.enabled=false).

  3. On On StreamsGroupTopologyDescriptionUpdate, the broker checks the READ ACL on the group and that a plugin is configured, then validates the MemberId: an empty MemberId is rejected with INVALID_REQUEST, a non-existing streams group with GROUP_ID_NOT_FOUND, and a MemberId not matching any current member with UNKNOWN_MEMBER_ID. The broker enforces no size limit; the plugin decides what it is willing to store. The broker then calls setTopology on the plugin. On success it writes a metadata record setting StoredTopologyEpoch = pushedEpoch and the response carries NONE. On InvalidRequestException or StreamsTopologyDescriptionTooLargeException PluginPermanentFailureException it writes LastFailedTopologyEpoch = pushedEpoch so subsequent heartbeats at the same epoch do not re-solicit. Any other exception maps to On PluginTransientFailureException or any other exception it writes no metadata record, arms the per-group back-off, and the next heartbeat re-solicits once the window elapses. In both failure cases the response carries STREAMS_TOPOLOGY_DESCRIPTION_UPDATE_FAILED, is logged at WARN, and is treated as transient — no metadata record is written and the next heartbeat re-solicits with the plugin's exception message in ErrorMessage; the permanent-vs-transient split is broker-internal state. If the plugin call succeeds but the metadata-record write fails, the next heartbeat sees the drift, re-solicits, and the idempotent re-push closes the gap.

  4. On DeleteGroups, the broker calls deleteTopology on the plugin before writing the group tombstone, for each requested streams group with StoredTopologyEpoch != -1. On plugin success the group is tombstoned and the per-group ErrorCode is NONE. On plugin failure the group is not tombstoned and the per-group ErrorCode is set to STREAMS_TOPOLOGY_DESCRIPTION_DELETE_FAILED ( with the plugin's exception is message in ErrorMessage (also logged at WARN on the broker); the operator retries the request once the plugin recovers, or unsets the plugin config to bypass it. Other groups in the same batch are unaffected — the failure is reported per group. This ordering matches the natural-expiration cleanup below, which also defers tombstoning until plugin.deleteTopology succeeds.

  5. On StreamsGroupDescribe with IncludeTopologyDescription=true, the broker calls getTopology on the plugin only when StoredTopologyEpoch == currentTopologyEpoch for that group; otherwise it reports NOT_STORED without making a plugin call. 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). With no plugin configured the broker returns NOT_STORED. A getTopology call returning null (plugin-side data loss) surfaces as NOT_STORED and is logged at WARN; subsequent describes keep returning NOT_STORED until the topology epoch advances or an operator clears plugin state.
  6. When a plugin is configured, the broker runs a periodic topology-description cleanup every offsets.retention.check.interval.ms. Each fire fans out a read-only query across the broker's hosted __consumer_offsets partitions to identify streams groups eligible for cleanup: isEmpty && allOffsetsExpired && StoredTopologyEpoch != -1. This is the same eligibility predicate the shard's offset-expiration sweep uses to delete consumer/share groups, with the additional StoredTopologyEpoch != -1 filter. For each eligible group, the broker calls plugin.deleteTopology(groupId) and, on success, writes a metadata record setting StoredTopologyEpoch = -1. Plugin failures leave the field set; the same group is retried on the next cycle. Once StoredTopologyEpoch = -1, the shard's offset-expiration sweep tombstones the (now flag-cleared) group on a subsequent cycle. When no plugin is configured, the periodic cleanup does not run and the shard's offset-expiration sweep expires streams groups normally, ignoring StoredTopologyEpoch; operators that disable a previously-configured plugin are responsible for cleaning up plugin-side state out-of-band.

...

  1. At startup, if topology.description.push.enabled=true, the Streams client converts the topology returned by Topology#describe() to the wire format and stores it internally. The mapping is one-to-one with the RPC schema; predecessor edges are not sent on the wire and the read side reconstructs them by inverting each node's successor list. When topology.description.push.enabled=false, no description is stored and the feature is disabled on this client.
  2. The Streams client records the TopologyDescriptionRequired flag from each heartbeat response.
  3. On each consumer background-thread poll, the client sends StreamsGroupTopologyDescriptionUpdate to the coordinator when a coordinator is known, the flag is set, a stored topology description is available, the client has a non-empty member ID assigned by the coordinator, and no prior request is in flight. The member ID populated on the request is the same one carried on StreamsGroupHeartbeat. The push runs on the consumer background thread and never blocks user-facing Kafka Streams APIs; the push is best-effort.
  4. Completion handling on the push response is keyed on the error code. NOT_COORDINATOR and COORDINATOR_NOT_AVAILABLE trigger coordinator rediscovery and leave the flag set. COORDINATOR_LOAD_IN_PROGRESS and network exceptions leave the flag set for retry on the next poll. UNKNOWN_MEMBER_ID means the

    broker no longer recognizes this member (group deleted, or member dropped)

    member has been dropped from the group: the client clears the flag and relies on the existing membership-management path to trigger a clean rejoin on the next heartbeat. All other errors (STREAMS_TOPOLOGY_DESCRIPTION

    _TOO

    _

    LARGE, STREAMS_TOPOLOGY_DESCRIPTION_

    UPDATE_FAILED, INVALID_REQUEST, UNSUPPORTED_VERSION, GROUP_ID_NOT_FOUND, GROUP_AUTHORIZATION_FAILED) clear the flag and log at WARN with the response's ErrorMessage; the client does not retry on its own, and a re-attempt happens only when the broker re-sets the flag via a subsequent heartbeat. A non-zero ThrottleTimeMs on the response delays the next push attempt by that amount, as with other request managers.

Plugin Implementation Guidelines

...

  • Treat setTopology (on (groupId, topologyEpoch)) and deleteTopology (on (groupId)) as idempotent; the broker may re-issue an identical call when an earlier call's bookkeeping write failed.
  • Be thread-safe under concurrent invocation: setTopology may be called by multiple members of the same group in the same heartbeat cycle, and the periodic-cleanup path may invoke deleteTopology while a member is mid-push.
  • Reject payloads the plugin will not accept by completing the the setTopology future with StreamsTopologyDescriptionTooLargeException or InvalidRequestException PluginPermanentFailureException; the broker persists the rejection at the epoch level via LastFailedTopologyEpoch and stops re-soliciting at the same epoch. The exception message reaches the client in ErrorMessage.
  • Signal transient storage-layer failures by completing the setTopology future with PluginTransientFailureException (or any other exception); the broker's per-group back-off (30 s → 1 h, exponential) throttles re-solicitation.
  • Return Return null from getTopology when the plugin has lost the description; the broker reports NOT_STORED on the describe response. Exceptions are reserved for transient backend failures and surface as ERROR.
  • Avoid blocking coordinator threads (plugin methods may be invoked on them) and complete futures within seconds, not minutes; the broker applies no wall-clock deadline to plugin calls, so the coordinator's responsiveness is bounded by what the plugin does. Bound plugin-side state explicitly — the plugin shares the broker heap.

...

MBeanTypeDescription
kafka.server:type=group-coordinator-metrics,name=streams-group-topology-description-set-success-{rate,count}MeterSuccessful plugin.setTopology calls.
kafka.server:type=group-coordinator-metrics,name=streams-group-topology-description-set-error-{rate,count}Meter

Failed plugin.setTopology calls. An error increments this sensor regardless of whether it was

StreamsTopologyDescriptionTooLargeException

PluginPermanentFailureException,

InvalidRequestException

PluginTransientFailureException, or any other exception.

kafka.server:type=group-coordinator-metrics,name=streams-group-topology-description-delete-success-{rate,count}MeterSuccessful plugin.deleteTopology calls (covers both the explicit DeleteGroups and periodic-cleanup paths).
kafka.server:type=group-coordinator-metrics,name=streams-group-topology-description-delete-error-{rate,count}MeterFailed plugin.deleteTopology calls.
kafka.server:type=group-coordinator-metrics,name=streams-group-topology-description-get-success-{rate,count}MeterSuccessful plugin.getTopology calls.
kafka.server:type=group-coordinator-metrics,name=streams-group-topology-description-get-error-error-{rate,count}MeterFailed plugin.getTopology calls.
kafka.server:type=group-coordinator-metrics,name=streams-group-topology-description-cleanup-cycle-{rate,count}MeterFailed plugin.getTopology calls.Periodic topology-description cleanup cycles that actually ran.
kafka.server:type=group-coordinator-metrics,name=streams-group-topology-description-cleanup-cycleeligible-{rate,count}MeterPeriodic topology-description cleanup cycles that actually ran.
kafka.server:type=group-coordinator-metrics,name=streams-group-topology-description-cleanup-eligible-{rate,count}MeterStreams group IDs identified as eligible for topology-description cleanup, summed across partitions.

No client-side metrics are introduced.

Compatibility, Deprecation, and Migration Plan

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.

Streams group IDs identified as eligible for topology-description cleanup, summed across partitions.


No client-side metrics are introduced.

Compatibility, Deprecation, and Migration Plan

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.

DeleteGroupsRequest and DeleteGroupsResponse are bumped to version 3, which adds a per-group ErrorMessage field on the response and introduces the new DELETE_FAILED error code. Older admin clients negotiate version 2 and never see the new field; they continue to receive only the per-group ErrorCode. DeleteGroups can newly fail on streams groups when a topology description plugin is configured and its deleteTopology call fails. Brokers without a configured plugin keep today's behaviour. Older clients that do receive DELETE_FAILED (because they understand version 3 but predate this KIP's error-code addition) decode the code as UNKNOWN_SERVER_ERROR via the standard forward-compatibility fallback in Errors.forCode; the group is still not tombstoned, so an idempotent retry of DeleteGroups converges once the plugin recovers. No client-side change is requiredWithout a configured plugin, no flags are set and no topology descriptions are sent. There is no behavioral change for existing deployments.

Rolling Upgrades

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 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.DeleteGroups can newly fail on streams groups when a topology description plugin is configured and its deleteTopology call fails. Brokers without a configured plugin keep today's behaviour. Older admin clients receive the new STREAMS_TOPOLOGY_DESCRIPTION_DELETE_FAILED code (136) as UNKNOWN_SERVER_ERROR via the standard forward-compatibility fallback in Errors.forCode; the group is still not tombstoned, so an idempotent retry of DeleteGroups converges once the plugin recovers. No client-side change is required.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.

Future Work

Hash-based mismatch detection. A future enhancement could introduce a topology hash to detect clients on different topology descriptions reporting the same topology epoch.

...

Test Plan

Integration Tests

Broker

  • StreamsGroupTopologyDescriptionUpdateRequestTest (new) — push happy path; permanent and transient plugin failures PluginPermanentFailureException (ratchets LastFailedTopologyEpoch) and PluginTransientFailureException (arms back-off) both surface as STREAMS_TOPOLOGY_DESCRIPTION_UPDATE_FAILED with the cause in ErrorMessage; heartbeat-flag gating; member/group/MemberId fencing; explicit DeleteGroups plugin-success path (tombstone written) and plugin-failure path (STREAMS_TOPOLOGY_DESCRIPTION_DELETE_FAILED returned with ErrorMessage, group not tombstoned, retry converges after plugin recovery).
  • StreamsGroupTopologyDescriptionUpdateNoPluginRequestTest (new) — broker without plugin: pushes rejected, flag never set, describe returns NOT_STORED.
  • AuthorizerIntegrationTestREAD ACL on the GROUP resource for the new RPC.

...