ETL 很适合检验一门语言的数据抽象是不是只停留在语法层。
真实管道不会只有一个整齐的 JSON 数组。输入可能同时来自 CSV、文件目录和 HTTP;中间会遇到脏字段、无效记录和远程失败;输出既要给人看,也要给下游程序继续消费。
DataFlow ETL 读取客户 CSV 与事件 JSON 目录,统一字段,过滤停用客户,并发请求画像 API,再按最低消费筛选、按部门汇总,最后原子写出 JSON 报告和 CSV 明细。

先把不同输入变成明确的数据
客户 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 })

自测固定验证 3 条合格客户、2 个事件文件和 2 个部门汇总,还检查邮箱清洗、消费排序、region/tier 补全以及 JSON/CSV 两份输出。
sh practical-projects/dataflow-etl/self-test.sh
我从这条 Flow 得到的判断
ETL 的复杂度通常不在某个 map,而在边界:什么时候清洗,失败是否保留,远程调用是否受控,输出能否安全替换。
HHY 的价值也在这些边界上。Pipe 让数据方向清楚,Stream 表达筛选与分组,attempt 把失败变成数据,Runtime 则统一超时和原子副作用。最后得到的不是一段“像 SQL 的语法”,而是一条可以运行、失败、解释和重跑的工程管道。