【发布时间】:2016-11-04 00:37:41
【问题描述】:
但是当我刚刚打开我的计算机并运行 celery 时,它每 50 秒运行一次任务,我看到一些跳过,比如 1 小时。除了意外的跳过之外,它实际上执行得很好。为什么会这样?如何解决?
这是我的工人 -l 信息中跳过日志的示例
2016-11-03 10:13:36,264: INFO/MainProcess] Task core.tasks.sample[8efcedc5-1e08-41c4-80b9-1f82a9ddbaad] succeeded in 1.062010367s: None
[2016-11-03 11:14:19,751: INFO/MainProcess] Received task: core.tasks.sample[ca9d6ef4-2cdc-4546-a9fb-c413541a80ee]
这是我的节拍 -l 信息中跳过日志的示例
[2016-11-03 10:13:35,199: INFO/MainProcess] Scheduler: Sending due task core.tasks.sample (core.tasks.sample)
[2016-11-03 11:14:19,748: INFO/MainProcess] Scheduler: Sending due task core.tasks.sample (core.tasks.sample)
这是我的任务代码:
# 50 seconds
@periodic_task(run_every=timedelta(**settings.XXX_XML_PERIODIC_TASK))
def sample():
global GLOBAL_CURRENT_DATE
if cache.get('XXX_xml_today_saved_data') is None:
cache.set('XXX_xml_today_saved_data', [])
saved_data = cache.get('XXX_xml_today_saved_data')
ftp = FTP('xxxxx')
ftp.login(user='xxxxx', passwd='xxxxx')
ftp.cwd('XXX')
date_dir = GLOBAL_CURRENT_DATE.replace("-", "")
try:
ftp.cwd(date_dir)
except:
ftp.cwd(str(int(date_dir) - 1))
_str = StringIO()
files = ftp.nlst()
if (GLOBAL_CURRENT_DATE != datetime.now().strftime("%Y-%m-%d") and
files == saved_data):
GLOBAL_CURRENT_DATE = datetime.now().strftime("%Y-%m-%d")
cache.delete('XXX_xml_today_saved_data')
return
print files
print "-----"
print saved_data
unsaved = list(set(files) - set(saved_data))
print "-----"
print unsaved
if unsaved:
file = min(unsaved)
# modified_time = ftp.sendcmd('MDTM '+ file)
print file
ftp.retrbinary('RETR ' + file, _str.write)
xml = '<root>'
xml += _str.getvalue()
xml += '</root>'
if cache.get('XXX_provider_id') is None:
cache.set('XXX_provider_id', Provider.objects.get(code="XXX").id)
_id = cache.get('XXX_provider_id')
_dict = xmltodict.parse(xml, process_namespaces=True,
dict_constructor=dict, attr_prefix="")
row = _dict['root']['row']
if type(_dict['root']['row']) == dict:
_dict['root']['row'] = []
_dict['root']['row'].append(row)
row = _dict['root']['row']
for x in row:
if cache.get('XXX_data_type_' + x['dataType']) is None:
obj, created = DataType.objects.get_or_create(code=x['dataType'])
obj, created = ProviderDataType.objects.get_or_create(provider_id=_id, data_type=obj)
if created:
cache.set('XXX_data_type_' + x['dataType'], obj.id)
_id = cache.get('XXX_data_type_' + x['dataType'])
obj, created = Transaction.objects.get_or_create(data=x, file_name=file,
provider_data_type_id=_id)
if created:
if x['dataType'] == "BR":
print "Transact"
br_transfer(**x)
else:
print "Not transacting"
saved_data.append(file)
cache.set('XXX_xml_today_saved_data', saved_data)
ftp.close()
这是我在 settings.py 中的 CELERY CONFIGS:
BROKER_URL = 'redis://localhost:6379'
CELERY_RESULT_BACKEND = 'redis://localhost:6379'
CELERY_ACCEPT_CONTENT = ['application/json']
CELERY_TASK_SERIALIZER = 'json'
CELERY_RESULT_SERIALIZER = 'json'
CELERY_TIMEZONE = 'Africa/Nairobi'
XXX_XML_PERIODIC_TASK = {'seconds': 50}
CACHES = {
'default': {
'BACKEND': 'redis_cache.RedisCache',
'LOCATION': 'localhost:6379',
'TIMEOUT': None,
},
}
有什么解释或建议吗?
我正在使用 python 2.7.10 和 django 1.10
【问题讨论】:
-
您是否尝试过添加更多工人?如果在您的任务触发时没有可用的,它必须等到一个可用。
-
如何添加工人?我是这个任务运行背景的新手
-
docs.celeryproject.org/en/latest/userguide/workers.html - 先试试 --concurrency
-
谢谢!但这能解决我的问题吗?
-
那是你告诉我们的!