Versions Compared

Key

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

...

Code Block
languagejs
linenumberstrue
{ "name": "ErrorMessage", "type": "string", "versions": "3+", "nullableVersions": "3+", "ignorable": true, "default": "null",
  "about": "AThe description of the deletion errorerror message, or null if the deletion succeeded orthere was no description is availableerror." }

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

...

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 PluginPermanentFailureExceptionStreamsTopologyDescriptionPermanentFailureException} — 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 PluginTransientFailureExceptionStreamsTopologyDescriptionTransientFailureException} 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);
}

...

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();
    }

    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();
    }
}

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 PluginPermanentFailureExceptionStreamsTopologyDescriptionPermanentFailureException extends ApiException {
    public PluginPermanentFailureExceptionStreamsTopologyDescriptionPermanentFailureException(String message) { super(message); }
    public PluginPermanentFailureExceptionStreamsTopologyDescriptionPermanentFailureException(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 PluginTransientFailureExceptionStreamsTopologyDescriptionTransientFailureException extends ApiException {
    public PluginTransientFailureExceptionStreamsTopologyDescriptionTransientFailureException(String message) { super(message); }
    public PluginTransientFailureExceptionStreamsTopologyDescriptionTransientFailureException(String message, Throwable cause) { super(message, cause); }
}

...

  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 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 PluginPermanentFailureException StreamsTopologyDescriptionPermanentFailureException it writes LastFailedTopologyEpoch = pushedEpoch so subsequent heartbeats at the same epoch do not re-solicit. On PluginTransientFailureException StreamsTopologyDescriptionTransientFailureException 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 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 DELETE_FAILED with the plugin's exception message in ErrorMessage (also logged at WARN); 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.

...

  • 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 setTopology future with PluginPermanentFailureException StreamsTopologyDescriptionPermanentFailureException; 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 StreamsTopologyDescriptionTransientFailureException (or any other exception); the broker's per-group back-off (30 s → 1 h, exponential) throttles re-solicitation.
  • 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 PluginPermanentFailureException StreamsTopologyDescriptionPermanentFailureException, PluginTransientFailureException StreamsTopologyDescriptionTransientFailureException, 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-{rate,count}MeterFailed plugin.getTopology calls.
kafka.server:type=group-coordinator-metrics,name=streams-group-topology-description-cleanup-cycle-{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.

...

  • StreamsGroupTopologyDescriptionUpdateRequestTest (new) — push happy path; PluginPermanentFailureException StreamsTopologyDescriptionPermanentFailureException (ratchets LastFailedTopologyEpoch) and PluginTransientFailureException StreamsTopologyDescriptionTransientFailureException (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 (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.

...