【问题标题】:Django celery task keep global stateDjango celery 任务保持全局状态
【发布时间】:2013-11-20 11:30:02
【问题描述】:

我目前正在开发一个基于 django-tenants-schema 的 Django 应用程序。您不需要查看模块的实际代码,但想法是它具有当前数据库连接的全局设置,定义用于应用程序租户的架构,例如

tenant = tenants_schema.get_tenant()

对于设置

tenants_schema.set_tenant(xxx)

对于 一些 任务,我希望他们记住在实例化期间选择的当前全局租户,例如理论上:

class AbstractTask(Task):
    '''
    Run this method before returning the task future
    '''
    def before_submit(self):
         self.run_args['tenant'] = tenants_schema.get_tenant()

    '''
    This method is run before related .run() task method
    '''
    def before_run(self):
         tenants_schema.set_tenant(self.run_args['tenant'])

在芹菜中有没有优雅的做法?

【问题讨论】:

    标签: django celery


    【解决方案1】:

    Celery(从 3.1 开始)具有 signals,您可以使用它来执行此操作。您可以更改传入的kwargs,另一方面,在将更改交给实际任务之前撤消更改:

    from celery import shared_task
    from celery.signals import before_task_publish, task_prerun, task_postrun
    from threading import local
    
    current_tenant = local()
    
    @before_task_publish.connect
    def add_tenant_to_task(body=None, **unused):
        body['kwargs']['tenant_middleware.tenant'] = getattr(current_tenant, 'id', None)
        print 'sending tenant: {t}'.format(t=current_tenant.id)
    
    
    @task_prerun.connect
    def extract_tenant_from_task(kwargs=None, **unused):
        tenant_id = kwargs.pop('tenant_middleware.tenant', None)
        current_tenant.id = tenant_id
        print 'current_tenant.id set to {t}'.format(t=tenant_id)
    
    
    @task_postrun.connect
    def cleanup_tenant(**kwargs):
        current_tenant.id = None
        print 'cleaned current_tenant.id'
    
    
    @shared_task
    def get_current_tenant():
        # Here is where you would do work that relied on current_tenant.id being set.
        import time
        time.sleep(1) 
        return current_tenant.id
    

    如果您运行任务(不显示来自工作人员的日志记录):

    In [1]: current_tenant.id = 1234; ct = get_current_tenant.delay(); current_tenant.id = 5678; ct.get()
    sending tenant: 1234
    Out[1]: 1234
    
    In [2]: current_tenant.id
    Out[2]: 5678
    

    如果没有发送消息,则不会调用信号(当您直接调用任务函数时,没有delay()apply_async())。如果要过滤任务名称,它在before_task_publish 信号处理程序中作为body['task'] 可用,task 对象本身在task_preruntask_postrun 处理程序中可用。

    我是 Celery 新手,所以我无法确定这是否是在 Celery 中做“中间件”类型的东西的“幸运”方式,但我认为它对我有用。

    【讨论】:

      【解决方案2】:

      我不确定你的意思,before_submit 是在客户端调用任务之前执行的吗?

      在这种情况下,我宁愿在这里使用 with 语句:

      from contextlib import contextmanager
      
      @contextmanager
      def set_tenant_db(tenant):
          prev_tenant = tenants_schema.get_tenant()
          try:
              tenants_scheme.set_tenant(tenant)
              yield
          finally:
              tenants_schema.set_tenant(prev_tenant)
      
      
      @app.task
      def tenant_task(tenant=None):
          with set_tenant_db(tenant):
              do_actions_here()
      
      
      tenant_task.delay(tenant=tenants_scheme.get_tenant())
      

      您当然可以创建一个自动执行此操作的基本任务, 例如,您可以在Task.__call__ 中应用上下文,但我不确定 如果你可以显式地使用 with 语句,这样可以为你节省很多。

      【讨论】:

      • before_submit() 必须在task.apply_async()task() 期间或在启动任务的任何其他情况下、在当前进程中或其他情况下调用。正如我所看到的:它只是任务*args, **kwargs 的过滤器,它添加了一个新参数,以便将其传递给任务运行程序进程。我需要将当前全局变量隐式保存在Task基类中,然后在将要运行任务的进程中隐式恢复它。
      猜你喜欢
      • 2012-02-11
      • 1970-01-01
      • 2014-06-13
      • 2016-10-30
      • 2013-05-14
      • 2018-05-24
      • 2014-03-16
      • 2020-06-30
      • 2017-06-16
      相关资源
      最近更新 更多