数据加载概述

数据加载是 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: 加载策略、目标存储写入、性能优化与失败恢复机制。