如何在Airflow数据感知调度中使用SQL Server表作为Dataset
能否将SQL Server表用作Airflow Dataset?
当然可以将SQL Server表作为Airflow的Dataset实现数据感知调度。核心是构造能唯一标识目标表的URI,并在生产者任务中标记Dataset更新,触发消费者DAG运行。
核心配置与代码示例
1. 定义SQL Server表对应的Dataset
Dataset的URI需明确指向目标SQL Server表,可基于SQLAlchemy连接字符串扩展,附加表名作为唯一标识:
from airflow import Dataset # 格式:mssql+pyodbc://<用户名>:<密码>@<服务器地址>/<数据库名>?driver=<ODBC驱动>#<模式名>.<表名> sql_server_dataset = Dataset( uri="mssql+pyodbc://db_user:db_pass@sql_server_01/my_database?driver=ODBC+Driver+17+for+SQL+Server#dbo.target_table" )
2. 消费者DAG(监听Dataset更新)
将上述Dataset设为调度触发器,当Dataset被标记更新时,该DAG自动启动:
from airflow import DAG from airflow.operators.empty import EmptyOperator from datetime import datetime with DAG( dag_id="sql_server_dataset_consumer", schedule=[sql_server_dataset], # 监听Dataset更新事件 start_date=datetime(2024, 1, 1), catchup=False ): start_processing = EmptyOperator(task_id="start_processing") # 可在此添加处理SQL Server表数据的任务(如SqlOperator读取数据)
3. 生产者DAG(标记Dataset更新)
在修改SQL Server表的任务完成后,必须标记Dataset更新,才能触发消费者DAG:
from airflow import DAG from airflow.providers.microsoft.mssql.operators.mssql import MsSqlOperator from datetime import datetime with DAG( dag_id="sql_server_dataset_producer", schedule="@daily", start_date=datetime(2024, 1, 1), catchup=False ): # 执行修改SQL Server表的任务(如插入/更新数据) update_table = MsSqlOperator( task_id="update_target_table", mssql_conn_id="sql_server_default", # 提前在Airflow配置的SQL Server连接ID sql=""" INSERT INTO dbo.target_table (col1, col2) VALUES ('value1', 'value2'); """ ) # 标记Dataset更新,触发消费者DAG update_table >> sql_server_dataset
注意事项
- 确保已安装
apache-airflow-providers-microsoft-mssql包,用于SQL Server相关操作 - 提前在Airflow连接管理中配置好SQL Server连接(对应
mssql_conn_id) - Dataset的URI需全局唯一,避免不同表的标识冲突
- 生产者任务必须显式标记Dataset更新,否则消费者DAG不会被触发
内容的提问来源于stack exchange,提问作者Tomato Master
相关产品推荐
相关产品推荐

