6 minutes
数据加载(Load):目标存储与写入
数据加载概述
数据加载是 ETL 流程的最后一公里。经过抽取和转换的数据,最终需要写入目标存储系统,供分析师、数据科学家和业务系统使用。
加载环节看似简单——不就是把数据写进去吗?但实际上,加载策略的选择直接影响数据的一致性和可用性。是全量刷新还是增量更新?用 INSERT 还是 COPY?数据量大时如何保证写入性能?加载过程中出错了怎么办?这些问题都必须在设计阶段想清楚。
加载策略
全量加载
全量加载每次都将整个数据集写入目标表。它最简单直接,适合数据量不大或者需要完全重建分析视图的场景。
-- 全量加载:先清空目标表,再写入全部数据
TRUNCATE TABLE dwd_sales_daily;
INSERT INTO dwd_sales_daily
SELECT * FROM clean_sales_data;
优点:逻辑简单、数据完全一致、不需要处理变更追踪。
缺点:数据量大时耗时很长、消耗大量资源、加载期间数据不可用。
增量加载
增量加载只写入自上次加载以来发生变化的数据。它高效但实现复杂,需要可靠的变更追踪机制。
import psycopg2
import pandas as pd
from datetime import datetime, timedelta
def incremental_load(source_conn_str, target_conn_str, table_name, date_column):
"""增量加载:只加载最近一天的数据"""
yesterday = (datetime.now() - timedelta(days=1)).strftime("%Y-%m-%d")
# 从源系统抽取增量数据
source_conn = psycopg2.connect(source_conn_str)
query = f"SELECT * FROM {table_name} WHERE {date_column} >= '{yesterday}'"
df = pd.read_sql(query, source_conn)
source_conn.close()
print(f"抽取到 {len(df)} 条增量记录")
if df.empty:
print("没有增量数据,跳过加载")
return
# 写入目标系统
target_conn = psycopg2.connect(target_conn_str)
cursor = target_conn.cursor()
for _, row in df.iterrows():
cursor.execute(
f"""
INSERT INTO {table_name} ({', '.join(df.columns)})
VALUES ({', '.join(['%s'] * len(row))})
ON CONFLICT (id) DO UPDATE SET
{', '.join([f"{col} = EXCLUDED.{col}" for col in df.columns if col != 'id'])}
""",
list(row.values)
)
target_conn.commit()
cursor.close()
target_conn.close()
print(f"成功写入 {len(df)} 条记录到目标表")
UPSERT(MERGE)策略
UPSERT 是增量加载的核心模式——如果记录已存在则更新,不存在则插入。几乎所有主流数据库都支持这种操作。
-- PostgreSQL 的 UPSERT(INSERT ... ON CONFLICT)
INSERT INTO dim_customer (customer_id, name, email, tier, updated_at)
VALUES
(1001, '张三', 'zhangsan@example.com', 'gold', CURRENT_TIMESTAMP),
(1002, '李四', 'lisi@example.com', 'silver', CURRENT_TIMESTAMP)
ON CONFLICT (customer_id) DO UPDATE SET
name = EXCLUDED.name,
email = EXCLUDED.email,
tier = EXCLUDED.tier,
updated_at = CURRENT_TIMESTAMP;
-- Snowflake 的 MERGE
MERGE INTO dim_customer AS target
USING (
SELECT 1001 AS customer_id, '张三' AS name, 'gold' AS tier
UNION ALL
SELECT 1002, '李四', 'silver'
) AS source
ON target.customer_id = source.customer_id
WHEN MATCHED THEN UPDATE SET
target.name = source.name,
target.tier = source.tier,
target.updated_at = CURRENT_TIMESTAMP
WHEN NOT MATCHED THEN INSERT (customer_id, name, tier, updated_at)
VALUES (source.customer_id, source.name, source.tier, CURRENT_TIMESTAMP);
-- BigQuery 的 MERGE
MERGE INTO `project.dataset.dim_customer` AS target
USING (
SELECT 1001 AS customer_id, '张三' AS name, 'gold' AS tier
UNION ALL
SELECT 1002, '李四', 'silver'
) AS source
ON target.customer_id = source.customer_id
WHEN MATCHED THEN
UPDATE SET name = source.name, tier = source.tier
WHEN NOT MATCHED THEN
INSERT (customer_id, name, tier) VALUES (source.customer_id, source.name, source.tier);
目标存储系统
数据仓库加载
数据仓库是 ETL 最经典的目标系统。以 PostgreSQL 为例:
import psycopg2
from psycopg2.extras import execute_values
def load_to_postgresql(df, table_name, connection_string):
"""批量加载 DataFrame 到 PostgreSQL"""
conn = psycopg2.connect(connection_string)
cursor = conn.cursor()
# 使用 execute_values 进行高效批量插入
columns = list(df.columns)
values = [list(row) for row in df.itertuples(index=False)]
insert_query = f"""
INSERT INTO {table_name} ({', '.join(columns)})
VALUES %s
"""
try:
execute_values(cursor, insert_query, values)
conn.commit()
print(f"成功加载 {len(df)} 行到 {table_name}")
except Exception as e:
conn.rollback()
print(f"加载失败: {e}")
raise
finally:
cursor.close()
conn.close()
Snowflake 加载
Snowflake 原生支持从 Parquet/CSV 文件批量加载,也支持通过 Python SDK 写入。
import snowflake.connector
from snowflake.connector.pandas_tools import write_pandas
def load_to_snowflake(df, table_name, connection_params):
"""使用 write_pandas 加载数据到 Snowflake"""
conn = snowflake.connector.connect(**connection_params)
try:
success, num_chunks, num_rows, output = write_pandas(
conn=conn,
df=df,
table_name=table_name.upper(),
quote_identifiers=False
)
if success:
print(f"成功加载 {num_rows} 行到 Snowflake 表 {table_name}")
else:
print(f"加载可能不完整: {output}")
finally:
conn.close()
数据湖(Parquet 格式)加载
对于数据湖架构,数据加载的目标不是数据库表,而是文件系统中的结构化文件。
import pandas as pd
import pyarrow as pa
import pyarrow.parquet as pq
from pathlib import Path
from datetime import datetime
def load_to_parquet_lake(df, base_path, table_name, partition_cols=None):
"""
将 DataFrame 以 Parquet 格式写入数据湖
支持按列分区存储(如按日期分区)
"""
# 创建表目录
table_path = Path(base_path) / table_name
table_path.mkdir(parents=True, exist_ok=True)
if partition_cols:
# 按分区列逐组写入
for keys, group_df in df.groupby(partition_cols):
if isinstance(keys, tuple):
partition_path = table_path / "/".join(
[f"{col}={val}" for col, val in zip(partition_cols, keys)]
)
else:
partition_path = table_path / f"{partition_cols[0]}={keys}"
partition_path.mkdir(parents=True, exist_ok=True)
# 写入 Parquet 文件
file_name = f"data_{datetime.now().strftime('%H%M%S')}.parquet"
file_path = partition_path / file_name
table = pa.Table.from_pandas(group_df)
pq.write_table(table, file_path)
print(f"写入 {file_path},{len(group_df)} 行")
else:
# 不分区,直接写入
file_name = f"data_{datetime.now().strftime('%Y%m%d_%H%M%S')}.parquet"
file_path = table_path / file_name
table = pa.Table.from_pandas(df)
pq.write_table(table, file_path)
print(f"写入 {file_path},{len(df)} 行")
print(f"数据加载完成,目录: {table_path}")
# 使用示例
load_to_parquet_lake(
df=transformed_data,
base_path="/data/lake/",
table_name="sales",
partition_cols=["year", "month"]
)
加载性能优化
批量处理
逐行 INSERT 的性能非常差。批量提交可以将性能提升一到两个数量级。
def batch_load(df, conn, table_name, batch_size=5000):
"""批量加载,控制每批的大小"""
cursor = conn.cursor()
columns = list(df.columns)
placeholders = ", ".join(["%s"] * len(columns))
column_names = ", ".join(columns)
insert_sql = f"INSERT INTO {table_name} ({column_names}) VALUES ({placeholders})"
rows_processed = 0
for start in range(0, len(df), batch_size):
batch = df.iloc[start:start + batch_size]
values = [tuple(row) for row in batch.itertuples(index=False)]
try:
cursor.executemany(insert_sql, values)
conn.commit()
rows_processed += len(batch)
print(f"已加载 {rows_processed}/{len(df)} 行")
except Exception as e:
conn.rollback()
print(f"批次加载失败(起始行 {start}): {e}")
raise
cursor.close()
COPY 命令
数据库的 COPY 命令是加载 CSV 文件最快的方式,比 INSERT 快 5 到 10 倍。
# PostgreSQL COPY 命令
psql -h localhost -U etl_user -d dw -c "\
COPY dwd_sales FROM '/data/export/sales_20260407.csv' \
WITH (FORMAT CSV, HEADER true, DELIMITER ',') \
"
# MySQL LOAD DATA
mysql -h localhost -u etl_user -p dw \
-e "LOAD DATA LOCAL INFILE '/data/export/sales_20260407.csv' \
INTO TABLE dwd_sales \
FIELDS TERMINATED BY ',' \
ENCLOSED BY '\"' \
LINES TERMINATED BY '\n' \
IGNORE 1 ROWS"
# Snowflake COPY INTO
# 先上传到内部阶段,再用 COPY 加载
COPY INTO dwd_sales
FROM @%dwd_sales/sales_20260407.csv
FILE_FORMAT = (TYPE = CSV FIELD_OPTIONALLY_ENCLOSED_BY = '"')
ON_ERROR = 'CONTINUE';
# 用 Python 执行 PostgreSQL COPY
import io
import csv
def copy_dataframe_to_pg(df, table_name, conn):
"""使用 COPY 从 DataFrame 加载数据,性能最佳"""
buffer = io.StringIO()
df.to_csv(buffer, index=False, header=False, quoting=csv.QUOTE_MINIMAL)
buffer.seek(0)
cursor = conn.cursor()
try:
cursor.copy_expert(
f"COPY {table_name} FROM STDIN WITH CSV",
buffer
)
conn.commit()
print(f"COPY 加载完成: {len(df)} 行")
except Exception as e:
conn.rollback()
print(f"COPY 失败: {e}")
raise
finally:
cursor.close()
并行加载
对于非常大的数据集,可以并行写入多个文件或分片。
from concurrent.futures import ThreadPoolExecutor, as_completed
def parallel_load_to_parquet(df, base_path, num_partitions=4):
"""将 DataFrame 分成多个分区并行写入"""
# 将数据分成多个块
chunks = np.array_split(df, num_partitions)
def write_chunk(chunk, index):
file_path = f"{base_path}/part_{index:04d}.parquet"
table = pa.Table.from_pandas(chunk)
pq.write_table(table, file_path)
return f"Part {index}: {len(chunk)} 行 -> {file_path}"
with ThreadPoolExecutor(max_workers=num_partitions) as executor:
futures = [
executor.submit(write_chunk, chunk, i)
for i, chunk in enumerate(chunks)
]
for future in as_completed(futures):
print(future.result())
加载失败处理
数据加载不可靠是常态,而不是例外。一个生产级的加载流程必须有容错机制。
class EtlLoadManager:
"""管理 ETL 加载的事务和断点续传"""
def __init__(self, conn, load_table, checkpoint_table="load_checkpoints"):
self.conn = conn
self.load_table = load_table
self.checkpoint_table = checkpoint_table
self._create_checkpoint_table()
def _create_checkpoint_table(self):
cursor = self.conn.cursor()
cursor.execute(f"""
CREATE TABLE IF NOT EXISTS {self.checkpoint_table} (
load_id SERIAL PRIMARY KEY,
target_table VARCHAR(255),
batch_id VARCHAR(100),
start_time TIMESTAMP,
end_time TIMESTAMP,
rows_loaded INTEGER,
status VARCHAR(20)
)
""")
self.conn.commit()
cursor.close()
def load_with_checkpoint(self, df, batch_size=5000):
"""带检查点的增量加载,失败后可恢复"""
import uuid
load_id = str(uuid.uuid4())[:8]
cursor = self.conn.cursor()
for start in range(0, len(df), batch_size):
batch = df.iloc[start:start + batch_size]
batch_id = f"{load_id}_{start}"
# 检查这个批次是否已经加载成功
cursor.execute(
f"SELECT status FROM {self.checkpoint_table} "
f"WHERE target_table = %s AND batch_id = %s",
(self.load_table, batch_id)
)
existing = cursor.fetchone()
if existing and existing[0] == "completed":
print(f"批次 {batch_id} 已加载,跳过")
continue
# 执行加载
try:
cursor.execute(
f"INSERT INTO " + self.checkpoint_table +
" (target_table, batch_id, start_time, status) "
"VALUES (%s, %s, CURRENT_TIMESTAMP, 'running')",
(self.load_table, batch_id)
)
# ... 执行实际的 INSERT 或 COPY ...
cursor.execute(
f"UPDATE {self.checkpoint_table} SET status = 'completed', "
"end_time = CURRENT_TIMESTAMP, rows_loaded = %s "
"WHERE batch_id = %s",
(len(batch), batch_id)
)
self.conn.commit()
print(f"批次 {batch_id} 加载成功")
except Exception as e:
self.conn.rollback()
print(f"批次 {batch_id} 加载失败: {e}")
# 记录失败状态
cursor.execute(
f"UPDATE {self.checkpoint_table} SET status = 'failed' "
"WHERE batch_id = %s",
(batch_id,)
)
self.conn.commit()
raise
cursor.close()
加载策略选择决策表
| 场景 | 推荐策略 | 理由 |
|---|---|---|
| 小表(< 100 万行),每天一次 | 全量加载 | 简单可靠,写入时间短 |
| 大表,有时间戳字段 | 增量加载(时间戳) | 效率高,数据量可控 |
| 大表,无时间戳但有主键 | 增量加载(CDC) | 需要变更捕获机制 |
| 实时性要求高 | 流式加载(Kafka + Flink) | 端到端秒级延迟 |
| 数据量极大(TB 级) | COPY + 并行分片 | 充分利用并行 IO |
| 目标为数据湖 | Parquet 分区写入 | 列式压缩,查询优化 |
小结
数据加载是 ETL 流程的收尾环节,直接影响数据的使用体验和下游系统的性能。本文介绍了全量加载、增量加载和 UPSERT 三种主要加载策略,展示了针对 PostgreSQL、Snowflake 和数据湖(Parquet)的加载实现。性能优化方面,我们讨论了批量处理、COPY 命令和并行加载三种手段。最后,通过带检查点的加载管理器,我们实现了失败恢复的容错机制。记住,选择合适的加载策略需要在数据一致性、加载性能、实现复杂度三者之间找到平衡。
至此,ETL 的三个核心环节(抽取、转换、加载)已全部介绍完毕。下一篇文章我们将从更宏观的视角审视 ETL,讨论 ETL 架构模式以及与 ELT 的对比。
Summary: 加载策略、目标存储写入、性能优化与失败恢复机制。