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

Airflow DAG每小时调度异常:加载未来5天数据而非实时更新求助

解决方案

问题根源

  • API接口特性:调用的forecast接口默认返回未来5天、每3小时一次的全量预报数据,并非每小时的增量更新数据。
  • 重复数据未处理:脚本未检查数据库中是否已存在相同city+date_time的记录,首次插入后,后续重复插入会因唯一约束(若表中隐含该约束)失败,导致数据不再更新;若表无约束,则会插入大量冗余数据。
  • 任务逻辑冗余:fetch_weather_task的结果未被insert_forecast_task利用,反而insert_forecast_into_db函数内重复调用接口,造成资源浪费。
  • 缺失必要导入:代码使用了datetime和timedelta但未导入,会导致运行报错。

具体修复步骤

1. 补充缺失的导入

在代码顶部添加:

from datetime import datetime, timedelta

2. 处理重复数据(保留预报场景)

如果需要继续获取5天预报数据,避免重复插入或实现数据更新:

  • 先给数据库表添加唯一约束:
ALTER TABLE weather_forecast ADD CONSTRAINT unique_city_datetime UNIQUE (city, date_time);
  • 修改插入语句,使用INSERT ... ON CONFLICT逻辑,存在则更新,不存在则插入:
cursor.execute("""
    INSERT INTO weather_forecast (city, temperature, humidity, description, date_time)
    VALUES (%s, %s, %s, %s, %s)
    ON CONFLICT (city, date_time) DO UPDATE 
    SET temperature = EXCLUDED.temperature,
        humidity = EXCLUDED.humidity,
        description = EXCLUDED.description;
""", (
    city,
    forecast['main']['temp'],
    forecast['main']['humidity'],
    forecast['weather'][0]['description'],
    date_time
))

3. 切换到实时天气接口(增量更新场景)

如果需求是每小时获取当前实时天气而非预报,改用OpenWeather的weather接口:

def fetch_current_weather(api_key, city):
    url = f'http://api.openweathermap.org/data/2.5/weather?q={city}&appid={api_key}&units=metric'
    response = requests.get(url)
    if response.status_code == 200:
        return response.json()
    else:
        print(f'Failed to fetch data: {response.status_code}')
        return None

同时修改插入逻辑,获取实时数据中的对应字段:

current_data = fetch_current_weather(api_key, city)
if current_data:
    date_time = datetime.now()  # 或使用接口返回的dt字段转时间
    cursor.execute("""
        INSERT INTO weather_forecast (city, temperature, humidity, description, date_time)
        VALUES (%s, %s, %s, %s, %s)
        ON CONFLICT (city, date_time) DO NOTHING;  # 避免同一小时重复插入
    """, (
        city,
        current_data['main']['temp'],
        current_data['main']['humidity'],
        current_data['weather'][0]['description'],
        date_time
    ))

4. 优化Airflow任务逻辑

移除冗余的fetch_weather_task,直接在插入任务中调用接口,避免重复请求:

# 移除原有的fetch_weather_task循环,只保留插入任务
for city in cities:
    insert_forecast_task = PythonOperator(
        task_id=f'insert_forecast_into_postgres_{city.lower().replace(" ", "_")}',
        python_callable=insert_forecast_into_db,
        op_kwargs={'api_key': api_key, 'city': city},
        dag=dag,
    )

完整修复后的代码示例

import requests
import psycopg2
from psycopg2 import Error
from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from datetime import datetime, timedelta

# Define the DAG parameters
default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 6, 28),
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

dag = DAG(
    'metro_cities',
    default_args=default_args,
    description='Fetch weather forecast data from OpenWeatherMap API and insert into PostgreSQL',
    schedule_interval=timedelta(hours=1),
)

# Function to fetch weather forecast data from OpenWeather API
def fetch_weather_forecast(api_key, city):
    url = f'http://api.openweathermap.org/data/2.5/forecast?q={city}&appid={api_key}&units=metric'
    response = requests.get(url)
    if response.status_code == 200:
        return response.json()
    else:
        print(f'Failed to fetch data: {response.status_code}')
        return None

# Function to insert forecast data into PostgreSQL
def insert_forecast_into_db(api_key, city, **kwargs):
    connection = None
    cursor = None
    try:
        connection = psycopg2.connect(
            user="postgres",
            password="password",
            host="localhost",
            port="5432",
            database="postgres"
        )
        cursor = connection.cursor()

        forecast_data = fetch_weather_forecast(api_key, city)
        if forecast_data:
            forecast_list = forecast_data.get('list', [])
            for forecast in forecast_list:
                date_time = datetime.fromtimestamp(forecast['dt'])
                
                # 处理重复数据,存在则更新
                cursor.execute("""
                    INSERT INTO weather_forecast (city, temperature, humidity, description, date_time)
                    VALUES (%s, %s, %s, %s, %s)
                    ON CONFLICT (city, date_time) DO UPDATE 
                    SET temperature = EXCLUDED.temperature,
                        humidity = EXCLUDED.humidity,
                        description = EXCLUDED.description;
                    """, (
                        city,
                        forecast['main']['temp'],
                        forecast['main']['humidity'],
                        forecast['weather'][0]['description'],
                        date_time
                    ))

            connection.commit()
            print("Forecast data inserted/updated successfully")
    except (Exception, Error) as error:
        print("Error while connecting to PostgreSQL", error)
        if connection:
            connection.rollback()
    finally:
        if cursor:
            cursor.close()
        if connection:
            connection.close()
            print("PostgreSQL connection is closed")

# Define cities for forecast
cities = ['Seattle', 'Oakland', 'Atlanta']
api_key ='api'

# Define tasks in the DAG for each city
for city in cities:
    insert_forecast_task = PythonOperator(
        task_id=f'insert_forecast_into_postgres_{city.lower().replace(" ", "_")}',
        python_callable=insert_forecast_into_db,
        op_kwargs={'api_key': api_key, 'city': city},
        dag=dag,
    )

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 20:54:54