Celery任务链中如何避免第二个任务早于首个任务数据库更新执行?
解决Celery任务链中数据库更新后的竞态条件问题
问题背景
我有一条Celery任务链,其中第二个任务不仅需要等待第一个任务执行完成,还必须等第一个任务触发的所有数据库更新操作全部完成后才能执行。目前测试中我通过while循环等待链内第一个chord结束来满足需求,但这种阻塞式等待的方式并不合理。我也曾尝试使用transaction.on_commit,但并未解决问题。
原代码
# My simple model class TestModel(BaseModel): # The datetime the item occurred occurred_at = models.DateTimeField(null=False, blank=False, db_index=True) # A fake quantity quantity = models.DecimalField( max_digits=20, decimal_places=10, null=True, blank=True ) @shared_task def race_condition_tester() -> None: zulu = dateTime.utcnow().strftime('%Y-%m-%dT%H:%M:%SZ') workflow = chain( # This chord should asynchronously write a bunch of things to the database chord( generate_list_of_test_items.s( zulu=zulu, ), spawn_tasks_from_list_of_test_items_and_write_each_item_to_database.s(), ), # This should run after the chord, but also needs to run after the database # updates are complete test_if_items_are_in_database.s( zulu=zulu, ), ) workflow.apply_async() return None @shared_task def generate_list_of_test_items( results=None, zulu: str = None, ) -> List: logger.info(f"generate_list_of_test_items: {zulu=}") task_param_lists = [] for i in range(10): task_param_lists.append([i, zulu]) return task_param_lists @shared_task def spawn_tasks_from_list_of_test_items_and_write_each_item_to_database( results=None, ) -> None: task_param_list = results[0] # I use a chord here to asynchronously write all the items to the db chord_task = chord( [write_item_to_db.si( i=task_params[0], zulu=task_params[1], ) for task_params in task_param_list], chord_finisher.s() ) # Using this we get a resulting_sum of 36 which is incorrect # chord_task() # Using this we get a resulting_sum of 36 which is incorrect # chord_task.apply_async() # Using this we get a resulting_sum of 36 which is incorrect # chord_task.apply_async() # time.sleep(5) # Using this we get a resulting_sum of 36 which is incorrect # from django.db import transaction # transaction.on_commit(lambda: chord_task.apply_async()) # This works... we get 45 result = chord_task.apply_async() while not result.ready(): time.sleep(1) return None @shared_task def chord_finisher(*args, **kwargs): # noqa """ Just a simple task for the chord """ return "OK" @shared_task def write_item_to_db(i, zulu): import time from apps.pw import logics logger.info(f"write_item_to_db: {i=}, {zulu=}") time.sleep(i) models.TestModel.objects.create( occurred_at=dateTime.strptime( zulu, '%Y-%m-%dT%H:%M:%SZ' ).replace(tzinfo=TimeZone.utc), quantity=Decimal(i), ) return None @shared_task def test_if_items_are_in_database( results=None, zulu: str = None, ): from apps.pw import models from apps.pw import logics logger.info(f"test_if_items_are_in_database") activities = models.TestModel.objects.filter( occurred_at=dateTime.strptime( zulu, '%Y-%m-%dT%H:%M:%SZ' ).replace(tzinfo=TimeZone.utc) ) resulting_sum = Decimal(0) for activity in activities: resulting_sum += activity.quantity # resulting_sum should be 45 if resulting_sum != Decimal(45): logger.error(f"test_if_items_are_in_database: {resulting_sum=} but should be 45") else: logger.info(f"test_if_items_are_in_database: {resulting_sum=}") return None
问题根源
- 嵌套异步任务的依赖断裂:外层
chain中的第一个chord执行完spawn_tasks_from_list_of_test_items_and_write_each_item_to_database后,就会立即触发后续的test_if_items_are_in_database任务。但spawn_tasks内部只是异步启动了另一个chord,并没有等待这个内部chord的所有写库任务完成,导致查询任务启动时数据库数据还未完全写入。 - while循环的不合理性:阻塞等待内部
chord完成会占用Celery Worker进程,浪费资源,违背异步任务的设计初衷。 - transaction.on_commit无效:这里的问题并非数据库事务未提交,而是任务执行顺序的依赖错误。
write_item_to_db每个任务本身都在独立事务中执行,事务提交后数据就会写入,但查询任务启动太早,根本没等到这些任务执行完毕。
正确解决方案
利用Celery的任务签名传递特性:如果一个Celery任务返回另一个任务(或任务组/Chord)的签名,Celery会自动调度这个返回的任务,并等待其完成后再继续执行后续的任务链。
1. 修改核心任务逻辑
不再阻塞等待内部chord,而是返回其签名让外层链自动等待:
@shared_task def spawn_tasks_from_list_of_test_items_and_write_each_item_to_database( results=None, ): task_param_list = results[0] # 构建内部chord的签名 chord_task = chord( [write_item_to_db.si( i=task_params[0], zulu=task_params[1], ) for task_params in task_param_list], chord_finisher.s() ) # 返回chord签名,Celery会自动执行并等待其完成 return chord_task
2. (可选)简化工作流结构
如果不需要中间的spawn_tasks任务,可直接将查询任务作为内部chord的回调,进一步简化依赖:
@shared_task def race_condition_tester() -> None: zulu = dateTime.utcnow().strftime('%Y-%m-%dT%H:%M:%SZ') workflow = chain( generate_list_of_test_items.s(zulu=zulu), # 直接用生成的列表构建写库任务组,查询任务作为chord的回调 chord( # 引用上一步返回的任务参数列表,生成写库任务 (write_item_to_db.si(i=p[0], zulu=p[1]) for p in bind(0)), test_if_items_are_in_database.s(zulu=zulu) ) ) workflow.apply_async() return None
效果验证
修改后,外层chain会等待内部chord的所有写库任务执行完成(包括所有数据库事务提交),才会触发test_if_items_are_in_database任务,此时查询到的总和会正确为45,彻底避免竞态条件。
内容的提问来源于stack exchange,提问作者Brett Elliot
相关产品推荐
相关产品推荐

