【问题标题】:Django, Celery with Recursion and Twitter APIDjango,带有递归和 Twitter API 的 Celery
【发布时间】:2013-02-08 13:39:15
【问题描述】:

我正在使用 Django 1.4 和 Celery 3.0 (rabbitmq) 构建任务组合,用于获取和缓存 Twitter API 1.1 的查询。我正在尝试实现的一件事是任务链,其中最后一个基于迄今为止的响应和最近检索到的响应中的响应数据,对两个节点的任务进行递归调用。具体来说,这允许应用程序遍历用户时间线(最多 3200 条推文),考虑到任何给定的请求最多只能产生 200 条推文(Twitter API 的限制)。

可以看到我的 tasks.py 的关键组件here,但在粘贴之前,我将展示我从我的 Python shell 调用的链(但最终将通过最终 Web 应用程序中的用户输入启动)。给定:

>>request(twitter_user_id='#1010101010101#, 
  total_requested=1000, 
  max_id = random.getrandbits(128) #e.g. arbitrarily large number)

我打电话:

>> res = (twitter_getter.s(request) | 
        pre_get_tweets_for_user_id.s() | 
        get_tweets_for_user_id.s() | 
        timeline_recursor.s()).apply_async()

关键是timeline_recursor 可以启动可变数量的get_tweets_for_user_id 子任务。当timeline_recursor 处于其基本情况时,它应该返回一个响应字典,定义如下:

@task(rate_limit=None)
def timeline_recursor(request):
    previous_tweets=request.get('previous_tweets', None) #If it's the first time through, this will be None
    if not previous_tweets:
        previous_tweets = [] #so we initiate to empty array
    tweets = request.get('tweets', None) 

    twitter_user_id=request['twitter_user_id']
    previous_max_id=request['previous_max_id']
    total_requested=request['total_requested']
    pulled_in=request['pulled_in']

    remaining_requested = total_requested - pulled_in
    if previous_max_id:
        remaining_requested += 1 #this is because cursored results will always have one overlapping id

    else:
        previous_max_id = random.getrandbits(128) # for first time through loop

    new_max_id = min([tweet['id'] for tweet in tweets])
    test = lambda x, y: x<y

    if remaining_requested < 0:  #because we overshoot by requesting batches of 200
        remaining_requested = 0

    if tweets:
        previous_tweets.extend(tweets)

    if tweets and remaining_requested and (pulled_in > 1) and test(new_max_id, previous_max_id):

        request = dict(user_pk=user_pk,
                    twitter_user_id=twitter_user_id,
                    max_id = new_max_id,
                    total_requested = remaining_requested,
                    tweets=previous_tweets)

        #problem happens in this part of the logic???

        response = (twitter_getter_config.s(request) | get_tweets_for_user_id.s() | timeline_recursor.s()).apply_async()

    else: #if in base case, combine all tweets pulled in thus far and send back as "tweets" -- to be 
          #saved in db or otherwise consumed
        response = dict(
                    twitter_user_id=twitter_user_id,
                    total_requested = total_requested,
                    tweets=previous_tweets)
    return response

因此,我对 res.result 的预期响应是一个字典,其中包含 twitter 用户 ID、请求的推文数量以及在连续调用中提取的推文集。 然而,在递归任务领域,一切都不是很好。当我运行上面确定的链时,如果我在启动链后立即输入 res.status,它表示“成功”,即使在我的 celery worker 的日志视图中,我也可以看到链式递归调用到 twitter api 正在按预期使用正确的参数进行制作。即使正在执行链式任务,我也可以立即运行 result.result。 res.result 产生一个 AsyncResponse 实例 id。即使在递归链式任务完成运行之后, res.result 仍然是 AsyncResult id。

另一方面,我可以通过访问 res.result.result.result.result['tweets'] 访问我的完整推文集。我可以推断出每个连锁的连锁子任务确实正在发生,我只是不明白为什么 res.result 没有预期的结果。当 timeline_recursor 获得其基本情况时应该发生的递归返回似乎没有按预期传播。

对可以做什么有什么想法吗? Celery 中的递归可以变得非常强大,但至少对我来说,我们应该如何考虑使用 Celery 的递归和递归函数以及这如何影响链式任务中的 return 语句的逻辑并不完全清楚。

很高兴根据需要澄清,并提前感谢您的任何建议。

【问题讨论】:

  • 为了它的价值,这里是芹菜工人的日志:pastebin.com/M4SkYBBb
  • 如何获得 3200 条推文?不应该是 3000 - 当前的速率限制是每 15 分钟窗口 15 个请求(每个请求 x 200 条推文)
  • 3200 与 Twitter 在任何给定时间为给定用户提供的最大推文数量有关。对于速率限制,这就是您所说的,我的任务配置了以下选项:@task(rate_limit="12/m", max_retries=3) 此外,您每 15 分钟允许对该端点进行 180 次调用,所以你的号码是关闭的:dev.twitter.com/docs/api/1.1/get/statuses/user_timeline
  • 啊好的。我以为你在调用 get_home_timline (15) ,而不是 get_user_timeline (180)
  • .apply_async 返回一个 AsyncResult 实例。这是一个承诺,所以结果可能还没有准备好,但你可以用它来等待任务的结果,或者检查它的进度。但是 - 您应该永远不要等待子任务完成,因为这可能会导致资源匮乏并最终导致死锁。我想这里的问题是如何使用 Celery 进行递归:您可以通过使用另一个回调来做到这一点:(twitter_getter_config.s(request) | get_tweets_for_user_id.s() | timeline_recursor.s() | CONTINUE.s())

标签: python django twitter rabbitmq celery


【解决方案1】:

apply_async 返回什么(如对象类型)?

我不知道 celery,但在 Twisted 和许多其他异步框架中......当任务被推迟到队列。

再次,不知道 celery ,我猜这正在发生:

你是:立即将 response 定义为异步延迟 task,但随后尝试对其采取行动,就好像结果已经到来一样

您想成为:定义一个callback 例程以在任务完成后在结果上运行并返回一个值

查看 celery 文档,apply_async 通过 link 接受回调 - 我找不到任何人试图从中捕获返回值的示例。

【讨论】:

  • 虽然我同意你关于如何实现的广泛逻辑,但这确实(据我所知)如何在 Celery 中实现回调。要链接回调,您可以使用您引用的链接选项,但您也可以简单地将任务链接在一起(docs.celeryproject.org/en/latest/userguide/canvas.html#chains),就像我在 res = (twitter_getter.s(request) | pre_get_tweets_for_user_id.s() | get_tweets_for_user_id .s() | timeline_recursor.s()).apply_async()
  • 在 celery 中,可以将任务的结果分配给变量。我相信预期的行为是 var.status 为“PENDING”,直到链中的每个任务都返回。实际上,如果您调用 var.get() 并且子任务仍处于挂起状态,则在执行完成之前您不会得到回报。 FWIW,我观察到不使用递归的链中链和回调的预期行为 - 意外结果似乎是由递归返回驱动的,没有按预期流入最终结果。
  • jonathan,您在分配响应后对响应采取行动的问题绝对正确。我现在在该链上使用 .get() 而不是 apply_async() ,这似乎会生成预期的递归调用,但不幸的副作用是导致 Celery 挂起。
  • 这是做什么的get_tweets_for_user_id.s() - 通常使用延迟调用,您会发送函数或函数+ args 以便稍后执行。 .s() 是 @task 装饰器的延迟操作吗?
  • get_tweets_for_user_id.s() 使用链中前一个节点的返回值应用具有该名称的任务。在 celery 中,.s() 是使用星形参数创建子任务签名的捷径(来源:docs.celeryproject.org/en/latest/userguide/canvas.html#subtasks
猜你喜欢
  • 2013-03-03
  • 2015-03-06
  • 2018-05-02
  • 2012-01-10
  • 2020-12-26
  • 2019-01-18
  • 2019-07-02
  • 2014-08-20
  • 1970-01-01
相关资源
最近更新 更多