【问题标题】:How to stop the execution of a long process if something changes in the db?如果数据库发生变化,如何停止执行长进程?
【发布时间】:2021-07-30 02:21:07
【问题描述】:

我有一个向RabbitMQ 队列发送消息的视图。

message = {'origin': 'Bytes CSV',
           'data': {'csv_key': str(csv_entry.key),
                    'csv_fields': csv_fields
                    'order_by': order_by,
                    'filters': filters}}

...

queue_service.send(message=message, headers={}, exchange_name=EXCHANGE_IN_NAME,
                   routing_key=MESSAGES_ROUTING_KEY.replace('#', 'bytes_counting.create'))

在我的消费者身上,我有一个很长的过程来生成 CSV。

def create(self, data):
    csv_obj = self._get_object(key=data['csv_key'])
    if csv_obj.status == CSVRequestStatus.CANCELED:
        self.logger.info(f'CSV {csv_obj.key} was canceled by the user')
        return

    result = self.generate_result_data(filters=data['filters'], order_by=data['order_by'], csv_obj=csv_obj)
    csv_data = self._generate_csv(result=result, csv_fields=data['csv_fields'], csv_obj=csv_obj)
    file_key = self._post_csv(csv_data=csv_data, csv_obj=csv_obj)

    csv_obj.status = CSVRequestStatus.READY
    csv_obj.status_additional = CSVRequestStatusAdditional.SUCCESS
    csv_obj.file_key = file_key
    csv_obj.ready_at = timezone.now()
    csv_obj.save(update_fields=['status', 'status_additional', 'ready_at', 'file_key'])

    self.logger.info(f'CSV {csv_obj.name} created')

长过程发生在self._generate_csv 内部,因为self.generate_result_data 返回一个queryset,这是惰性的。

如您所见,如果用户在消息开始被消费之前通过端点更改csv_request 的状态,则不会评估进程。我的目标是在执行self._generate_csv 期间让这种情况发生。

到目前为止,我尝试使用Threading,但没有成功。

我怎样才能实现我的目标?

非常感谢!

【问题讨论】:

  • "我的目标是在执行 self._generate_csv 期间让这种情况发生。"你这是什么意思?
  • @LordElrond 如果csv_obj 状态更改为CSVRequestStatus.CANCELED,我想停止执行此函数(即使它已经启动)。

标签: python django multithreading rabbitmq


【解决方案1】:

你为什么不检查 Celery 库?将celery with djangoRabbitMQ backend 一起使用比直接利用rabbitmq 队列要容易得多。

Celery 有一个内置函数 revoke 来终止正在进行的任务:

>>> from celery.task.control import revoke
>>> revoke(task_id, terminate=True)

对于您的用例,您可能需要类似(代码 sn-ps):

## celery/tasks.py
from celery import app

@app.task(queue="my_queue")
def create_csv(message):
    # ...snip...
    pass

## main.py
from celery import uuid, current_app

def start_task(task_id, message):
    current_app.send_task(
        "create_csv",
        args=[message],
        task_id=task_id,
    )

def kill_task(task_id):
    current_app.control.revoke(task_id, terminate=True)

## signals.py

from django.dispatch import receiver
from .models import MyModel
from .main import kill_task

# choose appropriate signal to listen for DB change
@receiver(models.signals.post_save, sender=MyModel)
def handler(sender, instance, **kwargs):
    kill_task(instance.task_id)
  • 使用celery.uuid生成可以存储在DB或缓存中的任务ID,并使用相同的任务ID来控制任务,即请求终止。

【讨论】:

  • 虽然我同意 celery 会简化一些事情,但我担心仅仅为了这个功能而改变一切会是一种矫枉过正......
  • 请注意,celery.task 已被贬值,在 v5.x 中根本不起作用。
  • @MuriloSitonio 最初我同意你的看法,因为单独使用 celery 来完成这项任务将是矫枉过正。然而,在尝试找到一个“纯粹”的解决方案后,我得出了这样的结论:自己实现它的维护开销将远远超过 celery 所需的开销,因为它可能需要数百行要实现的代码(大概使用pika)。 TLDR; 我会推荐芹菜!
  • @LordElrond 感谢您努力寻找“纯粹”的解决方案!我仍在努力寻找一种不会超过芹菜所需开销的方法,但恐怕我会同意你的看法。无论如何,如果你想用这个“纯粹”的解决方案分享你的思路,我会非常感谢你!
【解决方案2】:

由于 self._generate_csv 是最慢的,显而易见的解决方案是使用此函数。

为此,您可以将 csv 文件的创建分成几个部分。创建每个片段后,检查状态,看看是否可以继续创建文件。最后,将所有部分粘贴到一个完成的文件中。

Here is a method for combining multiple files into one.

【讨论】:

  • 这确实是一个好主意,但self._generate_csv 不会仅生成 csv,它还会评估queryset。我可以为我最终加入的每个较小的 csv 切片 queryset,但是我会多次访问数据库,并且数据可能会在两者之间发生变化。或者,也许我只是一次评估整个queryset,然后我可以将其像列表一样切片。这可能确实有效。我要测试一下。谢谢!
  • 您可以生成一次查询集,只访问数据库一次,然后对其进行切片或使用迭代器。例如,PostgreSQL 使用服务器端游标从数据库中流式传输结果,而无需将整个结果集加载到内存中。参考:docs.djangoproject.com/en/2.2/ref/models/querysets/#iterator
  • @TrevorCox 哇,我从未听说过迭代器!但是我仍然需要为 N 个切片敲击 db N 次,因为在每个切片上我都必须检查 csv_obj 的状态是否已更改...
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多