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

