You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何定位Apache Airflow中SQLite任务对应的数据库文件?

问题定位与解决方案

核心原因

你使用的SqliteOperator默认调用Airflow内置的sqlite_default连接,这个连接指向的并非Airflow元数据库airflow.db,而是默认路径下的sqlite_default.db文件(通常位于Airflow工作目录或你启动Airflow服务的目录中)。这就是任务执行成功但airflow.db中找不到目标表的根本原因。

一、定位实际执行的数据库文件

  1. 查看sqlite_default连接的具体配置:
    执行以下命令获取连接详情:

    airflow connections get sqlite_default
    

    输出中的conn_uri字段会显示该连接对应的SQLite文件路径,例如sqlite:///sqlite_default.db(相对路径,对应Airflow运行目录下的文件)。

  2. 验证数据库内容:
    用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文件执行任务,按以下步骤操作:

  1. 创建自定义连接:
    执行命令创建新的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字段填写数据库文件的绝对路径。

  2. 修改DAG中的SqliteOperator:
    在每个任务中添加conn_id='my_custom_sqlite'参数,运行DAG后,表会自动创建在my_custom.db中。

内容的提问来源于stack exchange,提问作者ademusire

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.03 07:01:06