Celery starts task each hour unexpectedly

Viewed 52

I am facing strange behaviour of Celery that I didn't expect and can't understand. Here is schedule I user:

CELERYBEAT_SCHEDULE = {
    'task-5872-accrual': {
        'task': 'Task5872Accrual',
        'schedule': crontab(minute='49', hour='16'),
    },
}

So starting celery beat:

$ celery -A app:celery beat --loglevel=DEBUG

I see that schedule taken correctly and beat triggers task at 16:49:

LocalTime -> 2021-08-20 16:46:35
Configuration ->
    . broker -> redis://localhost:6379/0
    . loader -> celery.loaders.app.AppLoader
    . scheduler -> celery.beat.PersistentScheduler
    . db -> celerybeat-schedule
    . logfile -> [stderr]@%DEBUG
    . maxinterval -> 5.00 minutes (300s)
[2021-08-20 16:46:35,682: DEBUG/MainProcess] Setting default socket timeout to 30
[2021-08-20 16:46:35,682: INFO/MainProcess] beat: Starting...
[2021-08-20 16:46:35,692: DEBUG/MainProcess] Current schedule:
<ScheduleEntry: task-5872-accrual Task5872Accrual() <crontab: 49 16 * * * (m/h/d/dM/MY)>
[2021-08-20 16:46:35,692: DEBUG/MainProcess] beat: Ticking with max interval->5.00 minutes
[2021-08-20 16:46:35,692: DEBUG/MainProcess] beat: Waking up in 2.41 minutes.
[2021-08-20 16:49:00,088: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 16:49:00,100: INFO/MainProcess] Scheduler: Sending due task task-5872-accrual (Task5872Accrual)
[2021-08-20 16:49:00,110: DEBUG/MainProcess] Task5872Accrual sent. id->58e8b8a0-8a53-4c14-aceb-c45c810b78b2
[2021-08-20 16:49:00,110: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 16:54:00,137: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 16:54:00,138: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 16:59:00,238: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 16:59:00,239: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 17:04:00,252: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 17:04:00,257: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 17:09:00,356: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 17:09:00,362: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 17:14:00,462: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 17:14:00,463: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 17:19:00,563: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 17:19:00,568: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 17:24:00,669: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 17:24:00,674: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 17:29:00,774: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 17:29:00,775: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 17:34:00,875: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 17:34:00,880: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 17:39:00,908: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 17:39:00,914: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 17:44:01,014: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 17:44:01,015: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 17:49:01,078: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 17:49:01,083: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 17:54:01,174: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 17:54:01,176: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 17:59:01,272: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 17:59:01,278: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 18:04:01,376: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 18:04:01,381: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 18:09:01,472: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 18:09:01,477: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 18:14:01,557: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 18:14:01,559: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 18:19:01,600: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 18:19:01,606: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 18:24:01,628: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 18:24:01,634: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 18:29:01,734: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 18:29:01,735: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 18:34:01,803: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 18:34:01,808: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 18:39:01,826: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 18:39:01,831: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 18:44:01,928: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 18:44:01,930: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 18:49:02,015: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 18:49:02,020: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.
[2021-08-20 18:54:02,097: DEBUG/MainProcess] beat: Synchronizing schedule...
[2021-08-20 18:54:02,103: DEBUG/MainProcess] beat: Waking up in 5.00 minutes.

but worker I starting:

$ celery -A app:celery worker -E --loglevel=ERROR

behaves strange:

 -------------- celery@al-desktop v5.0.5 (singularity)
--- ***** ----- 
-- ******* ---- Linux-5.11.0-27-generic-x86_64-with-glibc2.29 2021-08-20 16:46:39
- *** --- * --- 
- ** ---------- [config]
- ** ---------- .> app:         web:0x7f616ae3ab80
- ** ---------- .> transport:   redis://localhost:6379/0
- ** ---------- .> results:     redis://localhost:6379/0
- *** --- * --- .> concurrency: 64 (prefork)
-- ******* ---- .> task events: ON
--- ***** ----- 
 -------------- [queues]
                .> celery           exchange=celery(direct) key=celery
                

[2021-08-20 16:49:00,124: ERROR/ForkPoolWorker-62] Task5872Accrual started.time_limit: 86400
[2021-08-20 16:49:00,125: ERROR/ForkPoolWorker-62] Integrator Initialization. __init__.py. file: ./models/5872_accrual_BD_etl2.json 
[2021-08-20 16:49:00,130: ERROR/ForkPoolWorker-62] DB.py, __init__. db connection string: postgresql://test:********@localhost:5432/db
[2021-08-20 16:51:09,570: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 52910 - ['$52909']
[2021-08-20 16:51:10,143: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 53331 - ['$420']
[2021-08-20 16:51:10,891: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 53827 - ['$495']
[2021-08-20 16:51:34,785: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 57821 - ['$3993']
[2021-08-20 16:53:10,281: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 80626 - ['$22804']
[2021-08-20 16:53:17,659: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 83517 - ['$2890']
[2021-08-20 16:53:29,830: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 87511 - ['$3993']
[2021-08-20 16:54:34,153: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 109561 - ['$22049']
[2021-08-20 17:39:23,443: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 230152 - ['$120590']
[2021-08-20 17:49:53,123: ERROR/ForkPoolWorker-63] Task5872Accrual started.time_limit: 86400
[2021-08-20 17:49:53,124: ERROR/ForkPoolWorker-63] Integrator Initialization. __init__.py. file: ./models/5872_accrual_BD_etl2.json 
[2021-08-20 17:49:53,130: ERROR/ForkPoolWorker-63] DB.py, __init__. db connection string: postgresql://test:********@localhost:5432/etldb2
[2021-08-20 17:52:02,562: ERROR/ForkPoolWorker-63] accrual 5872 model_accrual.py. Broken string exception cautch at line 52910 - ['$52909']
[2021-08-20 17:52:03,129: ERROR/ForkPoolWorker-63] accrual 5872 model_accrual.py. Broken string exception cautch at line 53331 - ['$420']
[2021-08-20 17:52:03,874: ERROR/ForkPoolWorker-63] accrual 5872 model_accrual.py. Broken string exception cautch at line 53827 - ['$495']
[2021-08-20 17:52:26,980: ERROR/ForkPoolWorker-63] accrual 5872 model_accrual.py. Broken string exception cautch at line 57821 - ['$3993']
[2021-08-20 17:54:02,567: ERROR/ForkPoolWorker-63] accrual 5872 model_accrual.py. Broken string exception cautch at line 80626 - ['$22804']
[2021-08-20 17:54:09,939: ERROR/ForkPoolWorker-63] accrual 5872 model_accrual.py. Broken string exception cautch at line 83517 - ['$2890']
[2021-08-20 17:54:22,011: ERROR/ForkPoolWorker-63] accrual 5872 model_accrual.py. Broken string exception cautch at line 87511 - ['$3993']
[2021-08-20 17:55:28,189: ERROR/ForkPoolWorker-63] accrual 5872 model_accrual.py. Broken string exception cautch at line 109561 - ['$22049']
[2021-08-20 18:03:51,994: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 256604 - ['$26451']
[2021-08-20 18:29:55,192: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 288335 - ['$31730']
[2021-08-20 18:33:14,232: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 303105 - ['$14769']
[2021-08-20 18:33:45,468: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 303388 - ['$282']
[2021-08-20 18:33:47,998: ERROR/ForkPoolWorker-62] accrual 5872 model_accrual.py. Broken string exception cautch at line 303609 - ['$220']
[2021-08-20 18:39:51,124: ERROR/ForkPoolWorker-63] accrual 5872 model_accrual.py. Broken string exception cautch at line 230152 - ['$120590']
[2021-08-20 18:49:53,667: ERROR/ForkPoolWorker-64] Task5872Accrual started.time_limit: 86400
[2021-08-20 18:49:53,668: ERROR/ForkPoolWorker-64] Integrator Initialization. __init__.py. file: ./models/5872_accrual_BD_etl2.json 
[2021-08-20 18:49:53,675: ERROR/ForkPoolWorker-64] DB.py, __init__. db connection string: postgresql://test:********@localhost:5432/db
[2021-08-20 18:52:05,560: ERROR/ForkPoolWorker-64] accrual 5872 model_accrual.py. Broken string exception cautch at line 52910 - ['$52909']
[2021-08-20 18:52:06,134: ERROR/ForkPoolWorker-64] accrual 5872 model_accrual.py. Broken string exception cautch at line 53331 - ['$420']
[2021-08-20 18:52:06,874: ERROR/ForkPoolWorker-64] accrual 5872 model_accrual.py. Broken string exception cautch at line 53827 - ['$495']
[2021-08-20 18:52:30,275: ERROR/ForkPoolWorker-64] accrual 5872 model_accrual.py. Broken string exception cautch at line 57821 - ['$3993']
[2021-08-20 18:54:08,581: ERROR/ForkPoolWorker-64] accrual 5872 model_accrual.py. Broken string exception cautch at line 80626 - ['$22804']
[2021-08-20 18:54:16,090: ERROR/ForkPoolWorker-64] accrual 5872 model_accrual.py. Broken string exception cautch at line 83517 - ['$2890']
[2021-08-20 18:54:28,787: ERROR/ForkPoolWorker-64] accrual 5872 model_accrual.py. Broken string exception cautch at line 87511 - ['$3993']
[2021-08-20 18:55:35,893: ERROR/ForkPoolWorker-64] accrual 5872 model_accrual.py. Broken string exception cautch at line 109561 - ['$22049']

Here is my Base_task:

class BaseTask(Task):
    abstract = True
    acks_late = True
    reject_on_worker_lost = True
    ignore_result = False
    max_retries = int(os.getenv('CELERY_MAX_RETRIES_TIME', 0))
    time_limit = 86400
    soft_time_limit = 86400
    default_retry_delay = int(os.getenv('CELERY_TIME_COUNTDOWN_RETRY', 86400))
    expires = 86400
    validation_class = ''
    description = ''
    on_retry = retry_clbk
    on_failure = failure_clbk
    after_return = return_clbk

And here is task I am running by schedule:

class Task5872Accrual(BaseTask):
    name = 'Task5872Accrual'
    model_file = None

    def __init__(self, *args, **kwargs):
        super(self.__class__, self).__init__(*args, **kwargs)

    def run(self, model_file=None):
        self.model_file = model_file
        logger.error(f'Task5872Accrual started.time_limit: {self.time_limit}')
        try:
            integrator_run()
        except Exception:
            logger.error(f'Task5872Accrual ERROR {traceback.format_exc()}')

CELERY_MAX_RETRIES_TIME and CELERY_TIME_COUNTDOWN_RETRY not used by me at all so they are not defined.

I don't understand why worker starts new instance of same task each hour? Any ideas?

0 Answers
Related