lucasbru commented on code in PR #22622:
URL: https://github.com/apache/kafka/pull/22622#discussion_r3467093696
##########
group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java:
##########
@@ -786,6 +794,144 @@ private void
throwIfStreamsGroupTopologyDescriptionUpdateInvalid(
throwIfNull(request.topologyDescription(), "TopologyDescription can't
be null.");
}
+ /**
+ * Build one topology-description cleanup cycle: read every shard for
streams groups
+ * eligible for plugin-side cleanup (empty + all offsets expired +
storedEpoch != -1), call
+ * {@code plugin.deleteTopology} for each via the manager, then for every
group whose
+ * plugin call succeeded write a conditional metadata record that clears
+ * {@code StoredDescriptionTopologyEpoch} only if the persisted value
still matches the
+ * epoch we observed at scan time (so a concurrent {@code setTopology}
that has advanced
+ * the field is preserved). Failed plugin calls retry on the next cycle;
the next sweep
+ * then tombstones the now-empty group.
+ *
+ * <p>The single-flight guard, periodic timer scheduling, and {@code
running} flag live
+ * on {@link StreamsGroupTopologyDescriptionManager#startCleanupCycle};
this method is
+ * the cycle body it invokes, returning a future that the manager joins to
release the
+ * in-flight flag.
+ *
+ * <p><b>Concurrent setTopology race vs plugin.deleteTopology.</b> {@code
plugin.deleteTopology}
+ * is keyed only on {@code groupId}. If a new member joins between the
+ * eligibility scan and the cycle's plugin call and pushes a fresh
topology, the plugin's
+ * row is removed regardless of the new epoch — the conditional clear
above no-ops on the
+ * metadata side, but the plugin-side data the member just wrote is gone.
A subsequent
+ * {@code describe} → {@code getTopology} returns null and surfaces {@code
NOT_STORED} with
+ * a warn log; this is the graceful-degradation path accepted under the
label
+ * "plugin-side data loss". The {@code isEmpty} requirement on the scan
keeps the window
+ * narrow — concurrent setTopology requires a member to join an empty,
fully-expired group
+ * between scan and delete — and the next heartbeat at the same epoch will
not re-solicit
+ * (storedEpoch in metadata still reflects the new push), so the group
converges on
+ * NOT_STORED without churn rather than chasing the lost plugin row.
+ */
+ // Visible for testing.
+ CompletableFuture<?> runOneStreamsTopologyCleanupCycle() {
+ if (!streamsGroupTopologyDescriptionManager.isPluginConfigured()) {
+ return CompletableFuture.completedFuture(null);
+ }
+ groupCoordinatorMetrics.recordSensor(
+
GroupCoordinatorMetrics.STREAMS_GROUP_TOPOLOGY_DESCRIPTION_CLEANUP_CYCLE_RUNS_SENSOR_NAME);
+
+ List<CompletableFuture<Map<String, Integer>>> partitionFutures =
runtime.scheduleReadAllOperation(
+ "list-streams-groups-needing-topology-cleanup",
+ GroupCoordinatorShard::listStreamsGroupsNeedingTopologyCleanup
+ );
+
+ // ConcurrentLinkedQueue because per-partition .handle callbacks can
append concurrently
+ // from whichever thread completed each runtime read.
+ Queue<CompletableFuture<?>> perGroupFutures = new
ConcurrentLinkedQueue<>();
+ List<CompletableFuture<Void>> partitionDoneFutures = new
ArrayList<>(partitionFutures.size());
+ for (CompletableFuture<Map<String, Integer>> partitionFuture :
partitionFutures) {
+ partitionDoneFutures.add(partitionFuture.handle((eligible,
throwable) -> {
+ if (throwable != null) {
+ log.warn("Topology-description cleanup read failed for one
partition.", throwable);
+ return null;
+ }
+ if (eligible == null || eligible.isEmpty()) return null;
+ // Shutdown started after the per-partition read was
scheduled. Skip the
+ // plugin dispatch so we do not issue plugin.deleteTopology
calls into a
+ // manager whose plugin is about to be closed.
+ if (!isActive.get()) return null;
+ groupCoordinatorMetrics.recordSensor(
+
GroupCoordinatorMetrics.STREAMS_GROUP_TOPOLOGY_DESCRIPTION_CLEANUP_ELIGIBLE_GROUPS_SENSOR_NAME,
+ eligible.size()
+ );
+ perGroupFutures.add(streamsGroupTopologyDescriptionManager
+ .invokeDeleteTopologies(eligible.keySet())
+ .thenCompose(failures -> {
+ recordPluginDeleteOutcome(eligible.size(),
failures.size());
+ // Shutdown can have started between the plugin call
and the
+ // follow-up writes. Skip the conditional clears so we
do not
+ // schedule writes against a runtime that is being
closed.
+ if (!isActive.get()) return
CompletableFuture.completedFuture(null);
+ List<CompletableFuture<Void>> clearFutures = new
ArrayList<>(eligible.size());
+ eligible.forEach((groupId, expectedStoredEpoch) -> {
+ if (failures.containsKey(groupId)) {
+ // Plugin failed: leave both stored epoch and
the push-path
+ // back-off in place. Eligibility's "group is
empty" snapshot
+ // only held at scan time; a member can rejoin
between scan
+ // and now, and the existing back-off
correctly throttles their
+ // set-topology attempt against the
still-broken plugin
+ // instead of letting it re-attack at
attempts=0 every join.
+ return;
+ }
+ // Plugin succeeded; the group will be tombstoned
in the next sweep
+ // once the stored epoch is cleared. Drop the
broker-wide back-off
+ // entry — it is no longer load-bearing for any
future state of
+ // this groupId. A member that re-creates the same
id afterwards
+ // is a fresh lifecycle and will arm a fresh
back-off chain.
+
streamsGroupTopologyDescriptionManager.clearBackoffGroup(groupId);
+
clearFutures.add(clearStoredDescriptionTopologyEpochAsync(groupId,
expectedStoredEpoch));
+ });
+ return
CompletableFuture.allOf(clearFutures.toArray(new CompletableFuture<?>[0]));
+ }));
+ return null;
+ }));
+ }
+
+ return CompletableFuture.allOf(partitionDoneFutures.toArray(new
CompletableFuture<?>[0]))
+ .thenCompose(__ ->
CompletableFuture.allOf(perGroupFutures.toArray(new CompletableFuture<?>[0])));
+ }
+
+ /**
+ * Record per-call outcomes from a batched {@code plugin.deleteTopology}
invocation
+ * against the shared {@code delete-success} / {@code delete-error}
sensors. Used by
+ * the periodic cleanup cycle and the explicit {@code DeleteGroups} flow
so a single
+ * pair of meters tracks every {@code plugin.deleteTopology} the broker
drives,
+ * regardless of trigger.
+ */
+ private void recordPluginDeleteOutcome(int attempted, int errors) {
+ int successes = attempted - errors;
+ if (successes > 0) {
+ groupCoordinatorMetrics.recordSensor(
+
GroupCoordinatorMetrics.STREAMS_GROUP_TOPOLOGY_DESCRIPTION_DELETE_SUCCESS_SENSOR_NAME,
successes);
+ }
+ if (errors > 0) {
+ groupCoordinatorMetrics.recordSensor(
+
GroupCoordinatorMetrics.STREAMS_GROUP_TOPOLOGY_DESCRIPTION_DELETE_ERROR_SENSOR_NAME,
errors);
+ }
+ }
+
+ /**
+ * Conditional metadata write that clears {@code
StoredDescriptionTopologyEpoch} for
+ * {@code groupId} only when the persisted value still equals {@code
expectedStoredEpoch}.
+ * Mismatches and missing groups are silently ignored by the shard-side
method. Runtime
+ * write failures (NOT_COORDINATOR etc.) are logged here and swallowed so
a single failed
+ * write does not poison the cycle's allOf — the next cycle will retry
naturally because
+ * the persisted storedEpoch is still non-default.
+ */
+ private CompletableFuture<Void>
clearStoredDescriptionTopologyEpochAsync(String groupId, int
expectedStoredEpoch) {
+ return runtime.<Void>scheduleWriteOperation(
+ "clear-stored-topology-epoch",
+ topicPartitionFor(groupId),
+ coordinator ->
coordinator.clearStoredDescriptionTopologyEpoch(groupId, expectedStoredEpoch)
Review Comment:
In most cases I don't expect us to clean up many topologies at once, but
yeah for corner cases where for some reason we get many topologies expiring in
the same window, that could be a useful optimization.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]