Django中Celery调用类内方法报错:'do_work2'无'do_work3'属性
Celery类任务调用同类方法报错问题
问题重现
我是Celery新手,在Django中用它处理后台任务时遇到报错,代码如下:
from celery import shared_task, chain, group, chord, signature from django.conf import settings from modules.SVD.createVideo import CreateVideo import os import time import numpy as np from oralvis.celery import app class TestCelery: @app.task(bind=True) def do_work1(self, value): self.value = value value = value * 2 value = np.array(value) print('do_work1', value) return value @app.task(bind=True) def do_work2(self, list_of_numpyArrays): self.do_work3() # 尝试调用同类方法 print('process do work 2') print(list_of_numpyArrays) return {'status': True} def do_work3(self): print('from do work 3') print(self.value) # 创建类实例 testCeleryInstance = TestCelery() @shared_task def testCelery(value): # 编排任务链 ret = chain( (group(testCeleryInstance.do_work1.s(i) for i in range(value))), testCeleryInstance.do_work2.s()).apply_async()
运行后出现错误:
[2023-09-29 20:52:30,553: ERROR/MainProcess] Task SVDExperimentCompare.tasks.do_work2[de14846d-7c2d-43b0-9e83-50988d0cf8c0] raised unexpected: AttributeError("'do_work2' object has no attribute 'do_work3'") Traceback (most recent call last): File "/home/ubuntu/webapp/oralvis/env/lib/python3.10/site-packages/celery/app/trace.py", line 477, in trace_task R = retval = fun(*args, **kwargs) File "/home/ubuntu/webapp/oralvis/env/lib/python3.10/site-packages/celery/app/trace.py", line 760, in __protected_call__ return self.run(*args, **kwargs) File "/home/ubuntu/webapp/oralvis/oralvis/SVDExperimentCompare/tasks.py", line 21, in do_work2 self.do_work3() # Call the instance method using self AttributeError: 'do_work2' object has no attribute 'do_work3'
想知道是不是Celery运行类内函数时,无法访问同一类中的其他方法?
原因解析
不是Celery不能访问类内方法,而是你用@app.task(bind=True)装饰类方法时,Celery会把这个方法转换成Task子类的实例,而非原TestCelery类的实例。也就是说,do_work2执行时的self是Celery的Task对象,不是你创建的testCeleryInstance,自然找不到do_work3方法。
另外,Celery任务是跨进程执行的,类实例的状态无法序列化传递到worker进程,就算你强行绑定实例,也会因为状态丢失导致各种问题。
解决办法
方法1:把依赖方法改成独立函数或Celery任务
如果do_work3不需要依赖类实例状态,直接改成独立函数:
def do_work3(value): print('from do work 3') print(value) class TestCelery: @app.task(bind=True) def do_work1(self, value): value = value * 2 value = np.array(value) print('do_work1', value) return value @app.task(bind=True) def do_work2(self, list_of_numpyArrays): # 直接调用独立函数,传入需要的参数 do_work3(123) # 替换成实际需要的参数 print('process do work 2') print(list_of_numpyArrays) return {'status': True}
如果do_work3也需要异步执行,给它加上Celery任务装饰器:
class TestCelery: @app.task(bind=True) def do_work1(self, value): value = value * 2 value = np.array(value) print('do_work1', value) return value @app.task(bind=True) def do_work2(self, list_of_numpyArrays): # 异步调用同类任务方法 self.do_work3.delay(123) print('process do work 2') print(list_of_numpyArrays) return {'status': True} @app.task(bind=True) def do_work3(self, value): print('from do work 3') print(value)
方法2:使用Celery类任务(继承Task)
如果一定要用类组织任务逻辑,可以继承celery.Task,把整个类作为一个任务,内部方法可以正常调用:
from celery import Task class TestCelery(Task): name = 'test_celery_task' def do_work3(self, value): print('from do work 3') print(value) def do_work1(self, value): value = value * 2 value = np.array(value) print('do_work1', value) return value def do_work2(self, list_of_numpyArrays): self.do_work3(123) print('process do work 2') print(list_of_numpyArrays) return {'status': True} def run(self, value): # 在这里编排任务逻辑 chain( group(self.app.task(self.do_work1).s(i) for i in range(value)), self.app.task(self.do_work2).s() ).apply_async() # 注册任务到Celery app.tasks.register(TestCelery())
方法3:放弃类封装,用独立任务函数
最省心的方式是直接用独立的shared_task装饰函数,避免类实例带来的序列化问题:
@shared_task def do_work1(value): value = value * 2 value = np.array(value) print('do_work1', value) return value def do_work3(value): print('from do work 3') print(value) @shared_task def do_work2(list_of_numpyArrays): do_work3(123) print('process do work 2') print(list_of_numpyArrays) return {'status': True} @shared_task def testCelery(value): ret = chain( group(do_work1.s(i) for i in range(value)), do_work2.s() ).apply_async()
内容的提问来源于stack exchange,提问作者jxw
相关产品推荐
相关产品推荐

