Airflow Param参数取值异常求助:指定ID无法过滤CSV数据
问题分析
你遇到的核心问题是未正确获取Airflow UI触发时传入的动态参数:
直接在op_kwargs中使用dag.params["target_id"]只会读取DAG定义时的默认值(0),不会获取用户运行时输入的参数,因此无论你设置什么值,函数里拿到的永远是0,最终执行输出全部数据的逻辑。
解决方案
有两种可行的修正方式,任选其一即可:
方式1:通过任务上下文的params获取参数
修改read_csv函数,从任务上下文的params中读取运行时传入的参数,同时调整PythonOperator的配置:
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.models.param import Param from datetime import datetime, timedelta import pandas as pd default_args = { "owner": "bedirhan", "depends_on_past": False, "start_date": datetime(2023, 12, 22), "email_on_failure": False, "email_on_retry": False, "retries": 0, "retry_delay": timedelta(minutes=5), } def read_csv(**kwargs): file_path = kwargs.get("file_path") # 从上下文的params中获取运行时传入的target_id target_id = kwargs['params'].get("target_id", 0) df = pd.read_csv(file_path) if target_id != 0: filtered_df = df[df["id"] == target_id] print(f"Rows with id {target_id}:") print(filtered_df) else: # 用else替代重复if判断,逻辑更严谨 print("CSV file content:") print(df) with DAG( default_args=default_args, dag_id="csv_read_file", tags=["test"], schedule_interval="@daily", params={ "target_id": Param( 0, type="integer", minimum=0, ), }, ) as dag: csv_task = PythonOperator( task_id="read_csv", python_callable=read_csv, op_kwargs={ "file_path": "dags/test.csv", }, # 显式开启上下文传递,Airflow 2.x默认关闭 provide_context=True )
方式2:使用Airflow模板变量传递参数
直接在op_kwargs中用模板语法获取运行时参数,无需修改函数逻辑:
# 其余代码不变,仅修改PythonOperator的配置 csv_task = PythonOperator( task_id="read_csv", python_callable=read_csv, op_kwargs={ "file_path": "dags/test.csv", # 用模板变量读取运行时的target_id "target_id": "{{ params.target_id }}", }, # 开启模板渲染以解析变量 templates_dict={"target_id": "{{ params.target_id }}"} )
额外优化建议
- 将第二个独立的
if target_id == 0改为else,避免逻辑冗余,同时覆盖非预期值的分支处理。 - 如果CSV中的
id字段是字符串类型,需将target_id转换为字符串后再执行过滤匹配,避免类型不匹配导致的过滤失效。
内容的提问来源于stack exchange,提问作者Bedirhan Sahin
相关产品推荐
相关产品推荐

