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
相关产品推荐
相关产品推荐

