Repository navigation
Conversation
| # run. If the task really did die, heartbeat detection will fail it once the heartbeat stops. | ||
| heartbeat_timeout = conf.getint("scheduler", "task_instance_heartbeat_timeout") | ||
| ti_alive = ( | ||
| ti.state == TaskInstanceState.RUNNING |
There was a problem hiding this comment.
This change looks like a no-op to me, on L1450 we have
if state in (TaskInstanceState.QUEUED, TaskInstanceState.RUNNING):
....
continue
which means ti_alive will never be true as the first part of the condition is going to be False.
There was a problem hiding this comment.
Thanks for digging into this! I don't think it's a no-op though(but I could be mistaken) — the two conditions look at different things.
The state we unpack on L1488 comes out of event_buffer.pop(buffer_key). That buffer is filled by the executor — every success()/fail()/queued()/running_state() call writes (state, info) into it. So state here is "what the executor reported", not the TI's DB state. buffer_key is just the
lookup key into that dict.
ti_alive, on the other hand, checks ti.state, which we loaded from the database in the session.scalars(...) query above.
So passing the L1490 continue only tells us the executor didn't report QUEUED/RUNNING — it says nothing about what's in the DB. By the time we reach ti_alive we're handling a terminal executor event (SUCCESS/FAILED), and at that moment ti.state can still be RUNNING. That's could be exactly the
duplicate-dispatch case this guards against: a live worker is actually running the task and heartbeating, while a stale terminal event lands from the duplicate that already died. Without the check we'd fail the live run out from under the worker.
I can add a short comment making the state (executor) vs ti.state (DB) distinction explicit if that'd help future readers.
There was a problem hiding this comment.
@ashb I went over the code and it looks to me that @dsuhinin is right. These are different.
The state in your code snippet is the state that we get from the executor as a response.
state, info = event_buffer.pop(buffer_key)
if state in (TaskInstanceState.QUEUED, TaskInstanceState.RUNNING):
ti.external_executor_id = info
cls.logger().info("Setting external_executor_id for %s to %s", ti, info)
continueWhile the state in these changes, is the one that we read from the DB.
locked_query = with_row_locks(query, of=TI, session=session, skip_locked=True)
tis: Iterator[TI] = session.scalars(locked_query)
for ti in tis:
xBis7
left a comment
There was a problem hiding this comment.
The original root cause of the issue was that multiple schedulers could queue the same task and then send it to separate workers for execution. The upstream PRs that have already been merged focused on preventing the scheduler from queueing the task
but the current code will still kill a healthy running task. Assuming that it was somehow queued by multiple schedulers, if the current scheduler is the last one that tried to queue it and wrote its ID in the DB, then the task will be terminated.
The scheduler decides on a last writer wins basis and it's good to also use the new heartbeat check as an additional safeguard to determine if the task is still alive.
| "Task instance already running on another worker; standing down without failing it", | ||
| workload_id=str(ti.id), | ||
| ) | ||
| return 0 |
There was a problem hiding this comment.
As far as I know, Celery is the only executor currently handling TaskAlreadyRunningError and when it does, it ignores the error.
This except is equally applied to all executors but it's also running before Celery's except statement. So the error won't be handled in Celery anymore but the behavior is conflicting.
Returning 0 means a success while Celery ignores the error.
See
and
I think the behavior should be consistent.
We can re-raise the error here and let it propagate to the executors. Celery will catch it and treat the task failure as a no event.
There was a problem hiding this comment.
I've been running tests and I think that the changes in the supervisor should be reverted entirely.
Without the heartbeat check in the scheduler, this doesn't add any value and makes things worse for celery because it overrides the existing behavior.
For a local executor, this returns 0 and then the worker sends a success event back to the scheduler. The scheduler compares the reported event from the executor (state SUCCESS) with the value in the DB (state RUNNING) and because there is a mismatch, it updates the DB so that the running task is marked as FAILED. The worker picks up the update during the next heartbeat and kills the task.
For local and the other executors, this isn't any different than the existing behavior where the executor sends a state FAILED event back to the scheduler. In both cases (success/failure), there is a mismatch between the executor reported state and the state in the DB.
For celery, the existing executor event is ignored and therefore the scheduler doesn't have something to react on. But with the change, it gets back a SUCCESS event and goes back to killing the running task. With the heartbeat check, there isn't an issue but without it, this is a regression.
| # run. If the task really did die, heartbeat detection will fail it once the heartbeat stops. | ||
| heartbeat_timeout = conf.getint("scheduler", "task_instance_heartbeat_timeout") | ||
| ti_alive = ( | ||
| ti.state == TaskInstanceState.RUNNING |
There was a problem hiding this comment.
@ashb I went over the code and it looks to me that @dsuhinin is right. These are different.
The state in your code snippet is the state that we get from the executor as a response.
state, info = event_buffer.pop(buffer_key)
if state in (TaskInstanceState.QUEUED, TaskInstanceState.RUNNING):
ti.external_executor_id = info
cls.logger().info("Setting external_executor_id for %s to %s", ti, info)
continueWhile the state in these changes, is the one that we read from the DB.
locked_query = with_row_locks(query, of=TI, session=session, skip_locked=True)
tis: Iterator[TI] = session.scalars(locked_query)
for ti in tis:| # A running task that's still sending heartbeats is alive -- a worker is running it right now. | ||
| # This event is probably from a duplicate that already lost and died, so don't fail the live | ||
| # run. If the task really did die, heartbeat detection will fail it once the heartbeat stops. | ||
| heartbeat_timeout = conf.getint("scheduler", "task_instance_heartbeat_timeout") |
There was a problem hiding this comment.
We are in a for ti in tis: loop and we are reading the config again for every iteration. We should read it only once, outside of the loop.
| and ti.last_heartbeat_at >= timezone.utcnow() - timedelta(seconds=heartbeat_timeout) | ||
| ) | ||
|
|
||
| if ti_queued and not ti_requeued and not ti_alive: |
There was a problem hiding this comment.
I think a metric to monitor how many times the ti_alive condition has been satisfied, might be useful for the future. This could be a follow-up.
xBis7
left a comment
There was a problem hiding this comment.
I've been running some tests and I think that this PR should only include the heartbeat check in the scheduler and then probably there should be some follow-up PRs updating the rest of the executors to behave the same way as celery.
I've added a comment explaining why I think the supervisor change should be removed.
| "Task instance already running on another worker; standing down without failing it", | ||
| workload_id=str(ti.id), | ||
| ) | ||
| return 0 |
There was a problem hiding this comment.
I've been running tests and I think that the changes in the supervisor should be reverted entirely.
Without the heartbeat check in the scheduler, this doesn't add any value and makes things worse for celery because it overrides the existing behavior.
For a local executor, this returns 0 and then the worker sends a success event back to the scheduler. The scheduler compares the reported event from the executor (state SUCCESS) with the value in the DB (state RUNNING) and because there is a mismatch, it updates the DB so that the running task is marked as FAILED. The worker picks up the update during the next heartbeat and kills the task.
For local and the other executors, this isn't any different than the existing behavior where the executor sends a state FAILED event back to the scheduler. In both cases (success/failure), there is a mismatch between the executor reported state and the state in the DB.
For celery, the existing executor event is ignored and therefore the scheduler doesn't have something to react on. But with the change, it gets back a SUCCESS event and goes back to killing the running task. With the heartbeat check, there isn't an issue but without it, this is a regression.
When a task instance gets dispatched twice, only one worker actually runs it — the duplicate loses the race and dies. The problem was that the scheduler treated the leftover event from that dead duplicate as a failure and killed the live run that was happily executing.
This fixes both sides of the race:
possible fix for: #57041