【问题标题】:Celery, periodic task execution, with concurrencyCelery,周期性任务执行,具有并发性
【发布时间】:2014-07-07 16:33:27
【问题描述】:

我想每秒启动一个周期性任务,但前提是前一个任务结束(数据库轮询以将任务发送到 celery)。 在 Celery 文档中,他们使用 Django 缓存来锁定。

我尝试使用示例:

from __future__ import absolute_import

import datetime
import time

from celery import shared_task

from django.core.cache import cache
LOCK_EXPIRE = 60 * 5

@shared_task
def periodic():

    acquire_lock = lambda: cache.add('lock_id', 'true', LOCK_EXPIRE)
    release_lock = lambda: cache.delete('lock_id')

    a = acquire_lock()
    if a:
        try:
            time.sleep(10)
            print a, 'Hello ', datetime.datetime.now()
        finally:
            release_lock()
    else:
        print 'Ignore'

具有以下配置:

app.conf.update(
    CELERY_IGNORE_RESULT=True,
    CELERY_ACCEPT_CONTENT=['json'],
    CELERY_TASK_SERIALIZER='json',
    CELERY_RESULT_SERIALIZER='json',
    CELERYBEAT_SCHEDULE={
        'periodic_task': {
            'task': 'app_task_management.tasks.periodic',
            'schedule': timedelta(seconds=1),
        },
    },
)

但是在控制台中,我从来没有看到Ignore 消息,而且我每秒都有Hello。锁好像不好使。

我启动周期性任务:

celeryd -B -A my_app

和工人:

celery worker -A my_app -l info

你能纠正我的误解吗?

【问题讨论】:

  • 不确定,但也许 Celery Workers 正在运行不同的进程,因此锁不适用?
  • 你用的是什么缓存?
  • @fixmycode: CACHES = { 'default': { 'BACKEND': 'django.core.cache.backends.locmem.LocMemCache', } }

标签: python django celery


【解决方案1】:

来自关于local-memory cache 的 Django 缓存框架文档:

请注意,每个进程都有自己的私有缓存实例, 意味着不可能进行跨进程缓存。

所以基本上你的工作人员都在处理自己的缓存。如果您需要低资源成本的缓存后端,我会推荐基于文件的缓存或数据库缓存,两者都允许跨进程。

【讨论】:

  • 谢谢,这解释了为什么它不能正常工作。我很尴尬,因为基于文件或数据库的缓存对于数据库轮询来说代价高昂。最好的解决方案是 memcache,但这意味着一个新的守护进程。有没有办法用 celery 和 rabbitmq 队列生成锁系统?
  • 我不知道。也许你可以使用 Redis 来完成它,也可以使用 Redis 作为 Celery 的代理而不是 RabbitMQ。 Redis 原生支持锁定键。
猜你喜欢
  • 2017-09-25
  • 1970-01-01
  • 2017-06-15
  • 2018-11-28
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-10-28
  • 1970-01-01
相关资源
最近更新 更多