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

