Featured image of post DuckDB实战:高级数据清洗与ETL管道进阶

DuckDB实战:高级数据清洗与ETL管道进阶

掌握DuckDB高级数据清洗技巧:日志解析、混合格式处理、JSON清洗、数据质量验证框架,构建生产级ETL管道。

引言

上一篇文章我们介绍了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_extractregexp_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脚本
增量ETLUpsertMERGE, PROCEDURE

生产环境建议

  1. 分层处理:raw → staging → clean → mart,每层都有明确的质量门槛
  2. 保留审计:原始数据永远只读,所有清洗结果写入新表
  3. 质量门禁:关键规则(如主键唯一、必填字段非空)不通过则阻断下游
  4. 增量优先:尽量使用 MERGE 或窗口函数实现增量更新,避免全量重跑
  5. 监控告警:将质量报告接入告警系统,异常率超阈值自动通知

Architecture Diagram

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

SQL Execution Result

图:DuckDB CLI 执行增量ETL质量检查的输出示例

更多 DuckDB 实战技巧,请关注 DuckDB Lab(duckdblab.org)。

📺 Watch video tutorials → Olap Studio YouTube

Subscribe for more DuckDB & AI automation tutorials

使用 Hugo 构建
主题 StackJimmy 设计

⚠️ 本站为独立社区项目,与 DuckDB 基金会及 DuckDB 官方项目无任何从属、背书或赞助关系。

"DuckDB" 是 DuckDB 基金会的注册商标,本站仅以事实描述方式使用该名称。

本站内容仅供教育与社区推广用途,不构成任何商业服务。