【发布时间】: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', } }