Apache Spark — 慢作业的答案在执行计划和事件日志里
流处理是一连串小批次,由检查点负责记忆
一句话总结
Structured Streaming 是在无限增长的表上用小批次反复运行同一个查询的引擎,处理到哪里了由检查点记住,写出了什么由 sink 的记录记住。要等多久迟到的数据,由水位线决定。
为什么不能每天重新运行批处理
假设合作方的文件一天里有几十次落到落地目录。如果用批处理作业来解决,只有两条路。每次把整个目录重新读一遍(数据越积越多就越慢),或者自己管理已读文件列表(作业中途挂了,列表和结果就会对不上)。想把第二条路做对,最终还是要亲手编写一套“先写下决定处理的内容,写完之后再写下已结束”的机制。
Structured Streaming 就是把这套机制内置到了引擎里。入门文档把传入的数据看作不断追加行的输入表,把查询看作它之上的结果表。每个触发器只用新行来更新结果。查询的写法和批处理完全一样,增量运行的事交给引擎。根据概述文档,默认的执行方式是把这件事当作一连串小批处理作业来处理的微批处理(micro-batch)。
工作原理——文件源
API 文档中的文件源,会读取目录中新出现的文件。有四条规则需要知道。
- 文件按修改时间顺序处理。
- 文件必须原子地放进目录。在大多数文件系统里,做法是先在别处写完再移动(move)。如果就地慢慢写,就可能读到写到一半的文件。
- 是不是新文件,默认按完整路径判断(
fileNameOnly默认 false)。内容相同但名字不同,也是新文件。 maxFilesPerTrigger是一个触发器最多拿起的新文件数的上限,默认没有限制。maxFileAge默认是一周,但第一个批次里会把所有文件都视为有效。
第三条规则在实际工作中会造成重复。合作方把昨天发的文件包只改个名字重新发来,引擎就会把它当作新文件接收。文件层面的“一次”是保证的,但内容层面的一次并不保证。这要在后面用去重来解决。
检查点与 sink 的记录
入门文档的容错一节用一句话总结了这个设计。每个源都有表示读取位置的偏移量,引擎用检查点和预写日志(write-ahead log),在每个触发器里记录要处理的偏移量范围。sink 被设计成即使重复收到同一个批次,结果也一样(幂等)。可重新读取的源和幂等的 sink 相遇,就成了端到端精确一次。
打开检查点目录,就能在文件里看到这个设计。offsets/0 是第 0 号批次决定处理的范围,在处理之前写入。commits/0 是表示该批次已结束的标记,在处理之后写入。重启的查询会找到在 offsets 里有、在 commits 里没有的批次,并按同样的范围重新运行。文件 sink 一侧也有配对的记录。输出目录的 _spark_metadata/0 里记着第 0 号批次写出的文件列表,用 Spark 读取那个目录,只会看到这个列表里的文件。这就是死掉的批次留下的半截文件不会混进结果的原因。API 文档的表里,文件 sink 被标记为精确一次,也是多亏了这条记录。
所以删除检查点就是在删除记忆。用新的检查点启动同一个查询,会把目录里的所有文件从头再处理一遍。同一份文档还把在重启之间改变文件 sink 的输出路径或去重列,列为不被允许的更改。
触发器——什么时候运行批次
不指定触发器,就是上一个批次一结束就运行下一个批次。给定间隔,就按这个间隔运行。像本实验这样,只处理现有数据就停下的用法,适合用 availableNow。根据 API 文档,它会把运行时已有的数据全部处理完之后自行停止,并根据源选项(文件源就是 maxFilesPerTrigger)分成多个批次处理,而且会先处理之前运行中没能提交的批次。以前的一次性触发器(once)已预定废弃。
q = (spark.readStream.schema(schema).json("/root/landing")
.withWatermark("event_time", "10 minutes")
.dropDuplicates(["event_id", "event_time"])
.writeStream.format("parquet")
.option("path", "/root/out")
.option("checkpointLocation", "/root/chk")
.trigger(availableNow=True)
.start())
q.awaitTermination()
去重与水位线
流式的 dropDuplicates 含义与批处理相同,但代价不同。根据 dropDuplicates 文档,在流式中要把已经见过的键跨触发器地全部作为状态保存。没有水位线,这个状态就会无限增长。
水位线是“比这更晚到的数据,就不再等了”的线。withWatermark 文档把这条线定义为到目前为止见到的最大事件时间减去阈值。这条线做两件事。它告诉你窗口聚合里哪些窗口已经确定,并清除已确定窗口的状态。所以 Append 模式的窗口聚合,不会在窗口刚结束时立刻输出。要等水位线越过窗口的末尾之后,才会输出一次。
一定要记住,保证只有一个方向。API 文档保证,10 分钟的水位线绝不会丢弃在 10 分钟之内迟到的数据,但并没有说比这更晚的数据一定会被丢弃。通常会被丢弃,但也可能被聚合进去。水位线不是过滤迟到数据的过滤器,而是清除状态的依据。而且要在聚合里清除状态,必须把水位线加在聚合之前、加在聚合所用的同一个时间列上,并且输出模式必须是 Append 或 Update。
在现场相遇的样子
第一,把检查点放在临时文件夹里。一次重启,记忆就没了,查询会把所有文件重新处理一遍。检查点是和输出一样宝贵的数据。
第二,在落地目录里就地写入。上传工具在写文件的过程中触发器运行的话,就会读到半截文件。要让它先用临时名字写完,再移动过去。
第三,重发造成重复。文件源是靠路径来判断新文件的,所以改了名字重新发来的文件包,会原样进来。用唯一 ID 去重,并同时加上水位线,避免状态无限增长。
实际工作中真正重要的事
- 流式处理就是一连串小批次。查询的写法和批处理完全一样。
- offsets 在处理之前写入,commits 在处理之后写入。在两者之间挂掉,就会重新运行同样的范围。
- 文件 sink 的结果只有 _spark_metadata 里记录的文件。
- 文件源靠路径判断新文件。内容重复要用去重来防止。
- 水位线是清除状态的依据。阈值之内一定会接收,之外通常会丢弃。
下一项实验要做什么
把设备事件文件包放进落地目录,用 availableNow 运行第一个批次,确认检查点的 commits 和输出目录的 _spark_metadata 里出现了第 0 号记录。放入新的文件包,看到只有它被处理;一次放入两个文件包,看到它们被当作一个批次处理。数一数改了名字重发的文件包带来的重复 event_id,再用新检查点运行 dropDuplicates 流,只留下唯一的条数。最后,在保持修改时间复制过来的第二个落地目录里,每个触发器一个文件,运行 10 分钟水位线和 10 分钟窗口聚合,通过进度记录看水位线移动的样子,以及晚到 40 分钟左右的事件被丢弃。