FrankChen021 commented on code in PR #20387:
URL: https://github.com/apache/druid/pull/20387#discussion_r4071531278
##########
indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java:
##########
@@ -4488,6 +4488,18 @@ private void createNewTasks() throws
JsonProcessingException
continue;
}
+ // In bounded mode, do not top up replicas for a task group that has
already reached its end offsets.
+ // Otherwise, as completed replicas exit, this loop keeps recreating
replacement replicas whose start
+ // offsets are already at the bounded end, producing an endless churn of
tasks that complete instantly.
+ // This mirrors the completion guard on the task-group recreation path
in this method.
+ if (ioConfig.isBounded() && hasTaskGroupReachedBoundedEnd(groupId)) {
Review Comment:
[P2] Avoid metadata lookup for full replica groups
**Finding:** `hasTaskGroupReachedBoundedEnd(groupId)` is evaluated for every
non-empty active group before the `replicas > taskGroup.tasks.size()` check.
Thus a bounded supervisor whose groups already have all configured replicas
performs an additional datasource-metadata lookup per group on every supervisor
run (and, when the configs match, a second lookup through
`getOffsetsFromMetadataStorage()`), even though no top-up decision is needed.
With many task groups this adds synchronous SQL reads to every run and can
surface a transient metadata-store failure while no task creation was required.
**Suggestion:** Short-circuit on the replica-count comparison first and
invoke the bounded-end predicate only inside the top-up branch, or reuse the
metadata snapshot already fetched at the start of `createNewTasks()`.
##########
indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisorStateTest.java:
##########
@@ -2961,6 +2961,115 @@ public void
testComputeUnassignedServerPriorities_whenMultipleReplicasPerPriorit
verifyAll();
}
+ /**
+ * In bounded mode, a task group that has already reached its end offsets
must not have replacement replicas
+ * created for it. Without the guard in {@code createNewTasks()}'s replica
top-up loop, {@code replicas > tasks}
+ * would submit a replacement replica whose start offset is already at the
bounded end, so it completes instantly
+ * and re-triggers the top-up, churning tasks endlessly. This test flows
through {@code runInternal()} and asserts
+ * that no top-up task is submitted; it fails if the guard is removed.
+ */
+ @Test
+ public void testCreateNewTasks_boundedGroupReachedEnd_doesNotTopUpReplicas()
+ {
+ // replicas = 2, taskCount = 1, bounded with an empty range (start == end)
so the group has reached its end.
+ final BoundedStreamConfig boundedConfig = new BoundedStreamConfig(
Review Comment:
[P3] Exercise a completed non-empty bounded range
**Finding:** The replacement regression test sets `startSequenceNumbers ==
endSequenceNumbers`, so `hasTaskGroupReachedBoundedEnd` returns through its
empty-range branch and never checks the metadata offsets that define the
reported failure: a previously non-empty group whose committed offsets have
reached the configured end. The test can therefore pass while the
metadata-based completion path or its bounded-config matching is broken,
leaving the real replica-churn case unprotected.
**Suggestion:** Use a non-empty configured range and return matching end
offsets and bounded configuration from the metadata mock while retaining the
under-replicated active group, then assert that no task is submitted.
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]