【问题标题】:How to schedule my crawler function in django periodically using celery?如何使用 celery 定期在 django 中安排我的爬虫功能?
【发布时间】:2020-09-19 20:40:17
【问题描述】:

这里我有一个视图CrawlerHomeView,它用于从表单创建任务对象,现在我想用 celery 定期安排这个任务。

我想使用任务对象 search_frequency 并检查一些任务对象字段来安排这个 CrawlerHomeView 进程。

任务模型

class Task(models.Model):
    INITIAL = 0
    STARTED = 1
    COMPLETED = 2

    task_status = (
        (INITIAL, 'running'),
        (STARTED, 'running'),
        (COMPLETED, 'completed'),
        (ERROR, 'error')
    )

    FREQUENCY = (
        ('1', '1 hrs'),
        ('2', '2 hrs'),
        ('6', '6 hrs'),
        ('8', '8 hrs'),
        ('10', '10 hrs'),
    )

    name = models.CharField(max_length=255)
    scraping_end_date = models.DateField(null=True, blank=True)
    search_frequency = models.CharField(max_length=5, null=True, blank=True, choices=FREQUENCY)
    status = models.IntegerField(choices=task_status)

tasks.py

如果任务状态为 0 或 1 并且未超过任务抓取结束日期,我想运行下面定期发布的视图 [period=(task's search_frequency time]。但我卡在这里。我该怎么做?

@periodic_task(run_every=crontab(hour="task.search_frequency"))  # how to do  with task search_frequency value
def schedule_task(pk):
    task = Task.objects.get(pk=pk)
    if task.status == 0 or task.status == 1 and not datetime.date.today() > task.scraping_end_date:

        # perform the crawl function ---> def crawl() how ??

        if task.scraping_end_date == datetime.date.today():
            task.status = 2
            task.save()  # change the task status as complete.

views.py

我想定期运行这个视图。我该怎么做?

class CrawlerHomeView(LoginRequiredMixin, View):
    login_url = 'users:login'

    def get(self, request, *args, **kwargs):
        # all_task = Task.objects.all().order_by('-id')
        frequency = Task()
        categories = Category.objects.all()
        targets = TargetSite.objects.all()
        keywords = Keyword.objects.all()

        form = CreateTaskForm()
        context = {
            'targets': targets,
            'keywords': keywords,
            'frequency': frequency,
            'form':form,
            'categories': categories,
        }
        return render(request, 'index.html', context)

    def post(self, request, *args, **kwargs):

        form = CreateTaskForm(request.POST)
        if form.is_valid():

            # try:
            unique_id = str(uuid4()) # create a unique ID. 
            obj = form.save(commit=False)

            # obj.keywords = keywords
            obj.created_by = request.user
            obj.unique_id = unique_id
            obj.status = 0
            obj.save()
            form.save_m2m()

            keywords = ''
            # for keys in ast.literal_eval(obj.keywords.all()): #keywords change to csv
            for keys in obj.keywords.all():
                if keywords:
                    keywords += ', ' + keys.title
                else:
                    keywords += keys.title
            # tasks = request.POST.get('targets')
            # targets = ['thehimalayantimes', 'kathmandupost']
            # print('$$$$$$$$$$$$$$$ keywords', keywords)

            task_ids = [] #one Task/Project contains one or multiple scrapy task

            settings = {
                'spider_count' : len(obj.targets.all()),
                'keywords' : keywords,
                'unique_id': unique_id, # unique ID for each record for DB
                'USER_AGENT': 'Mozilla/5.0 (compatible; Googlebot/2.1; +http://www.google.com/bot.html)'
            }

            # res = ast.literal_eval(ini_list) 

            for site_url in obj.targets.all():
                domain = urlparse(site_url.address).netloc # parse the url and extract the domain
                spider_name = domain.replace('.com', '')
                task = scrapyd.schedule('default', spider_name, settings=settings, url=site_url.address, domain=domain, keywords=keywords)

            # task = scrapyd.schedule('default', spider_name , settings=settings, url=obj.targets, domain=domain, keywords=obj.keywords)
            return redirect('crawler:task-list')
            # except:
            #     return render(request, 'index.html', {'form':form})
        return render(request, 'index.html', {'form':form, 'errors':form.errors})

对于这个问题有什么建议或答案吗?

【问题讨论】:

  • Task.frequency是一个什么样的字段?持续时间字段?
  • @mark 显示你的scheduler.models.Task 实现
  • @ncopiy 已在 qiestion 中更新
  • @Ben 在问题中更新

标签: django web-scraping django-views celery


【解决方案1】:

对于错误,

Exception Type: EncodeError
Exception Value:    
Object of type timedelta is not JSON serializable 

而不是在 django 设置中定义以下变量,

CELERY_BEAT_SCHEDULE = {

    'task-first': {
    'task': 'scheduler.tasks.create_task',
    'schedule': timedelta(minutes=1)
   },

你可以在你的芹菜文件中尝试以下内容吗:

app.conf.beat_schedule = {
    'task-first': {
        'task': 'scheduler.tasks.create_task',
        'schedule': crontab(minute='*/1')
    }
}

这对我有用,芹菜服务器已启动并正在运行。

除此之外,您为什么要在每个任务之后重定向到'list_tasks',它究竟做了什么?另外,您已经从add_task_celery.delay(name,date,freq) 视图中调用了 celery 任务,除了使用 celery-beat 定义的定期任务之外,它只是另一种添加任务的方法吗?

编辑 1:

我的结构如下:

settings.py

CELERY_TIMEZONE = 'Asia/Kolkata'
CELERY_BROKER_URL = 'amqp://localhost'

芹菜.py

app.conf.beat_schedule = {
    'task1': {
        'task': '<app_name>.tasks.random_task',
        'schedule': crontab(minute=0, hour=0)
    },
}

这里你应该注意到我的应用文件夹中有一个名为 tasks 的文件,我在那里编写了一个共享任务,如下所示:

@shared_task
def random_task(total):
    ...

此外,除此之外,您应该同时启动 celery beat 和 celery worker 进程,如下所示:

celery -A <project_name>.celery worker -l error
celery -A <project_name>.celery beat -l error --scheduler django_celery_beat.schedulers:DatabaseScheduler

您可以使用任何您想要的调度程序,在生产中我使用DatabaseScheduler。对于测试,您可以尝试使用以下命令:

celery -A <project_name> beat -l info -S django

您应该从 Django 项目的项目文件夹中运行所有这些命令

【讨论】:

  • 感谢您的回答。这是我第一次使用芹菜,所以发生了这种情况,我也认为这不是调用任务的正确方法。你能帮我指导我如何在这里安排任务。一步一步的过程是什么。我想在数据库中每 1 分钟创建一次任务对象,并希望在 list_tasks 页面中显示 1 分钟后的结果
  • 更新了我的分步过程答案
  • 我已经编辑了我的问题,你可以看看这个吗?
  • 你已经完全改变了这个问题。你为什么不使用它来制作一个函数,它会获取你想要获取的所有参数并从视图和任务中调用该函数。此外,您应该从 celery 任务中调用该函数并提供所有必要的参数,而不是从任务中调用视图。基本上,与其在视图中编写所有代码,不如编写一个可以从 celery 和视图中调用的可重用函数。
【解决方案2】:

在 15k 任务/秒的设置中与 Celery 抗争了 5 年之后,我强烈建议您切换到 Dramatiq,,它有一个健全、可靠、高性能的代码库,不会拆分为多个复杂的包,并且可以完美地分为两个到目前为止我的新项目。

来自author's motivation

多年来,我一直专业地使用 Celery,我对它越来越感到沮丧,这也是我开发 Dramatiq 的原因之一。以下是 Dramatiq、Celery 和 RQ 之间的一些主要区别:

还有一个 Django 帮助程序包:https://github.com/Bogdanp/django_dramatiq

当然,你不会有内置的 celerybeat,但是调用 python 任务的 cron 无论如何都会更健壮,我们丢失了大量数据,因为 celerybeat 决定定期停止 :)


有两个项目旨在添加周期性任务创建:https://gitlab.com/bersace/periodiqhttps://apscheduler.readthedocs.io/en/stable/

我还没有使用这些包,您可以尝试使用 periodiq 选择您的数据库条目,遍历这些并为每个条目定义一个周期性任务(但这需要定期重新启动 periodiq 工作程序以获取更改) :

# tasks.py
from dramatiq import get_broker
from periodiq import PeriodiqMiddleware, cron

broker = get_broker()
broker.add_middleware(PeriodiqMiddleware(skip_delay=30))


for obj in Task.objects.all():
   @dramatiq.actor(periodic=cron(obj.frequency))
   def hourly(obj=obj):
       # import logic based on obj.name
       # Do something each hour…

【讨论】:

  • 似乎更容易,但我必须使用数据库表中的持续时间,我可以在 Dramatiq 中以这种方式安排任务吗?
  • @mark: 使用示例代码进行了更新,这只是一个想法如何与基于数据库的计划一起工作:)
  • 这里是 periodiq @dramatiq.actor(periodic=cron('0 * * * *))。 ** ?secs/minutes/hours/ 的值应该是多少?格式是什么?你能告诉我,因为我在文档中没有找到它。
  • @user12428848 请参阅 en.wikipedia.org/wiki/Cron 了解 cron 语法。 0 * * * * 基本上说“在第 0 分钟每整小时运行一次这个 cron”(01:00 am、02:00 am、03:00 am...)。
  • 谢谢。但是当我在终端上运行命令periodiq -v app ii 时会抛出这个错误[2020-07-20 22:01:07,160] [PID 2188] [MainThread] [periodiq] [CRITICAL] Unsupported system: alarm syscall is not available.
【解决方案3】:

我认为问题出在任务定义中的第二个和第三个参数,即freqdate。尽管从错误中,您发布了 Object of type timedelta is not JSON serializable,但它看起来像是在谈论 freq 字段,该字段是 DurationField 类型,返回 timedelta 对象。

理想情况下,两个字段都必须在传递给任务之前进行序列化。 一种简单的方法是 -

1)您可以显式序列化这些字段并传递给任务,并在任务中再次将其转换为 datetime / timedelta 对象。

或者,如果项目太多,您可以转储整个数据字典。

add_task_celery.delay(json.dumps(form.cleaned_data)),

然后在任务中做 -> json.loads(...)

2) 您可以尝试的另一件事是在参数中显式传递序列化程序。(使用apply_async 而不是delay

add_task_celery.apply_async((name, date, freq), serializer='json')

3) 你也可以设置值,如果你还没有设置CELERY_TASK_SERIALIZER = 'json'(默认值是'pickle')。

【讨论】:

  • add_task_celery() 在使用第二种方法时得到了一个意外的关键字参数“序列化器”
  • add_task_celery() argument after ** must be a mapping, not datetime.date 使用 apply_async 后
  • 我有一个表单来创建任务对象所以我用正确的输入填写表单(在使用芹菜之前工作)然后在提交后我抛出这个错误。芹菜版本是 4+
  • 他已经用json序列化了,否则他不会有json序列化错误。问题是您无法在运行时轻松注册另一个序列化程序,而必须在设置时使用 Kombu 进行注册。 DjangoJSONEncoder 已经解决了这个问题,因为它支持序列化时间。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-09-27
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-07-12
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多