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

Airflow每日调度DAG如何实现单日数据库数据提取?

解决方案:按Airflow运行日期提取单日数据

核心思路

放弃原有的startfrom偏移量追踪逻辑,直接利用Airflow任务的逻辑运行日期作为过滤条件,将日期参数传递给数据提取脚本,让脚本仅拉取对应单日的数据。

具体修改步骤

1. 调整Airflow DAG代码

修改PythonOperator以传递运行日期,移除不必要的startfrom变量维护逻辑:

''' task : extraction des données journalières '''

import os
from airflow.operators.python_operator import PythonOperator
import logging
import pendulum
from datetime import datetime, timedelta
from airflow import DAG
import subprocess

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2023, 12, 11),
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

dag = DAG(
    'daily_data_extraction',
    default_args=default_args,
    description='Run a Python script every day at 6:00 AM to extract daily data',
    schedule_interval='0 6 * * *',
    # 指定时区,避免日期偏移问题
    timezone=pendulum.timezone("Europe/Paris"), 
)

def run_my_script(**context):
    # 获取当前任务的逻辑运行日期(Airflow 2.x推荐使用)
    execution_date = context['logical_date']
    # 格式化为YYYY-MM-DD字符串,适配脚本参数格式
    target_date = execution_date.strftime("%Y-%m-%d")
    
    script_path = "billetiques/script_1.py"
    
    # 将日期作为参数传递给提取脚本
    result = subprocess.run(
        ["python", script_path, "--date", target_date], 
        capture_output=True, 
        text=True
    )
    
    # 错误处理:脚本执行失败时抛出异常,触发Airflow任务重试
    if result.returncode != 0:
        logging.error(f"Extraction script failed: {result.stderr}")
        raise Exception("Daily data extraction failed")
    logging.info(f"Extraction completed: {result.stdout}")

run_script_task = PythonOperator(
    task_id='run_daily_extraction',
    python_callable=run_my_script,
    # 允许函数接收Airflow上下文参数,获取运行日期
    provide_context=True,
    dag=dag,
)

2. 修改数据提取脚本script_1.py

更新脚本以接收--date参数,并在数据库查询中添加日期过滤条件:

import argparse
# 根据你的数据库类型替换对应的连接库,示例用PostgreSQL
import psycopg2

def extract_daily_data(target_date):
    # 建立数据库连接
    conn = psycopg2.connect(
        dbname="your_database",
        user="your_user",
        password="your_password",
        host="your_db_host"
    )
    cursor = conn.cursor()
    
    # 编写带日期过滤的查询语句(假设表中有create_date字段存储数据生成日期)
    query = """
        SELECT * FROM your_target_table
        WHERE DATE(create_date) = %s
    """
    
    cursor.execute(query, (target_date,))
    daily_data = cursor.fetchall()
    
    # 后续数据处理逻辑(如写入文件、同步到数据仓库等)
    # ...
    
    cursor.close()
    conn.close()

if __name__ == "__main__":
    parser = argparse.ArgumentParser(description='Extract daily business data')
    parser.add_argument('--date', required=True, help='Target date in YYYY-MM-DD format')
    args = parser.parse_args()
    
    extract_daily_data(args.date)

关键细节说明

  • 时区匹配:DAG中指定的timezone要和业务数据的时区一致,避免出现日期跨天的错误。
  • 日期调整:如果6:00运行的任务需要提取当天的数据,可将target_date改为(execution_date + timedelta(days=1)).strftime("%Y-%m-%d"),根据实际业务需求调整。
  • 任务可靠性:添加脚本执行失败的异常抛出,确保Airflow能正确识别任务状态并触发重试机制。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 06:17:34