Celery SoftTimeLimitExceeded and TimeLimitExceeded

Viewed 1991

Below is the code from documentation:

from celery.exceptions import SoftTimeLimitExceeded

@celery.task(soft_time_limit=15, time_limit=20)
def mytask():
    try:
        return do_work()
    except SoftTimeLimitExceeded:
        cleanup_in_a_hurry()

The question is how celery allows to catch the exception inside the function. If it executes my_task and raises SoftTimeLimitExceeded, how is this exception propagates inside the function?

Also, why it is not possible to catch TimeLimitExceeded inside the function?

Thank you.

2 Answers

In short, Celery interrupts your task via a signal and raises SoftTimeLimitExceeded.

Note that there are a few different ways Celery can be configured to run tasks (e.g. threads), but my answer is limited to process pools. In this case, your task is being executed in one of the worker's child processes. When these pool processes are created, Celery registers a signal handler to handle the SIGUSR1 signal. Upon timeout, SIGUSR1 is sent. This interrupts your tasks: wherever they were in their Python bytecode execution, they stop, add Celery's signal handler onto their stack frames, and execute Celery's handler, which raises SoftTimeLimitExceeded. The exception propagates up the stack frame (note this will be wherever your task was when the interruption occurred) until it is caught, presumably by your task code.

SoftTimeLimit exists for that exact reason - so you can catch the exception, and handle it. Hard limit is there to actually stop the task from running if the limit is reached. I think it was deliberately (and I would add rightfully) designed like that so we, developers do not mess things up.

Here is an example how to catch the SoftTimeLimitExceeded exception (https://github.com/scoringengine/scoringengine/blob/master/scoring_engine/engine/execute_command.py):

from scoring_engine.celery_app import celery_app
from celery.exceptions import SoftTimeLimitExceeded
import subprocess

from scoring_engine.logger import logger


@celery_app.task(name='execute_command', acks_late=True, reject_on_worker_lost=True, soft_time_limit=30)
def execute_command(job):
    output = ""
    # Disable duplicate celery log messages
    if logger.propagate:
        logger.propagate = False
    logger.info("Running cmd for " + str(job))
    try:
        cmd_result = subprocess.run(
            job['command'],
            shell=True,
            stdout=subprocess.PIPE,
            stderr=subprocess.STDOUT
        )
        output = cmd_result.stdout.decode("utf-8")
        job['errored_out'] = False
    except SoftTimeLimitExceeded:
        job['errored_out'] = True
    job['output'] = output
    return job
Related