Airflow 2.5.3跨调度保留XCom变量的可行方案求助
Airflow 2.5.3 自定义Sensor XCom存储替代方案(原2.0.2方案失效)
我在Airflow 2.0.2中曾用以下方案实现自定义Sensor:通过调用XCom.set方法,用原task_id拼接“_original”的不存在task_id存储requestId,以此避免任务调度清理XCom变量。但升级到2.5.3版本后,因为该task_id没有对应的task_instance,触发了外键约束错误(sqlalchemy.exc.IntegrityError),现在需要可行的替代方案。
原代码
def poke(self, context): request_id = context['ti'].xcom_pull(key="requestId", task_ids=self.task_id + "_original", dag_id=self.dag_id) if request_id is None: self.log.info("requestId not found, so we launch the Lambda function") request_id = self.launch_lambda() # To avoid task reschedule clean the xcom message, we use the XCom API directly to push the requestId # but with a taskId modified XCom.set( key="requestId", value=request_id, task_id=self.task_id + "_original", dag_id=self.dag_id, execution_date=context['ti'].execution_date) self.log.info("Pushed requestId %s to xcom", request_id)
错误信息
sqlalchemy.exc.IntegrityError: (psycopg2.errors.ForeignKeyViolation) insert or update on table "xcom" violates foreign key constraint "xcom_task_instance_fkey" DETAIL: Key (dag_id, task_id, run_id, map_index)=(performance_1_batchsensor_dev, eks-airflow-batchsensor_job_laucher_task_1_original, manual__2023-09-07T16:51:33.660188+00:00, -1) is not present in table "task_instance".
可行替代方案
方案1:使用Sensor实例保留状态(推荐)
Airflow Sensor的实例状态会在poke循环中保留,无需依赖XCom,彻底避免调度清理问题。修改代码如下:
def __init__(self, **kwargs): super().__init__(**kwargs) self.request_id = None def poke(self, context): if self.request_id is None: self.log.info("requestId not found, so we launch the Lambda function") self.request_id = self.launch_lambda() self.log.info("Generated requestId %s", self.request_id) # 编写检查Lambda任务状态的逻辑 task_completed = self.check_lambda_status(self.request_id) if task_completed: # 任务完成后,将requestId推送到正常XCom供后续任务使用 context['ti'].xcom_push(key="requestId", value=self.request_id) return True return False
方案2:使用DAG级XCom
Airflow支持存储不关联特定task_id的DAG级XCom,不会触发task_instance外键约束。修改XCom操作代码:
# 推送DAG级XCom XCom.set( key="requestId", value=request_id, task_id=None, # 不指定task_id,标记为DAG级别 dag_id=self.dag_id, execution_date=context['ti'].execution_date ) # 拉取DAG级XCom request_id = context['ti'].xcom_pull(key="requestId", task_ids=None, dag_id=self.dag_id)
DAG级XCom与DAG运行绑定,只要DAG运行存在就不会被清理。
方案3:添加辅助空Task
在DAG中新增一个与Sensor同属一个DAG的空Task(如DummyOperator),使用原来的xxx_original作为task_id,确保存在对应的task_instance:
# 在DAG定义代码中添加辅助任务 from airflow.operators.dummy import DummyOperator # 辅助任务仅用于生成对应task_instance,无需设置依赖关系 original_task = DummyOperator( task_id=self.task_id + "_original", dag=self.dag )
添加后原有的XCom代码可保持不变,因为xxx_original对应的task_instance会在DAG运行时自动创建。
内容的提问来源于stack exchange,提问作者Francisco Albert
相关产品推荐
相关产品推荐

