引言
上一篇文章我们介绍了DuckDB基础的数据清洗与ETL管道。本文将继续深入,探讨更复杂的真实场景:日志文件解析、混合数据格式清洗、JSON嵌套数据展平、以及生产级数据质量验证框架。这些模式在电商、金融、IoT等行业的日常数据处理中非常常见。
场景一:解析半结构化日志数据
问题背景
公司的应用服务器每天生成大量Nginx访问日志,格式如下:
192.168.1.100 - - [04/Sep/2026:10:15:30 +0800] "GET /api/products HTTP/1.1" 200 1234 "-" "Mozilla/5.0 (Windows NT 10.0; Win64; x64)"
10.0.0.55 - admin [04/Sep/2026:10:15:31 +0800] "POST /api/orders HTTP/1.1" 201 567 "https://example.com" "curl/7.81.0"
这些日志包含IP地址、时间戳、HTTP方法、路径、状态码、响应大小和User-Agent。我们需要从中提取结构化数据用于流量分析。
DuckDB解决方案
DuckDB提供了强大的正则表达式函数 regexp_extract 和 regexp_match,可以直接从日志行中提取字段:
-- 读取并解析Nginx日志
CREATE TABLE nginx_access_logs AS
SELECT
regexp_extract(line, '^(\\S+)', 1) AS client_ip,
regexp_extract(line, '^\\S+ \\S+ (\\S+)', 1) AS remote_user,
regexp_extract(line, '\\[(\\d{2}/\\w{3}/\\d{4}:\\d{2}:\\d{2}:\\d{2}) [+\\-]\\d{4}\\]', 1) AS timestamp_str,
regexp_extract(line, '"(\\w+) (\\S+) (\\S+)"', 1) AS http_method,
regexp_extract(line, '"\\w+ (\\S+) \\S+"', 1) AS request_path,
regexp_extract(line, '"\\w+ \\S+ (\\S+)"', 1) AS http_version,
regexp_extract(line, '"\\w+ \\S+ \\S+" (\\d{3})', 1) AS status_code,
regexp_extract(line, '"\\w+ \\S+ \\S+" \\d{3} (\\d+)', 1) AS response_size,
regexp_extract(line, '"([^"]*)" "[^"]*"$', 1) AS referer,
regexp_extract(line, '"[^"]*" "([^"]*)"', 1) AS user_agent
FROM read_csv_auto('logs/nginx_2026-09-04.log');
-- 查看解析结果
SELECT * FROM nginx_access_logs LIMIT 5;
图:Nginx日志解析流程架构图

图:Nginx日志正则解析流程——从原始文本到结构化列
时间戳转换与规范化
-- 将解析出的时间字符串转换为TIMESTAMP,并提取有用维度
CREATE TABLE parsed_access_logs AS
SELECT
client_ip,
remote_user,
TRY_CAST(timestamp_str AS TIMESTAMP) AS log_time,
http_method,
request_path,
status_code::INTEGER AS status,
response_size::BIGINT AS bytes_sent,
user_agent,
-- 提取URL路径参数
CASE
WHEN request_path LIKE '/api/%' THEN 'api'
WHEN request_path LIKE '/static/%' THEN 'static'
WHEN request_path LIKE '/images/%' THEN 'images'
ELSE 'other'
END AS path_category,
-- 判断是否为搜索引擎爬虫
CASE
WHEN user_agent ILIKE '%googlebot%' THEN 'google'
WHEN user_agent ILIKE '%bingbot%' THEN 'bing'
WHEN user_agent ILIKE '%baiduspider%' THEN 'baidu'
ELSE 'human'
END AS bot_type,
-- 计算响应时间分类
CASE
WHEN status BETWEEN 200 AND 299 THEN 'success'
WHEN status BETWEEN 300 AND 399 THEN 'redirect'
WHEN status BETWEEN 400 AND 499 THEN 'client_error'
WHEN status BETWEEN 500 AND 599 THEN 'server_error'
ELSE 'unknown'
END AS status_category
FROM nginx_access_logs
WHERE TRY_CAST(timestamp_str AS TIMESTAMP) IS NOT NULL;
-- 验证转换结果
SELECT
path_category,
bot_type,
status_category,
COUNT(*) AS request_count,
ROUND(AVG(bytes_sent)) AS avg_bytes
FROM parsed_access_logs
GROUP BY path_category, bot_type, status_category
ORDER BY request_count DESC;
图:时间戳转换与维度提取SQL执行结果

