如何定位Apache Airflow中SQLite任务对应的数据库文件?
核心原因
你使用的SqliteOperator默认调用Airflow内置的sqlite_default连接,这个连接指向的并非Airflow元数据库airflow.db,而是默认路径下的sqlite_default.db文件(通常位于Airflow工作目录或你启动Airflow服务的目录中)。这就是任务执行成功但airflow.db中找不到目标表的根本原因。
一、定位实际执行的数据库文件
查看
sqlite_default连接的具体配置:
执行以下命令获取连接详情:airflow connections get sqlite_default输出中的
conn_uri字段会显示该连接对应的SQLite文件路径,例如sqlite:///sqlite_default.db(相对路径,对应Airflow运行目录下的文件)。验证数据库内容:
用SQLite命令行工具打开该文件,检查是否存在创建的表:sqlite3 /path/to/sqlite_default.db在SQLite交互界面输入
.tables,即可看到tripdata_monthly_statistics表。
二、修改DAG,指定目标数据库
方法1:让任务使用Airflow元数据库(airflow.db)
在每个SqliteOperator中添加conn_id='airflow_db'参数,这个默认连接直接指向你的airflow.db。修改后的DAG示例:
from datetime import datetime, timedelta from airflow import DAG from airflow.providers.sqlite.operators.sqlite import SqliteOperator default_args = { "owner": "ademusire", "retries": 0, "retry_delay": timedelta(minutes=2) } with DAG( dag_id="dag_with_sqlite_operator_v06", default_args=default_args, start_date=datetime(2023, 12, 23), schedule_interval="@daily" ) as dag: task1 = SqliteOperator( task_id="create_table_sqlite", conn_id='airflow_db', # 添加指定连接 sql=r""" CREATE TABLE IF NOT EXISTS tripdata_monthly_statistics( id INTEGER PRIMARY KEY AUTOINCREMENT, month TEXT, sat_mean_trip_count NUMERIC, sat_mean_fare_per_trip NUMERIC, sat_mean_duration_per_trip NUMERIC, sun_mean_trip_count NUMERIC, sun_mean_fare_per_trip NUMERIC, sun_mean_duration_per_trip NUMERIC ); """, ) task2 = SqliteOperator( task_id="insert_into_table", conn_id='airflow_db', # 添加指定连接 sql=r""" INSERT INTO tripdata_monthly_statistics(id, month, sat_mean_trip_count, sat_mean_fare_per_trip, sat_mean_duration_per_trip, sun_mean_trip_count, sun_mean_fare_per_trip, sun_mean_duration_per_trip) VALUES(1, '2023-11', 7, 8, 9, 10, 11, 12); """, ) task3 = SqliteOperator( task_id="select_from_table", conn_id='airflow_db', # 添加指定连接 sql=r"""SELECT * FROM tripdata_monthly_statistics;""", ) task4 = SqliteOperator( task_id="show_tables", conn_id='airflow_db', # 添加指定连接 sql=r""" SELECT name FROM sqlite_schema WHERE type ='table' AND name NOT LIKE 'sqlite_%'; """, ) task1 >> task2 >> task3 >> task4
重新运行DAG后,即可在/home/ademusire/airflow/airflow.db中找到创建的表。
方法2:创建自定义SQLite数据库
若想单独使用新的SQLite文件执行任务,按以下步骤操作:
创建自定义连接:
执行命令创建新的SQLite连接,指定自定义数据库文件的绝对路径:airflow connections add 'my_custom_sqlite' \ --conn-type 'sqlite' \ --conn-host '/home/ademusire/airflow/my_custom.db'也可通过Airflow UI的
Admin -> Connections页面手动添加:选择Conn Type为SQLite,Host字段填写数据库文件的绝对路径。修改DAG中的
SqliteOperator:
在每个任务中添加conn_id='my_custom_sqlite'参数,运行DAG后,表会自动创建在my_custom.db中。
内容的提问来源于stack exchange,提问作者ademusire

