如何将数据框中的数据插入数据库?
我的项目里有一个包含三个任务的DAG。这三个任务调用了三个函数,分别位于不同的文件中。我在使用ELT/ETL过程。看起来一切都没问题,直到数据被注入时出现问题。我不知道为什么无法把数据插入数据库。
我已经确认数据已被提取并存在于我的DataFrame中,因此这不是问题所在。
下面是loader.py文件中的函数:
import os
import pandas as pd
from dotenv import load_dotenv
from config.database import engine
load_dotenv()
def load_weather_data(data):
try:
print(f"\n=== DÉBOGAGE ===")
print(f"hourly_df - shape: {data['hourly_df'].shape}, vide: {data['hourly_df'].empty}")
print(f"daily_df - shape: {data['daily_df'].shape}, vide: {data['daily_df'].empty}")
data["hourly_df"].to_sql("hourly_weather", engine, if_exists="append", index=False)
if not data['hourly_df'].empty:
print(f"hourly_df colonnes: {list(data['hourly_df'].columns)}")
print(f"hourly_df première ligne:\n{data['hourly_df'].head(1)}")
print("✓ hourly_weather chargé avec succès")
# data["hourly_df"].to_sql("hourly_weather", engine, if_exists="append", index=False)
print("⚠️ hourly_df est vide, rien ne sera chargé")
if not data['daily_df'].empty:
print(f"daily_df colonnes: {list(data['daily_df'].columns)}")
print(f"daily_df première ligne:\n{data['daily_df'].head(1)}")
data["daily_df"].to_sql("daily_weather", engine, if_exists="append", index=False)
if not data['daily_df'].empty:
print("✓ daily_weather chargé avec succès")
else:
print("⚠️ daily_df est vide, rien ne sera chargé")
print("Chargement terminé\n")
except Exception as e:
print(f"✗ Erreur lors du chargement: {e}")
if __name__ == "__main__":
test_data = {"current_df": pd.DataFrame(), "hourly_df": pd.DataFrame(), "daily_df": pd.DataFrame()}
load_weather_data(test_data)
下面是我的Airflow DAG:
from airflow import DAG
from airflow.decorators import task
from datetime import datetime
import logging
logging.basicConfig(level=logging.INFO)
# from airflow.providers.postgres.hooks.postgres import PostgresHook
from src.extract.open_meteo import extract_weather_data
from src.transformations.open_meteo import transform_weather_data
from src.loaders.loader import load_weather_data
logger = logging.getLogger(__name__)
with DAG("weather_pipeline",schedule="@hourly", start_date=datetime(2026, 6, 16), catchup=False) as dag:
@task
def extract():
logger.info("Extracting weather data...")
return extract_weather_data()
@task
def transform(extracted_response):
logger.info("Transforming weather data...")
return transform_weather_data(extracted_response)
@task
def load(transformed_data):
logger.info("Loading weather datappppppmpmpmpmmpmpmmpmpmpm...")
return load_weather_data(transformed_data)
extracted_data = extract()
transformed_data = transform(extracted_data)
load(transformed_data)
下面是来自airflow_logs的其中一条日志:
[2026-06-23T14:43:15.709+0000] {local_task_job_runner.py:120} INFO - ::group::Pre task execution logs
[2026-06-23T14:43:15.736+0000] {taskinstance.py:2076} INFO - Dependencies all met for dep_context=non-requeueable deps ti=<TaskInstance: weather_pipeline.load manual__2026-06-23T14:43:10.772224+00:00 [queued]>
[2026-06-23T14:43:15.741+0000] {taskinstance.py:2076} INFO - Dependencies all met for dep_context=requeueable deps ti=<TaskInstance: weather_pipeline.load manual__2026-06-23T14:43:10.772224+00:00 [queued]>
[2026-06-23T14:43:15.742+0000] {taskinstance.py:2306} INFO - Starting attempt 1 of 1
[2026-06-23T14:43:15.749+0000] {taskinstance.py:2330} INFO - Executing <Task(_PythonDecoratedOperator): load> on 2026-06-23 14:43:10.772224+00:00
[2026-06-23T14:43:15.758+0000] {standard_task_runner.py:64} INFO - Started process 6056 to run task
[2026-06-23T14:43:15.761+0000] {standard_task_runner.py:90} INFO - Running: ['airflow', 'tasks', 'run', 'weather_pipeline', 'load', 'manual__2026-06-23T14:43:10.772224+00:00', '--job-id', '6', '--raw', '--subdir', 'DAGS_FOLDER/weather_pipeline.py', '--cfg-path', '/tmp/tmpsk6asf74']
[2026-06-23T14:43:15.763+0000] {standard_task_runner.py:91} INFO - Job 6: Subtask load
[2026-06-23T14:43:15.805+0000] {task_command.py:426} INFO - Running <TaskInstance: weather_pipeline.load manual__2026-06-23T14:43:10.772224+00:00 [running]> on host 536028be57d8
[2026-06-23T14:43:16.164+0000] {taskinstance.py:2648} INFO - Exporting env vars: AIRFLOW_CTX_DAG_OWNER='airflow' AIRFLOW_CTX_DAG_ID='weather_pipeline' AIRFLOW_CTX_TASK_ID='load' AIRFLOW_CTX_EXECUTION_DATE='2026-06-23T14:43:10.772224+00:00' AIRFLOW_CTX_TRY_NUMBER='1' AIRFLOW_CTX_DAG_RUN_ID='manual__2026-06-23T14:43:10.772224+00:00'
[2026-06-23T14:43:16.164+0000] {taskinstance.py:430} INFO - ::endgroup::
[2026-06-23T14:43:16.165+0000] {weather_pipeline.py:28} INFO - Loading weather datappppppmpmpmpmmpmpmmpmpmpm...
[2026-06-23T14:43:16.165+0000] {logging_mixin.py:188} INFO -
=== DÉBOGAGE ===
[2026-06-23T14:43:16.165+0000] {logging_mixin.py:188} INFO - hourly_df - shape: (912, 9), vide: False
[2026-06-23T14:43:16.165+0000] {logging_mixin.py:188} INFO - daily_df - shape: (38, 13), vide: False
[2026-06-23T14:43:16.176+0000] {logging_mixin.py:188} INFO - hourly_df colonnes: ['date', 'temperature_2m', 'relative_humidity_2m', 'dew_point_2m', 'apparent_temperature', 'precipitation_probability', 'rain', 'snowfall', 'weather_code']
[2026-06-23T14:43:16.182+0000] {logging_mixin.py:188} INFO - hourly_df première ligne:
date temperature_2m ... snowfall weather_code
0 2026-05-22 22:00:00+00:00 25.437 ... 0.0 0.0
[1 rows x 9 columns]
[2026-06-23T14:43:16.185+0000] {logging_mixin.py:188} WARNING - /opt/airflow/src/loaders/loader.py:20 UserWarning: pandas only supports SQLAlchemy connectable (engine/connection) or database string URI or sqlite3 DBAPI2 connection. Other DBAPI2 objects are not tested. Please consider using SQLAlchemy.
[2026-06-23T14:43:16.185+0000] {logging_mixin.py:188} INFO - ✗ Erreur lors du chargement: 'Connection' object has no attribute 'cursor'
[2026-06-23T14:43:16.185+0000] {python.py:237} INFO - Done. Returned value was: None
[2026-06-23T14:43:16.186+0000] {taskinstance.py:441} INFO - ::group::Post task execution logs
[2026-06-23T14:43:16.194+0000] {taskinstance.py:1206} INFO - Marking task as SUCCESS. dag_id=weather_pipeline, task_id=load, run_id=manual__2026-06-23T14:43:10.772224+00:00, execution_date=20260623T144310, start_date=20260623T144315, end_date=20260623T144316
[2026-06-23T14:43:16.253+0000] {local_task_job_runner.py:243} INFO - Task exited with return code 0
[2026-06-23T14:43:16.271+0000] {taskinstance.py:3503} INFO - 0 downstream tasks scheduled from follow-on schedule check
[2026-06-23T14:43:16.273+0000] {local_task_job_runner.py:222} INFO - ::endgroup::
我不明白为什么里面会有一个WARNING:
[2026-06-23T14:43:16.185+0000] {logging_mixin.py:188} WARNING - /opt/airflow/src/loaders/loader.py:20 UserWarning: pandas only supports SQLAlchemy connectable (engine/connection) or database string URI or sqlite3 DBAPI2 connection. Other DBAPI2 objects are not tested. Please consider using SQLAlchemy.
就像每次脚本被运行时,都会因此而停止处理。问题在于我在另外一个文件中使用SQLAlchemy来配置数据库,并在loader.py中调用了它
import os
from dotenv import load_dotenv
from sqlalchemy import create_engine
load_dotenv()
DATABASE_URL = (
f"postgresql+psycopg2://"
f"{os.getenv('POSTGRES_USER')}:{os.getenv('POSTGRES_PASSWORD')}"
f"@{os.getenv('POSTGRES_HOST')}:{os.getenv('POSTGRES_PORT')}"
f"/{os.getenv('POSTGRES_DB')}"
)
engine = create_engine(DATABASE_URL)
所以我不懂。
我在用Docker搭数据库(不确定这是否可能是问题原因,所以就先说到这儿)
解决方案
警告只是一个症状。真正的问题在这里:
'Connection' object has no attribute 'cursor'
pandas.to_sql() 收到一个它无法识别为有效的SQLAlchemy引擎/连接的对象,因此回退将其当作DBAPI连接来处理并尝试调用 .cursor()。SQLAlchemy Connection 对象没有 .cursor(),因此发生错误。
另外,你的Airflow任务被标记为 SUCCESS,因为你捕获了异常却没有重新抛出它:
except Exception as e:
print(f"✗ Erreur lors du chargement: {e}")
从Airflow的角度来看,这个函数已经正常完成。
请使用真正的SQLAlchemy引擎/连接并重新抛出异常:
from config.database import engine
def load_weather_data(data):
hourly_df = data["hourly_df"]
daily_df = data["daily_df"]
print("\n=== DEBUG ===")
print(f"hourly_df shape: {hourly_df.shape}, empty: {hourly_df.empty}")
print(f"daily_df shape: {daily_df.shape}, empty: {daily_df.empty}")
try:
with engine.begin() as conn:
if not hourly_df.empty:
print(f"hourly_df columns: {list(hourly_df.columns)}")
print(hourly_df.head(1))
hourly_df.to_sql(
"hourly_weather",
con=conn,
if_exists="append",
index=False,
method="multi"
)
print("hourly_weather loaded successfully")
if not daily_df.empty:
print(f"daily_df columns: {list(daily_df.columns)}")
print(daily_df.head(1))
daily_df.to_sql(
"daily_weather",
con=conn,
if_exists="append",
index=False,
method="multi"
)
print("daily_weather loaded successfully")
except Exception:
# Important: let Airflow mark the task as failed
raise
另外把函数顶部的这一行也删除:
data["hourly_df"].to_sql("hourly_weather", engine, if_exists="append", index=False)
现在你在检查是否为空之前就插入了 hourly_df,之后又有一条被注释掉的插入语句。请把逻辑放在同一个地方处理。
然后在Airflow容器中检查 engine 实际上是什么。它应该是一个SQLAlchemy引擎,例如:
<class 'sqlalchemy.engine.base.Engine'>
如果不是,那么你就导入或创建了错误的对象。
站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。