TT Lab
开始
学习 学习路径 课程

Apache Spark — 慢作业的答案在执行计划和事件日志里

每个到达的文件只处理一次,并丢弃迟到的数据

在 TT Lab 中继续学习

目标

每当有文件到达落地目录,就用 Structured Streaming 处理并写入 Parquet sink,确认检查点和 sink 日志是如何记住“处理到了哪里”的。看到只是名字不同的重发文件造成重复之后,用 dropDuplicates 防住,再通过窗口聚合看水位线丢弃迟到数据的样子。

为什么重要

流式处理难的不是计算,而是记忆。作业挂了再起来,必须知道已经做了什么、还没做什么。Structured Streaming 把这些写在检查点文件夹里。在开始批次之前,把要读取的范围写进 offsets,把结果全部写进 sink 之后,写进 commits。文件 sink 会在自己文件夹的 _spark_metadata 里,留下每个批次写出的文件列表,读取一方只会把这个列表里的文件当作结果。 这份记忆是有限度的。文件源靠文件名来记住是否处理过,所以合作方把同样的内容用另一个名字再发一次,就会被当作新数据处理。要按内容去重,就必须把见过的键作为状态保存。 按事件时间做窗口聚合,又会产生另一个问题。窗口什么时候关闭?水位线是“到目前为止见到的最晚时间 − 允许延迟”,比它更旧的数据,当作迟到数据丢弃,只有末尾越过水位线的窗口才会作为结果输出(append 模式)。水位线只在批次之间移动,所以数据被分成几个批次进来,会改变结果。

步骤

  1. 创建落地目录 /root/spk/stream/in,把 /data/stream/batch-01.jsonl 复制进去。
  2. 用 /root/spk/stream/stream.py(应用 spk-stream-run)读取落地目录,把写入 /root/spk/stream/out/events 的 Parquet 的流,以检查点 /root/spk/stream/ckpt/events、trigger(availableNow=True) 运行一次。
  3. 把 batch-02.jsonl 加入落地目录,再运行同一个流。只有新文件应该成为新批次。
  4. 把 batch-03.jsonl 和 batch-04.jsonl 一起加入,再运行。两个文件必须被当作一个批次处理。
  5. 加入重发文件 batch-06.jsonl(内容与 batch-03 相同),再运行之后,把在 sink 里进入了两次的 event_id 数,以整数写入 /root/spk/stream/out/dups.txt。
  6. 用 /root/spk/stream/dedup.py(应用 spk-stream-dedup)把落地目录做 dropDuplicates(["event_id"]),运行写入 /root/spk/stream/out/dedup 的新流(检查点 /root/spk/stream/ckpt/dedup)。
  7. 在第二个落地目录 /root/spk/stream/late_in 里,用 cp -p 复制 batch-01 到 batch-05,并用 /root/spk/stream/window.py(应用 spk-stream-window),以 maxFilesPerTrigger=1、10 分钟水位线、10 分钟窗口计数聚合,以 append 模式写入 /root/spk/stream/out/windows(检查点 /root/spk/stream/ckpt/windows)。把进度记录写入 /root/spk/stream/out/progress.json,把迟到数据的分析写入 /root/spk/stream/out/late.json。
  8. 在 /root/spk/stream/report.md 中以 ## 한 번씩만 처리하기(韩文,意为“只处理一次”)、## 재전송과 중복(韩文,意为“重发与重复”)、## 늦은 자료(韩文,意为“迟到数据”)三节写成报告。第二节放入第 5 步的重复数,第三节放入第 7 步的迟到事件数。

参考

把第一个文件包放进落地目录

创建落地目录 /root/spk/stream/in,把 /data/stream/batch-01.jsonl 复制到里面(/root/spk/stream/in/batch-01.jsonl)。

流式文件源会盯着文件夹,把新出现的文件拿起来作为下一个批次。文件必须以完整的状态一次性出现——如果拿起了正在写入的文件,处理的就是半截文件(所以通常先写到别处,再移动过来)。

第一个批次——检查点与 sink 日志

以应用名 spk-stream-run 创建 /root/spk/stream/stream.py,用 schema(event_id STRING, device STRING, event_time TIMESTAMP, value INT)读取 /root/spk/stream/in,以 checkpointLocation=/root/spk/stream/ckpt/events、trigger(availableNow=True) 启动写入 /root/spk/stream/out/events 的 Parquet 流,并等它结束。

