Vivek1106-04 commented on code in PR #12277:
URL: https://github.com/apache/seatunnel/pull/12277#discussion_r3998540026
##########
seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/TaskExecutionService.java:
##########
@@ -1349,12 +1429,60 @@ public boolean runNewBusWork(boolean checkTaskQueue) {
BlockingQueue<Future<?>> futureBlockingQueue = new
LinkedBlockingQueue<>();
CooperativeTaskWorker cooperativeTaskWorker =
new CooperativeTaskWorker(taskQueue, this,
futureBlockingQueue);
- Future<?> submit =
executorService.submit(cooperativeTaskWorker);
+ sharedCooperativeWorkers.incrementAndGet();
+ Future<?> submit;
+ try {
+ submit = executorService.submit(cooperativeTaskWorker);
+ } catch (RuntimeException e) {
+ // The worker never runs, so it can never return its place
itself.
+ sharedCooperativeWorkers.decrementAndGet();
+ throw e;
+ }
futureBlockingQueue.add(submit);
return true;
}
return false;
}
+
+ /**
+ * Decides whether a worker running a slow task call may be promoted
to an exclusive worker.
+ *
+ * <p>A promotion costs one worker thread: the current worker stays
with the slow task and a
+ * replacement is started for the shared queue. The promotion is
therefore admitted by
+ * {@link CooperativeWorkerBudget} first. When the budget is exhausted
the worker keeps the
+ * slow task on the shared queue side and the caller retries later,
except that a
+ * replacement is still started when this is the last worker serving
the queue, so an
+ * exhausted budget can never stop queued tasks from reaching
readiness.
+ *
+ * @param worker the worker that is executing the slow task call
+ * @param taskTracker the slow task
+ * @return true when the worker was promoted, false when the budget
denied it
+ */
+ public boolean tryPromoteCooperativeWorker(
+ CooperativeTaskWorker worker, TaskTracker taskTracker) {
+ long jobId =
+ taskTracker
+ .taskGroupExecutionTracker
+ .taskGroup
+ .getTaskGroupLocation()
+ .getJobId();
+ if (!cooperativeWorkerBudget.tryAcquire(jobId)) {
+ logger.fine(
+ String.format(
+ "Promotion of a cooperative worker for job %d
was denied with reason BUDGET_EXHAUSTED, "
+ + "promoted workers: %d, denied
promotions: %d",
+ jobId,
+ cooperativeWorkerBudget.getPromotedWorkers(),
+
cooperativeWorkerBudget.getDeniedPromotions()));
+ if (sharedCooperativeWorkers.get() <= 1) {
Review Comment:
You are right, and thank you for reproducing it rather than only flagging
it. The counter was measuring the wrong thing: a worker that had not been
promoted was counted as serving the queue even while it was blocked inside a
task call, so with a per job budget of 1 and several blocked calls the check
saw workers that could not poll anything.
Fixed in f6d0251 by counting availability instead of promotion state:
- `queuePollingWorkers` counts workers currently waiting on the shared
queue. A worker increments it right before `takeFirst()` and decrements it
again as soon as it has a task, so a worker running a call is never counted.
- `pendingQueueWorkers` counts workers that were started but have not
reached the queue yet, so a worker on its way there is not missed. A worker
increments polling before it clears pending, so the sum never dips to zero
while it is in flight.
- `ensureQueueIsServed()` starts a worker when `queuePollingWorkers +
pendingQueueWorkers == 0`, under a lock with a double check, which also answers
the over-spawn concern raised in the other review: concurrent denials agree on
one worker instead of each starting its own.
I also took your point that the budget must not be advertised as a bound on
total threads. It bounds promotions; when calls block indefinitely the node
still ends up with one worker per blocked call, because that is what keeps
queued work starting. The English and Chinese deployment guides now state that
explicitly.
Regression added as
`testQueuedTaskStartsWhileDeniedPromotionsBlockEveryWorker`: three gated
cooperative tasks whose calls only return once a fourth task, queued behind
them, starts. I confirmed it reproduces your scenario by disabling
`ensureQueueIsServed()` locally, where it fails with "the queued task never
started, so denied promotions starved the shared queue" after 91 s, and passes
with the guard in place.
--
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]