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

Airflow仅运行最新DAG及增量ETL DAG配置技术问询

Hey there! Let's tackle your two Airflow questions one by one, with practical, actionable solutions tailored to your needs.


1. How to Run Only the Latest Airflow DAG Instance?

I’m guessing you want to avoid multiple old DAG runs cluttering your system or competing for resources while the latest instance executes. Here are a few reliable approaches:

Option 1: Limit Concurrent Runs with Core Configuration

Add max_active_runs=1 to your DAG definition. This tells Airflow to only allow one instance of the DAG to run at a time. If a new scheduled run comes up while an old one is still going, the new run will queue until the old one finishes. It’s simple and works for most basic use cases.

Option 2: Auto-Terminate Old Runs When a New One Starts

If you need to immediately cancel incomplete old runs when a new instance launches, use a custom PythonOperator to clean up prior runs at the start of your DAG:

from airflow.models import DagRun
from airflow.utils.state import State
from airflow.operators.python import PythonOperator

def cancel_stale_runs(**context):
    dag_id = context["dag"].dag_id
    current_run_id = context["run_id"]
    
    # Fetch all running/queued runs except the current one
    stale_runs = DagRun.find(dag_id=dag_id, state=[State.RUNNING, State.QUEUED])
    for run in stale_runs:
        if run.run_id != current_run_id:
            run.set_state(State.CANCELLED)
            print(f"Cancelled stale run: {run.run_id}")

# Add this task at the start of your DAG
cancel_stale_task = PythonOperator(
    task_id="cancel_stale_runs",
    python_callable=cancel_stale_runs,
    provide_context=True,
)

Option 3: Trigger the Latest Scheduled Run via CLI

For manual triggers, use this command to launch only the next scheduled instance of your DAG:

airflow dags trigger --run-id $(airflow dags next-execution your_dag_id) your_dag_id

2. Deploying a Hourly ETL DAG with Initial Large Data Loads

Your scenario—handling big backfilled data before settling into regular hourly runs—is super common. The key is to tie your data extraction range to the latest timestamp in DB2, rather than rigidly following the DAG’s scheduled time. Here’s how to build it:

Step 1: Core DAG Configuration

Start with settings that accommodate long initial runs:

from airflow import DAG
from datetime import datetime, timedelta

default_args = {
    'owner': 'your_team',
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
    'execution_timeout': timedelta(hours=72)  # Allow long initial runs to complete
}

with DAG(
    dag_id="db1_to_db2_etl",
    default_args=default_args,
    schedule_interval="@hourly",
    start_date=datetime(2024, 1, 1),
    catchup=False,  # We'll handle backfills via our custom logic
    max_active_runs=1,  # Prevent concurrent runs during long initial loads
) as dag:

Step 2: Define the ETL Pipeline

Break the workflow into three focused tasks:

1. Fetch the Latest Timestamp from DB2

First, get the last time data was inserted into DB2 (or use a dedicated update timestamp field). If the table is empty (first run), set a starting point (e.g., 1 month ago):

from airflow.operators.python import PythonOperator
import ibm_db  # Use DB2's official driver

def get_latest_db2_timestamp(**context):
    conn = ibm_db.connect("YOUR_DB2_CONN_STRING", "", "")
    query = "SELECT MAX(insert_timestamp) FROM your_target_table"
    stmt = ibm_db.exec_immediate(conn, query)
    latest_ts = ibm_db.fetch_tuple(stmt)[0]

    # Set fallback for first run
    if not latest_ts:
        latest_ts = datetime.now() - timedelta(days=30)
    
    # Pass timestamp to next task via XCom
    context['ti'].xcom_push(
        key='latest_db2_ts', 
        value=latest_ts.strftime('%Y-%m-%d %H:%M:%S')
    )
    ibm_db.close(conn)

get_latest_ts_task = PythonOperator(
    task_id="get_latest_db2_timestamp",
    python_callable=get_latest_db2_timestamp,
    provide_context=True,
)

2. Extract & Load Data from DB1 to DB2

Use the timestamp from the previous task to pull only new data from DB1, then load it into DB2:

def extract_load_data(**context):
    latest_ts = context['ti'].xcom_pull(
        key='latest_db2_ts', 
        task_ids='get_latest_db2_timestamp'
    )
    latest_ts_dt = datetime.strptime(latest_ts, '%Y-%m-%d %H:%M:%S')

    # Extract from DB1 (use your DB1's driver, e.g., psycopg2 for PostgreSQL)
    import psycopg2
    db1_conn = psycopg2.connect("YOUR_DB1_CONN_STRING")
    cursor = db1_conn.cursor()
    extract_query = f"""
        SELECT col1, col2, update_timestamp 
        FROM your_source_table 
        WHERE update_timestamp > '{latest_ts}'
    """
    cursor.execute(extract_query)
    data = cursor.fetchall()

    # Load to DB2
    db2_conn = ibm_db.connect("YOUR_DB2_CONN_STRING", "", "")
    insert_query = """
        INSERT INTO your_target_table (col1, col2, insert_timestamp)
        VALUES (?, ?, ?)
    """
    stmt = ibm_db.prepare(db2_conn, insert_query)
    
    for row in data:
        ibm_db.execute(stmt, (row[0], row[1], row[2]))
    
    ibm_db.commit(db2_conn)
    # Clean up connections
    ibm_db.close(db2_conn)
    db1_conn.close()

etl_task = PythonOperator(
    task_id="extract_load_db1_to_db2",
    python_callable=extract_load_data,
    provide_context=True,
)

3. Wire Up the Tasks

Link the tasks together to form the full pipeline:

get_latest_ts_task >> etl_task

Step 3: How It Handles Your Specific Workflow

  • First run: DB2 has no data, so we pull 1 month of data from DB1 and load it (takes ~72 hours)
  • Second run: Fetches the latest timestamp from the first load, pulls DB1 data from that point onward (~72 hours of data, takes ~7.2 hours)
  • Third run: Repeats the logic, pulling only the last ~7.2 hours of data (finishes within 1 hour)
  • Subsequent runs: Hourly pulls of the last hour’s data, completing on schedule

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 08:45:22