availableNow 的意思是“把现有的全部处理完就停止”。结束之后,检查点里会留下 offsets/0 和 commits/0,sink 文件夹的 _spark_metadata/0 里会留下这个批次写出的文件列表。评分器会读取那个列表里的文件,检查是否与 batch-01 的 event_id 相同。

只有新文件才成为下一个批次

把 /data/stream/batch-02.jsonl 加入落地目录,并把第 2 步的流原样再运行一遍。新批次里必须只有 batch-02 的事件。

因为检查点记得已经处理过 batch-01,所以不会再读。如果删除检查点,记忆也就没了,会从头再处理一遍,sink 里同样的数据就会多堆一份。评分器会在 sink 日志里找只装着 batch-02 的 id 的批次。

两个文件放进一个批次

把 batch-03.jsonl 和 batch-04.jsonl 一起加入落地目录,再运行流。两个文件的事件必须被一个批次处理。

一个批次拿起几个文件,由 maxFilesPerTrigger 决定,不指定的话,就把当时所有的新文件全部拿起来。批次边界会影响延迟和水位线,所以在第 7 步会再遇到。

只是名字不同的重发造成重复

把重发文件 /data/stream/batch-06.jsonl(内容与 batch-03 相同)加入落地目录,再运行流。然后读取 sink 中已提交的文件(每个批次的 _spark_metadata 列表),把出现两次以上的 event_id 的个数,以整数写入 /root/spk/stream/out/dups.txt。

文件源的记忆是文件名。名字是新的,所以当作新数据处理。直接读取 sink 文件夹也可以(Spark 会参照 _spark_metadata 来读取),但要记住,直接遍历目录里的 part 文件的工具,可能连失败批次留下的残渣也拿起来。

dropDuplicates——记住见过的 id

以应用名 spk-stream-dedup 创建 /root/spk/stream/dedup.py,读取同一个落地目录,把 dropDuplicates(["event_id"]) 的结果写入 /root/spk/stream/out/dedup,用 checkpointLocation=/root/spk/stream/ckpt/dedup、availableNow 运行这个新流。结果中的 event_id 必须每个只出现一次。

去重是有状态的运算。它在状态存储里保存见过的 id,遇到同样的 id 就丢弃。如果没有水位线,状态就会无限增长,所以在生产环境里要与 withWatermark 一起用,或者使用 dropDuplicatesWithinWatermark。检查点是每个流各用各的。

水位线——丢弃迟到数据,只输出已关闭的窗口

创建 /root/spk/stream/late_in,用 cp -p 复制 batch-01 到 batch-05,再以应用名 spk-stream-window 创建 /root/spk/stream/window.py,以 maxFilesPerTrigger=1 读取,在 withWatermark("event_time", "10 minutes") 之后按 10 分钟窗口计数,以 start、end、count 列用 append 模式写入 /root/spk/stream/out/windows(checkpointLocation=/root/spk/stream/ckpt/windows,availableNow)。结束之后,从 query.recentProgress 里,把每个批次的 batch、rows、watermark、dropped(状态算子的 numRowsDroppedByWatermark 之和),以列表写入 /root/spk/stream/out/progress.json,并把早于 batch-05 进来时水位线的 batch-05 事件数,以 {"watermark": "yyyy-MM-dd HH:mm:ss", "late_events": 정수, "dropped_rows_metric": 정수}(占位符为整数)写入 /root/spk/stream/out/late.json。

水位线在批次结束时升到“见到的最晚时间 − 10 分钟”,从下一个批次开始起作用。所以要让文件一个一个地作为批次放进来,batch-05 进来的时候水位线才已经在 09:29 前后。也看看迟到的事件有几十条,而 dropped 指标却是 1——因为在状态算子之前,部分聚合已经把它们按窗口合并成了一行。

用数字留下记忆、重复、迟到

在 /root/spk/stream/report.md 中写出 ## 한 번씩만 처리하기(韩文,意为“只处理一次”)、## 재전송과 중복(韩文,意为“重发与重复”)、## 늦은 자료(韩文,意为“迟到数据”)三节。第二节放入第 5 步的重复数,第三节放入第 7 步的 late_events 和 dropped_rows_metric。

第一节写检查点的 offsets、commits 和 sink 日志各自记住什么,第二节写为什么只改了名字的文件会成为重复,第三节写迟到事件数和指标为什么不同。