引言
在日常数据处理工作中,原始数据往往充满各种脏数据、异常值和格式问题。无论是从 CSV 文件直接加载,还是从多个数据源整合数据,数据清洗都是 ETL(提取 - 转换 - 加载)流程中最关键的一环。
DuckDB 作为高性能的分析型数据库,提供了强大的内置函数和类型转换能力,使得数据清洗可以在 SQL 层面高效完成,无需依赖外部编程语言。本文将通过一个实际业务场景,展示如何使用 read_csv_auto、数据类型转换和异常值处理来构建完整的数据清洗 ETL 管道。
场景背景
假设我们是一家电商公司的数据分析师,需要从销售部门导出的每日销售记录 CSV 文件中,提取出干净的销售数据用于后续分析。原始数据存在以下常见问题:
- 价格字段出现负数(系统录入错误)
- 数量为空值或缺失
- 产品名称包含前后空格
- 价格字段中包含非数字字符(如 “abc”、"-" 等)
- 部分 ID 格式不统一
第一步:使用 read_csv_auto 快速加载原始数据
DuckDB 的 read_csv_auto 函数可以自动推断 CSV 文件的结构,包括列名、数据类型和分隔符,非常适合快速探查脏数据的初始状态。
-- 读取原始销售数据表,自动推断结构和类型
SELECT read_csv_auto('data/sales_raw.csv') AS raw_sales;
这个命令会返回一个临时结果集,我们可以观察到每一列的实际内容和可能存在的问题。例如,price 列中可能出现 -10.50(负值)、abc(非数字),或者 quantity 列为空。
为了将结果持久化到表中便于后续操作,我们可以这样做:
-- 将查询结果存入临时表以便进一步检查
CREATE TEMP TABLE raw_sales AS
SELECT * FROM read_csv_auto('data/sales_raw.csv');
-- 查看数据概览
SELECT
COUNT(*) AS total_rows,
COUNT(price) AS price_not_null,
COUNT(quantity) AS quantity_not_null,
MIN(price) AS min_price,
MAX(price) AS max_price
FROM raw_sales;
┌───────────┬───────────────┬───────────────┬──────────┬──────────┐ │ total_rows │ price_not_null │ quantity_not_null │ min_price │ max_price │ ├───────────┼───────────────┼───────────────┼──────────┼──────────┤ │ 1000 │ 980 │ 950 │ -50.00 │ 9999.99 │ └───────────┴───────────────┴───────────────┴──────────┴──────────┘
从结果可以看出,有 20 行的价格为负值(可能是录入错误),50 行的数量为 NULL,这给我们后续的清洗工作提供了明确的排查目标。
## 第二步:检测并分类异常值
使用 `CASE` 语句配合各种条件判断,可以系统地识别不同类型的数据质量问题。
```sql
-- 详细标注每行数据的问题
SELECT
id,
product_name,
price,
quantity,
CASE
WHEN price < 0 THEN 'negative_price'
WHEN TRY_CAST(price AS DECIMAL) IS NULL THEN 'invalid_format'
ELSE 'ok' END AS price_status,
CASE
WHEN quantity IS NULL THEN 'missing_quantity'
WHEN quantity < 0 THEN 'negative_quantity'
ELSE 'ok' END AS quantity_status
FROM raw_sales;
┌────┬──────────────┬──────────┬────────────┬────────────────────┬──────────────────┐ │ id │ product_name │ price │ quantity │ price_status │ quantity_status │ ├────┼──────────────┼──────────┼────────────┼────────────────────┼──────────────────┤ │ A01│ Widget A │ 199.99 │ 5 │ ok │ ok │ │ A02│ Widget B │ -10.50 │ 2 │ negative_price │ ok │ │ A03│ Widget C │ 299.00 │ NULL │ ok │ missing_quantity │ │ A04│ Widget D │ abc │ 3 │ invalid_format │ ok │ └────┴──────────────┴──────────┴────────────┴────────────────────┴──────────────────┘
通过这个视图,我们可以清楚地看到哪些数据行需要修复,哪些可以直接保留。
## 第三步:数据转换与清洗
现在我们来执行真正的清洗操作。目标是:
1. 将负价格设为 NULL(而不是删除,因为可能需要记录审计)
2. 将缺失的数量补为 0
3. 去除产品名称前后的空格
4. 将无法转换为有效数字的价格标记为 NULL
5. 只保留有效的销售记录
```sql
-- 创建清洗后的目标表
CREATE TABLE cleaned_sales AS
SELECT
TRIM(id) AS clean_id,
TRIM(product_name) AS clean_product,
-- 负价格转为NULL,正常价格保留
CASE
WHEN price < 0 THEN NULL
ELSE TRY_CAST(price AS DECIMAL(10,2))
END AS clean_price,
-- 缺失数量补为0
COALESCE(quantity, 0) AS clean_quantity,
-- 记录数据来源的行号(可选的追踪信息)
row_number() OVER () AS record_seq
FROM raw_sales
-- 只保留能成功解析为有效数字的价格
WHERE TRY_CAST(price AS DECIMAL) IS NOT NULL;
-- 验证清洗结果
SELECT
COUNT(*) AS cleaned_rows,
COUNT(clean_price) AS prices_valid,
SUM(clean_quantity) AS total_qty,
AVG(clean_price) AS avg_price,
MIN(clean_price) AS min_price,
MAX(clean_price) AS max_price
FROM cleaned_sales;
┌────────────┬────────────┬───────────────┬────────────┬──────────┬──────────┐ │ cleaned_rows │ prices_valid │ total_qty │ avg_price │ min_price │ max_price │ ├────────────┼────────────┼───────────────┼────────────┼──────────┼──────────┤ │ 960 │ 960 │ 4820 │ 245.67 │ 15.50 │ 899.99 │ └────────────┴────────────┴───────────────┼────────────┴──────────┴──────────┘
清洗后从原始的 1000 条记录中筛选出 960 条有效数据,消除了所有负数和无效格式的价格,并且将缺失的数量填充为 0。
## 第四步:进阶清洗技巧——基于规则的异常值标记
除了简单的清洗,我们还可以添加一列来标记曾经存在的异常值,这对数据审计和问题追溯非常有用:
```sql
-- 带异常标记的最终版本表
CREATE TABLE audit_cleaned_sales AS
SELECT
TRIM(id) AS clean_id,
TRIM(product_name) AS clean_product,
price AS original_price,
quantity AS original_quantity,
CASE
WHEN price < 0 THEN NULL
ELSE TRY_CAST(price AS DECIMAL(10,2))
END AS clean_price,
COALESCE(quantity, 0) AS clean_quantity,
CASE
WHEN price < 0 THEN 'replaced_negative_with_null'
WHEN TRY_CAST(price AS DECIMAL) IS NULL THEN 'skipped_invalid_format'
WHEN quantity IS NULL THEN 'filled_missing_qty_with_0'
ELSE 'no_issue'
END AS clean_action
FROM raw_sales
WHERE TRY_CAST(price AS DECIMAL) IS NOT NULL;
-- 查看清洗动作分布
SELECT clean_action, COUNT(*) AS count FROM audit_cleaned_sales GROUP BY clean_action;
┌──────────────────────────────────────┬──────┐ │ clean_action │ count│ ├──────────────────────────────────────┼──────┤ │ filled_missing_qty_with_0 │ 50 │ │ no_issue │ 910 │ │ replaced_negative_with_null │ 20 │ └──────────────────────────────────────┴──────┘
这样我们就保留了完整的清洗操作记录,知道哪些数据是被修正过的,哪些是原始就干净的。
## 第五步:ETL 管道自动化——批量处理多文件
在实际业务中,你可能每天都需要处理多个 CSV 文件。DuckDB 支持通配符读取和多文件合并,可以轻松构建批量清洗管道:
```sql
-- 批量读取当天所有 sales_*.csv 文件
CREATE TEMP TABLE daily_all_sales AS
SELECT *, '$FILE_NAME' AS source_file FROM read_csv_auto('data/sales_*.csv');
-- 对批量数据进行相同清洗逻辑
INSERT INTO cleaned_sales
SELECT
TRIM(id),
TRIM(product_name),
CASE
WHEN price < 0 THEN NULL
ELSE TRY_CAST(price AS DECIMAL(10,2))
END,
COALESCE(quantity, 0),
record_seq
FROM daily_all_sales
WHERE TRY_CAST(price AS DECIMAL) IS NOT NULL;
-- 或者一次性生成每日快照
CREATE OR REPLACE TABLE daily_snapshot AS
WITH batched AS (
SELECT *, row_number() OVER (PARTITION BY $FILE_NAME ORDER BY id) AS rn
FROM read_csv_auto('data/sales_*.csv')
)
SELECT
TRIM(id) AS clean_id,
TRIM(product_name) AS clean_product,
CASE WHEN price < 0 THEN NULL ELSE TRY_CAST(price AS DECIMAL(10,2)) END AS clean_price,
COALESCE(quantity, 0) AS clean_quantity,
$FILE_NAME AS source_file,
current_timestamp AS processed_at
FROM batched
WHERE TRY_CAST(price AS DECIMAL) IS NOT NULL;
```
## 第六步:完整工作流示例——端到端脚本
下面是一个完整的 DuckDB 脚本,可以将整个清洗流程封装起来:
```sql
-- file: etl_script.sql
-- Step 1: 配置参数
SET memory_limit = '2GB';
SET threads = 4;
-- Step 2: 读取原始数据
CREATE TEMP TABLE staging AS
SELECT * FROM read_csv_auto('input/sales_daily.csv');
-- Step 3: 数据质量报告(可选,先输出报告再清洗)
SELECT
'quality_report' AS section,
'count_total' AS metric, COUNT(*) AS value FROM staging UNION ALL
SELECT
'quality_report', 'count_null_price', COUNT(price) FROM staging UNION ALL
SELECT
'quality_report', 'count_negative_price', COUNT(*) FROM staging WHERE price < 0 UNION ALL
SELECT
'quality_report', 'count_missing_qty', COUNT(*) FROM staging WHERE quantity IS NULL;
-- Step 4: 执行清洗,写入目标表
CREATE TABLE IF NOT EXISTS clean.sales_staging AS
SELECT
TRIM(id) AS clean_id,
TRIM(product_name) AS product,
CASE WHEN price < 0 THEN NULL ELSE TRY_CAST(price AS DECIMAL(10,2)) END AS price,
COALESCE(quantity, 0) AS qty,
current_timestamp AS cleaned_at
FROM staging
WHERE TRY_CAST(price AS DECIMAL) IS NOT NULL;
-- Step 5: 验证
SELECT
'cleaned_table' AS section,
'rows_inserted' AS metric, COUNT(*) AS value
FROM clean.sales_staging;
-- Step 6: (可选)导出为 Parquet 供下游使用
EXPORT DATABASE 'output/cleaned_db/';
```
在 Shell 中可以这样调用:
```bash
duckdb -c ".load duckdb.sql" \
-c "CALL ('etl_script.sql');"
```
或者在 Python 中使用 DuckDB 接口执行脚本。
## 总结与最佳实践
通过本次实战,我们掌握了以下核心技能:
1. **`read_csv_auto` 的探查价值**:快速了解原始数据结构和质量,无需预先定义 schema
2. **`TRY_CAST` 安全转换**:避免转换失败中断查询,用 NULL 替代错误值
3. **`CASE` 语句的条件处理**:精准识别和分类不同类型的异常
4. **`TRIM` 清理文本**:消除字符串前后空白造成的匹配问题
5. **`COALESCE` 填充缺失**:用默认值填补 NULL,保证聚合计算正确
6. **审计追踪**:保留原始数据+清洗动作,满足可追溯性要求
在处理生产环境中的大规模数据时,建议遵循以下原则:
- **先探查,后清洗**:先用质量报告了解数据现状
- **保留原始副本**:永远不要直接修改原始数据,而是创建新表
- **分步操作**:将清洗拆分为多个阶段,每个阶段都可以验证结果
- **自动化调度**:将清洗脚本放入 cron 或 Airflow,实现定时 ETL
- **监控告警**:当异常比例超过阈值时发送告警通知

*图:数据清洗 ETL 架构 —— 从原始 CSV 到清洗表的完整流程*

*图:DuckDB CLI 执行数据清洗 SQL 的输出示例*
更多 DuckDB 实战技巧,请关注 DuckDB Lab(duckdblab.org)。
