【问题标题】:Celery skips task by an hour芹菜跳过一个小时的任务
【发布时间】: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

【问题讨论】:

  • 您是否尝试过添加更多工人?如果在您的任务触发时没有可用的,它必须等到一个可用。
  • 如何添加工人?我是这个任务运行背景的新手
  • 谢谢!但这能解决我的问题吗?
  • 那是你告诉我们的!

标签: python django celery


【解决方案1】:

可能有几个问题。最有可能的是,当您的任务被触发时,您的工作人员正忙。你可以通过增加工人来防止这种情况。 docs 解释了单个工作进程的--concurrency 选项,以及运行多个工作进程的选项。

您还可以运行附加到不同项目的不同工作人员,以便将某些任务分配给某些项目。即,某些任务的专用队列:Starting worker with dynamic routing_key?

我还看到,工作人员可以预取任务并保留它们——但如果它当前运行的任务运行超过倒计时,您的任务可能会延迟。

你会想要阅读CELERYD_PREFETCH_MULTIPLIER

【讨论】:

  • 非常感谢!我将不得不去并发,顺便说一下,我们的互联网是间歇性的,有时会导致任务失败,这也是一个原因,对吗?
  • 添加更多日志并找出答案?
  • 我其实用的是芹菜原木
【解决方案2】:

Celery 工作人员在准备就绪时从队列中弹出任务,但如果任务有倒计时,它会同时弹出其他任务并通过执行其他操作等待时间到期。它不保证任务会在那个时候运行,至少在那个时候或更晚。

【讨论】:

  • 那么我的解决方案是什么?有没有办法始终遵循设置的周期任务时间?
  • 使用 cron 获得更好的保证
  • 对不起,我对这些东西不是很熟悉?你能建议或制作一个替代我的代码吗?
  • cron 是一个与 beat 完全一样的程序,但更好一点 - 超级稳定,一直存在,不难使用。我会朝那个方向走,否则你必须建立一个死信交换,这仍然无法让你保持这种一致性。
  • 我不知道 - 我以前没有使用过它,但我在这些情况下使用了 cron 很多次并且从来没有遇到过问题 - 它安装在大多数 linux 发行版上,并且有一个任务调度程序窗户也能正常工作。
猜你喜欢
  • 1970-01-01
  • 2018-05-11
  • 1970-01-01
  • 2016-08-19
  • 1970-01-01
  • 2014-12-17
  • 2014-12-04
  • 2014-02-03
  • 1970-01-01
相关资源
最近更新 更多