技术札记

实战 DataFlow ETL:用一条 HHY Flow 串起 CSV、JSON 与 HTTP

从客户 CSV、事件 JSON 到并发画像补全、部门聚合和原子双输出,我如何用 HHY 写一条完整 ETL 管道。

HOUHUIYANG.COM

扫码继续阅读

正在生成…

实战 DataFlow ETL:用一条 HHY Flow 串起 CSV、JSON 与 HTTP

houhuiyang.com/zh/notes/building-dataflow-etl-with-hhy

ETL 很适合检验一门语言的数据抽象是不是只停留在语法层。

真实管道不会只有一个整齐的 JSON 数组。输入可能同时来自 CSV、文件目录和 HTTP;中间会遇到脏字段、无效记录和远程失败;输出既要给人看,也要给下游程序继续消费。

DataFlow ETL 读取客户 CSV 与事件 JSON 目录,统一字段,过滤停用客户,并发请求画像 API,再按最低消费筛选、按部门汇总,最后原子写出 JSON 报告和 CSV 明细。

DataFlow ETL 的真实项目目录

先把不同输入变成明确的数据

客户 CSV 在进入主流程时完成 trim、邮箱小写化、布尔值和整数转换:

export fn load_customers(input_root) {
    return path_join(input_root, "customers.csv")
        |> read_lines
        |> parse_csv({ header: true })
        |> map { row ->
        return {
            id: trim(row.id),
            name: trim(row.name),
            email: lower(trim(row.email)),
            department: trim(row.department),
            active: lower(trim(row.active)) == "true",
            spend: to_int(row.spend)
        }
    }
        |> collect
}

事件则通过 files("*.json") 遍历目录。这里我没有做一个万能 Loader:CSV 和 JSON 的读取方式不同,但它们离开 source 模块时都已经是业务可以理解的 Map。

并发补全,但不丢失失败信息

停用客户没有必要请求远程画像,因此过滤发生在 HTTP 之前。有效客户进入 parallel(4),每个请求带 3 秒超时和一次重试。

let profile = attempt {
    http.get(join([api_base, "profiles", customer.id], "/"))
        |> timeout(3s)
        |> retry({ count: 1, backoff: 50ms })
        |> send
        |> response_body
        |> parse_json
}

补全失败的客户不会消失。结果仍保留原始字段,同时写入 enriched: false 和错误消息。最终报告只有在所有入选客户都补全成功时才 ok=true。这比 silently drop 更诚实,也让重跑和排障有依据。

Stream 负责关系,数组负责边界

最低消费筛选与排序是一条短 Flow:

customers
    |> stream
    |> where { customer -> customer.spend >= minimum_spend }
    |> sort_by({ order: "desc" }) { customer -> customer.spend }
    |> collect

部门统计则用 group_by 后对每组的 spend 求和。我的习惯是在可以惰性处理的阶段保持 Stream,到要并发、复用或写入报告的边界再 collect。这样生命周期一眼可见,也不会让同一条 Stream 被意外消费两次。

双输出不是复制两套逻辑

JSON 报告保存客户、事件、部门汇总和来源统计,适合诊断;CSV 只投影下游需要的七个字段。二者都使用原子写入,旧文件不会在进程中途失败时留下半份新内容。

report |> encode_json({ pretty: true }) |> save_text(report_path, { atomic: true })
customers |> csv_rows |> encode_csv({ header: true }) |> save_lines(csv_path, { atomic: true })

DataFlow ETL 的端到端自测结果

自测固定验证 3 条合格客户、2 个事件文件和 2 个部门汇总,还检查邮箱清洗、消费排序、region/tier 补全以及 JSON/CSV 两份输出。

sh practical-projects/dataflow-etl/self-test.sh

我从这条 Flow 得到的判断

ETL 的复杂度通常不在某个 map,而在边界:什么时候清洗,失败是否保留,远程调用是否受控,输出能否安全替换。

HHY 的价值也在这些边界上。Pipe 让数据方向清楚,Stream 表达筛选与分组,attempt 把失败变成数据,Runtime 则统一超时和原子副作用。最后得到的不是一段“像 SQL 的语法”,而是一条可以运行、失败、解释和重跑的工程管道。

参考

返回技术札记