图:DuckDB CLI 执行日志解析SQL的输出示例
场景二:清洗混合格式的CSV数据
问题背景
在一次数据迁移中,你收到了来自三个不同供应商的销售数据,每个供应商使用不同的CSV格式:
- 供应商A:标准格式,逗号分隔,日期为
YYYY-MM-DD - 供应商B:分号分隔,日期为
DD/MM/YYYY,价格使用欧洲格式(.作为千位分隔符,,作为小数点) - 供应商C:制表符分隔,包含多行字符串字段(用双引号包裹),日期字段可能缺失
统一清洗流程
-- 第一步:分别读取三种格式的数据
CREATE TEMP TABLE supplier_a AS
SELECT
order_id,
TRY_CAST(customer_id AS INTEGER) AS customer_id,
product_name,
TRY_CAST(REPLACE(price, '.', '') AS DECIMAL(12,2)) AS price, -- 移除千位分隔符
TRY_CAST(date AS DATE) AS order_date,
quantity,
'supplier_a' AS source
FROM read_csv_auto('data/supplier_a.csv', sep=',');
CREATE TEMP TABLE supplier_b AS
SELECT
order_id,
TRY_CAST(customer_id AS INTEGER) AS customer_id,
product_name,
TRY_CAST(REPLACE(REPLACE(price, '.', ''), ',', '.') AS DECIMAL(12,2)) AS price, -- 欧洲格式转换
TRY_CAST(strptime(date, '%d/%m/%Y') AS DATE) AS order_date,
TRY_CAST(quantity AS INTEGER) AS quantity,
'supplier_b' AS source
FROM read_csv_auto('data/supplier_b.csv', sep=';');
CREATE TEMP TABLE supplier_c AS
SELECT
order_id,
TRY_CAST(customer_id AS INTEGER) AS customer_id,
product_name,
TRY_CAST(price AS DECIMAL(12,2)) AS price,
TRY_CAST(date AS DATE) AS order_date,
TRY_CAST(quantity AS INTEGER) AS quantity,
'supplier_c' AS source
FROM read_csv_auto('data/supplier_c.tsv', sep='\t');
-- 第二步:合并并统一格式
CREATE TABLE unified_sales AS
SELECT order_id, customer_id, product_name, price, order_date, quantity, source
FROM supplier_a
UNION ALL
SELECT order_id, customer_id, product_name, price, order_date, quantity, source
FROM supplier_b
UNION ALL
SELECT order_id, customer_id, product_name, price, order_date, quantity, source
FROM supplier_c;
-- 第三步:数据质量检查
SELECT
source,
COUNT(*) AS total_records,
COUNT(order_id) AS valid_ids,
COUNT(customer_id) AS valid_customers,
COUNT(order_date) AS valid_dates,
COUNT(price) AS valid_prices,
COUNT(quantity) AS valid_quantities,
ROUND(AVG(price), 2) AS avg_price,
MIN(price) AS min_price,
MAX(price) AS max_price
FROM unified_sales
GROUP BY source;
处理重复订单与冲突数据
-- 检测同一订单ID在不同供应商中的重复记录
SELECT
order_id,
COUNT(*) AS occurrence_count,
GROUP_CONCAT(DISTINCT source) AS sources,
GROUP_CONCAT(DISTINCT customer_id) AS customer_ids,
SUM(price) AS total_price_sum
FROM unified_sales
GROUP BY order_id
HAVING COUNT(*) > 1;
-- 智能去重:优先保留数据最完整的记录
CREATE TABLE deduplicated_sales AS
WITH ranked AS (
SELECT *,
row_number() OVER (
PARTITION BY order_id
ORDER BY
-- 优先保留有日期的
CASE WHEN order_date IS NOT NULL THEN 0 ELSE 1 END,
-- 其次保留客户ID完整的
CASE WHEN customer_id IS NOT NULL THEN 0 ELSE 1 END,
-- 最后按source优先级
CASE source WHEN 'supplier_a' THEN 0 WHEN 'supplier_b' THEN 1 ELSE 2 END
) AS rank
FROM unified_sales
)
SELECT * FROM ranked WHERE rank = 1;
SELECT COUNT(*) AS original_count FROM unified_sales;
SELECT COUNT(*) AS deduplicated_count FROM deduplicated_sales;
场景三:清洗嵌套JSON数据
问题背景
某电商平台通过API返回的订单数据是嵌套JSON格式,包含订单基本信息、商品列表、收货地址和支付方式。我们需要将其展平为关系表。
展开嵌套JSON
-- 读取嵌套JSON订单数据
CREATE TABLE raw_orders_json AS
SELECT * FROM read_json_auto('data/orders.json');
-- 使用JSON展开函数处理嵌套结构
CREATE TABLE flattened_orders AS
SELECT
o.order_id,
o.customer_id,
o.order_time,
o.status,
-- 展开商品数组
item.product_id,
item.product_name,
item.quantity,
item.unit_price,
item.discount_percent,
-- 展开地址信息
addr.country,
addr.province,
addr.city,
addr.district,
addr.full_address,
-- 展开支付方式
payment.method,
payment.amount AS paid_amount,
payment.transaction_id
FROM raw_orders_json o,
LATERAL unpack(o.items) AS item,
LATERAL o.shipping_address AS addr,
LATERAL unpack(o.payments) AS payment;
-- 查看展开后的结果
SELECT * FROM flattened_orders LIMIT 10;
JSON数据清洗与验证
-- 清洗和验证展开后的数据
CREATE TABLE clean_orders AS
SELECT
order_id,
TRY_CAST(customer_id AS BIGINT) AS customer_id,
TRY_CAST(order_time AS TIMESTAMP) AS order_time,
-- 标准化订单状态
LOWER(TRIM(status)) AS order_status,
TRY_CAST(product_id AS BIGINT) AS product_id,
TRIM(product_name) AS product_name,
GREATEST(TRY_CAST(quantity AS INTEGER), 0) AS quantity,
GREATEST(TRY_CAST(unit_price AS DECIMAL(10,2)), 0) AS unit_price,
LEAST(GREATEST(TRY_CAST(discount_percent AS DECIMAL(5,2)), 0), 100) AS discount,
UPPER(country) AS country,
TRIM(province) AS province,
TRIM(city) AS city,
REPLACE(full_address, E'\n', ' ') AS address_clean,
UPPER(method) AS pay_method,
TRY_CAST(paid_amount AS DECIMAL(12,2)) AS paid_amount,
TRIM(transaction_id) AS txn_id
FROM flattened_orders
WHERE TRY_CAST(customer_id AS BIGINT) IS NOT NULL
AND TRY_CAST(unit_price AS DECIMAL(10,2)) IS NOT NULL;
-- 数据统计
SELECT
order_status,
pay_method,
COUNT(DISTINCT order_id) AS order_count,
COUNT(*) AS line_item_count,
ROUND(SUM(unit_price * quantity * (1 - discount/100)), 2) AS total_revenue,
ROUND(AVG(paid_amount), 2) AS avg_payment
FROM clean_orders
GROUP BY order_status, pay_method
ORDER BY order_count DESC;
场景四:生产级数据质量验证框架
质量规则定义
在生产环境中,数据清洗不仅仅是转换格式,还需要建立自动化的质量验证机制。DuckDB支持通过SQL定义质量规则,并自动生成报告:
-- 创建质量规则表
CREATE TABLE data_quality_rules (
rule_id VARCHAR PRIMARY KEY,
table_name VARCHAR,
column_name VARCHAR,
rule_type VARCHAR,
rule_expression VARCHAR,
severity VARCHAR, -- 'error', 'warning', 'info'
description VARCHAR
);
-- 插入质量规则
INSERT INTO data_quality_rules VALUES
('dq_001', 'clean_orders', 'order_id', 'not_null', 'order_id IS NOT NULL', 'error', '订单ID不能为空'),
('dq_002', 'clean_orders', 'customer_id', 'not_null', 'customer_id IS NOT NULL', 'error', '客户ID不能为空'),
('dq_003', 'clean_orders', 'unit_price', 'positive', 'unit_price > 0', 'error', '单价必须为正数'),
('dq_004', 'clean_orders', 'quantity', 'range', 'quantity >= 1 AND quantity <= 999', 'warning', '数量应在1-999之间'),
('dq_005', 'clean_orders', 'discount', 'range', 'discount >= 0 AND discount <= 100', 'warning', '折扣应在0-100%之间'),
('dq_006', 'clean_orders', 'order_time', 'not_null', 'order_time IS NOT NULL', 'error', '订单时间不能为空'),
('dq_007', 'clean_orders', 'paid_amount', 'match', 'ABS(paid_amount - unit_price * quantity * (1 - discount/100)) < 0.01', 'warning', '支付金额应与计算金额匹配'),
('dq_008', 'clean_orders', 'country', 'pattern', 'country REGEXP /^[A-Z]{2}$/', 'info', '国家代码应为2位大写字母');
-- 执行质量检查并生成报告
CREATE TABLE quality_report AS
WITH violations AS (
SELECT
r.rule_id,
r.table_name,
r.column_name,
r.rule_type,
r.severity,
r.description,
COUNT(*) AS violation_count,
COUNT(*) * 100.0 / (SELECT COUNT(*) FROM clean_orders) AS violation_rate_pct
FROM clean_orders o
JOIN data_quality_rules r ON TRUE
WHERE
CASE r.rule_type
WHEN 'not_null' THEN EVALUATE_EXPRESSION(r.rule_expression, o) IS FALSE
WHEN 'positive' THEN EVALUATE_EXPRESSION(r.rule_expression, o) IS FALSE
WHEN 'range' THEN EVALUATE_EXPRESSION(r.rule_expression, o) IS FALSE
WHEN 'pattern' THEN EVALUATE_EXPRESSION(r.rule_expression, o) IS FALSE
ELSE TRUE
END
GROUP BY r.rule_id, r.table_name, r.column_name, r.rule_type, r.severity, r.description
)
SELECT * FROM violations;
提示:DuckDB本身不内置
EVALUATE_EXPRESSION函数,上面的伪代码展示了思路。实际实现中,可以用Python脚本遍历规则表,动态构建SQL检查语句。下面展示实际可用的实现方式:
import duckdb
import json
con = duckdb.connect('orders.duckdb')
# 定义质量检查函数
def run_quality_checks(conn, table_name, rules):
"""执行数据质量检查"""
results = []
total_rows = conn.execute(f"SELECT COUNT(*) FROM {table_name}").fetchone()[0]
for rule in rules:
rule_id = rule['rule_id']
condition = rule['condition']
severity = rule['severity']
description = rule['description']
# 动态执行检查
query = f"SELECT COUNT(*) FROM {table_name} WHERE NOT ({condition})"
violation_count = conn.execute(query).fetchone()[0]
rate = violation_count * 100.0 / total_rows if total_rows > 0 else 0
results.append({
'rule_id': rule_id,
'table': table_name,
'violation_count': violation_count,
'violation_rate_pct': round(rate, 2),
'severity': severity,
'description': description,
'total_rows': total_rows
})
return results
# 定义规则
rules = [
{'rule_id': 'dq_001', 'condition': 'order_id IS NOT NULL', 'severity': 'error', 'description': '订单ID不能为空'},
{'rule_id': 'dq_002', 'condition': 'customer_id IS NOT NULL', 'severity': 'error', 'description': '客户ID不能为空'},
{'rule_id': 'dq_003', 'condition': 'unit_price > 0', 'severity': 'error', 'description': '单价必须为正数'},
{'rule_id': 'dq_004', 'condition': 'quantity >= 1 AND quantity <= 999', 'severity': 'warning', 'description': '数量应在1-999之间'},
{'rule_id': 'dq_005', 'condition': 'discount >= 0 AND discount <= 100', 'severity': 'warning', 'description': '折扣应在0-100%之间'},
]
# 执行检查
report = run_quality_checks(con, 'clean_orders', rules)
con.execute("CREATE TABLE quality_check_report AS SELECT * FROM report")
# 输出报告
for r in report:
print(f"[{r['severity'].upper()}] {r['rule_id']}: {r['violation_count']} violations ({r['violation_rate_pct']}%) - {r['description']}")
质量报告汇总查询
-- 汇总质量报告
SELECT
severity,
rule_id,
description,
violation_count,
violation_rate_pct,
total_rows,
CASE
WHEN severity = 'error' AND violation_count > 0 THEN 'BLOCK'
WHEN severity = 'warning' AND violation_rate_pct > 5 THEN 'REVIEW'
ELSE 'PASS'
END AS action_required
FROM quality_check_report
ORDER BY
CASE severity WHEN 'error' THEN 0 WHEN 'warning' THEN 1 ELSE 2 END,
violation_count DESC;
场景五:增量ETL与数据同步
增量更新模式
在实际生产环境中,数据通常是增量更新的。我们需要设计一个能够处理增量数据的ETL流程:
-- 创建增量ETL存储过程
CREATE OR REPLACE PROCEDURE incremental_etl_load()
LANGUAGE SQL
AS $$
BEGIN
-- Step 1: 创建增量 staging 表
CREATE TEMP TABLE IF NOT EXISTS etl_staging AS
SELECT * FROM read_csv_auto('data/incremental_orders_2026-09.csv')
WHERE 1=0; -- 只创建结构
-- Step 2: 加载增量数据
INSERT INTO etl_staging
SELECT * FROM read_csv_auto('data/incremental_orders_2026-09.csv');
-- Step 3: Upsert逻辑 —— 匹配已有记录则更新,否则插入
MERGE INTO clean_orders AS target
USING etl_staging AS source
ON target.order_id = source.order_id
WHEN MATCHED THEN
UPDATE SET
customer_id = source.customer_id,
order_time = source.order_time,
status = source.status,
updated_at = CURRENT_TIMESTAMP
WHEN NOT MATCHED THEN
INSERT (order_id, customer_id, product_id, product_name,
quantity, unit_price, discount, order_time, status)
VALUES (
source.order_id, source.customer_id, source.product_id,
source.product_name, source.quantity, source.unit_price,
source.discount, source.order_time, source.status
);
-- Step 4: 记录ETL元数据
INSERT INTO etl_audit_log (
load_type, source_file, rows_loaded, rows_updated,
rows_inserted, started_at, completed_at
)
SELECT
'incremental',
'incremental_orders_2026-09.csv',
(SELECT COUNT(*) FROM etl_staging),
(SELECT COUNT(*) FROM clean_orders WHERE updated_at > CURRENT_DATE - 1),
(SELECT COUNT(*) FROM etl_staging s LEFT JOIN clean_orders c ON s.order_id = c.order_id WHERE c.order_id IS NULL),
CURRENT_TIMESTAMP - INTERVAL '5' SECOND,
CURRENT_TIMESTAMP;
END;
$$;
-- 执行增量ETL
CALL incremental_etl_load();
ETL审计日志
-- 创建审计日志表
CREATE TABLE IF NOT EXISTS etl_audit_log (
log_id BIGINT AUTO_INCREMENT PRIMARY KEY,
load_type VARCHAR,
source_file VARCHAR,
rows_loaded INTEGER,
rows_updated INTEGER,
rows_inserted INTEGER,
started_at TIMESTAMP,
completed_at TIMESTAMP,
duration_seconds DOUBLE,
status VARCHAR
);
-- 计算执行耗时并更新
UPDATE etl_audit_log
SET
duration_seconds = (completed_at - started_at)::DOUBLE,
status = 'SUCCESS'
WHERE completed_at IS NOT NULL AND status IS NULL;
-- 查询ETL性能趋势
SELECT
DATE(started_at) AS load_date,
load_type,
COUNT(*) AS run_count,
AVG(duration_seconds) AS avg_duration,
SUM(rows_loaded) AS total_rows,
SUM(rows_inserted) AS total_inserted,
SUM(rows_updated) AS total_updated
FROM etl_audit_log
GROUP BY DATE(started_at), load_type
ORDER BY load_date DESC, load_type;
总结与最佳实践
本文涵盖了DuckDB在高级数据清洗和ETL管道中的五个核心场景:
| 场景 | 核心技术 | 关键函数 |
|---|---|---|
| 日志解析 | 正则表达式 | regexp_extract, regexp_match |
| 混合格式清洗 | 多源合并 | read_csv_auto, strptime, UNION ALL |
| JSON展平 | 嵌套展开 | unpack, LATERAL |
| 质量验证 | 规则引擎 | 动态SQL + Python脚本 |
| 增量ETL | Upsert | MERGE, PROCEDURE |
生产环境建议:
- 分层处理:raw → staging → clean → mart,每层都有明确的质量门槛
- 保留审计:原始数据永远只读,所有清洗结果写入新表
- 质量门禁:关键规则(如主键唯一、必填字段非空)不通过则阻断下游
- 增量优先:尽量使用
MERGE或窗口函数实现增量更新,避免全量重跑 - 监控告警:将质量报告接入告警系统,异常率超阈值自动通知

图:高级数据清洗ETL架构 —— 从多源脏数据到高质量数仓的完整流水线

图:DuckDB CLI 执行增量ETL质量检查的输出示例
更多 DuckDB 实战技巧,请关注 DuckDB Lab(duckdblab.org)。