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]

Reply via email to