笔记本上没事,到了 Pod 就被杀
目标
构建不把大文件加载到内存里就能处理它的工具 stream.py,并设定内存上限,用退出码证明它真的是流式处理。汇总、前 N 个、唯一值统计、外部排序和验证,全部以“手里拿着的东西的个数”来设计。
为什么重要
整体读入文件的代码,在文件小的时候不仅更短,而且更快。所以第一次写的时候没有任何理由怀疑它,而文件会变多大,由客户而不是我们来决定。容器超过内存限制时,内核会杀掉进程,而这种死亡不会被 Python 异常捕捉到,标准输出里也什么都不会留下。 设计的标准不是文件大小,而是在任何时刻手里拿着的东西的个数。汇总只需要一行,前 N 条只需要 N 个。可是统计唯一值,即使是流式处理,内存也与值的种类数成正比:不明白这个区别,就会变成“我明明写成了流式,怎么还是崩溃”。 而且“我运行过了,可以啊”不是证明,它只是说在那个时候、那个文件上可以。设定上限并通过,才是证明,而且结果必须以退出码而不是文字的形式留下,下一个人才不用再问。 评分器不会相信你的文字。它在临时目录里摆好评分器生成的日志,在内存上限之下真实运行你的工具,并核对结果。条数和金额每次运行都会变。
步骤
- 创建并运行 /root/stream/gen_events.py,生成 /root/stream/data/events.csv。共 30 万行。
- 在 /root/stream/stream.py 中实现
agg <파일>(占位符为文件),让它一行一行地读取,给出条数、合计、平均值、最小值、最大值。 - 增加
top <파일> <N>(占位符依次为文件、N),让它用堆只保留前 N 条。 - 增加
distinct <파일> <칼럼>(占位符依次为文件、列名),让它统计唯一值并同时给出当时的最大内存;然后比较种类数少的列和种类数多的列,写入 /root/stream/distinct.json。 - 创建 /root/stream/cap.py 和 /root/stream/slurp.py,在 64MiB 的内存上限之下,把两个版本并排运行,并把结果写入 /root/stream/limit.json。
- 增加
sortmerge <파일> <칼럼> <청크행수>(占位符依次为文件、列名、数据块行数),让它通过拆分与排序、归并来做外部排序。 - 增加
verify <파일> <칼럼>(占位符依次为文件、列名),让它给出是否有序以及与顺序无关的指纹,并对原始文件和排序后的文件进行核对,写入 /root/stream/verify.json。 - 把全部内容整理成一份报告,生成 /root/stream/stream_report.json 和 /root/stream/stream_report.md。
参考
- 运行契约:
python3 /root/stream/stream.py <명령> <파일> [인자](占位符依次为命令、文件、参数)。命令有 agg、top、distinct、sortmerge、verify 五个。成功时退出码为 0,文件不存在时为 3,命令或参数个数有误时为 2。 agg的响应:rows、sum_amount、mean_amount、min_amount、max_amount、peak_kib。平均值四舍五入到小数点后第二位。peak_kib是该进程使用的最大内存,通过resource.getrusage(resource.RUSAGE_SELF).ru_maxrss获得。在 Linux 上单位是 KiB。top的响应:{"n": 정수, "rows": [[event_id, amount], ...], "peak_kib": 정수}(占位符依次为整数、整数)。按金额降序排列,金额相同时 event_id 大的排在前面。distinct的响应:column、rows、exact、peak_kib。sortmerge的响应:rows、chunks、chunk_rows、output、peak_kib。数据块写入输入文件所在目录的chunks/,结果连同表头一起写入同一目录的sorted.csv。verify的响应:rows、ordered、sum_amount、digest、peak_kib。digest 是把去掉表头之后的每一行的sha256(줄 전체 바이트)(韩文,意为“整行的全部字节”)的前 8 个字节按整数读取并全部相加的值,对 2 的 64 次方取余,以 16 位小写十六进制输出。顺序不同也必须得到相同的值。- 排序键在值是整数的样子时为
(정수, 원래 문자열)(占位符依次为整数、原字符串),否则为(0, 원래 문자열)(占位符为原字符串)。这条规则是本实验的假设。 python3 /root/stream/cap.py <MiB> <명령> [인자...](占位符依次为 MiB 数、命令、参数)在地址空间上限之下运行该命令,并输出{"limit_mib": 정수, "argv": [...], "exit_code": 정수, "ok": true|false}(占位符依次为整数、整数)。即使子进程崩溃,cap.py 自己也以 0 结束。python3 /root/stream/slurp.py <파일>(占位符为文件)是把文件整个读入并放进列表的版本,是为了在上限之下让它失败而做的对照组。- 官方文档:heapq · resource · hashlib · POSIX sort
- 常见错误:以
read()或readlines()开头;把全部数据排序后再取前面的部分;为了验证把原始文件和结果都放进列表;不删除数据块文件。 - Pod 的资源是 2 个 CPU 核、2Gi 内存和 6Gi 临时磁盘。不要生成更大的文件。
生成一个无法一次读完的文件
创建并运行 /root/stream/gen_events.py,生成 /root/stream/data/events.csv。表头是 event_id,shop_id,kind,amount,ts,共 30 万行。
生成时也不要攒进列表,而要一行一行地写进文件。shop_id 的种类数要少,event_id 要让每一行都不同,这样之后才能看到统计唯一值时的内存差异。
一行一行地读取并汇总
在 /root/stream/stream.py 中实现 agg <파일>(占位符为文件),让它给出 rows、sum_amount、mean_amount、min_amount、max_amount、peak_kib。
直接把文件对象放进 for 语句,就是一行一行地读。一旦调用 read() 或 readlines(),这个性质就消失了。最小值和最大值只需与刚读到的值比较,平均值用合计和条数在最后算出。
手里只拿着前 N 条
增加 top <파일> <N>(占位符依次为文件、N),让它以 [[event_id, amount], ...] 的形式给出金额前 N 条。按金额降序排列,金额相同时 event_id 大的排在前面。
把全部数据排序后取前面的部分,整个数据都会加载到内存。维护一个大小为 N 的最小堆,堆里有 N 个之后,只在新值比堆底更大时才推入。heapq 的 heappushpop 一次就能完成。放进堆的值用(金额、event_id)这一对,平局的处理也就一并解决了。
统计唯一值时,内存与什么成正比
增加 distinct <파일> <칼럼>(占位符依次为文件、列名),让它给出 column、rows、exact、peak_kib。并把 shop_id 和 event_id 两列的结果,以 rows、low、high 写入 /root/stream/distinct.json。
即使一行一行地读,也必须记住已经见过的值,所以内存与值的种类数成正比。分别测量种类数少的列和每行都不同的列,这种差异就会以数字的形式看到。low 放种类数少的一方的结果,high 放种类数多的一方的结果。
设定上限来证明
创建 /root/stream/cap.py 和 /root/stream/slurp.py,在 64MiB 的上限之下分别运行 stream.py agg 和 slurp.py,并把 limit_mib、file_bytes、rows、stream_exit、slurp_exit 写入 /root/stream/limit.json。
Linux 对进程能占用的地址空间有上限,在 Python 中可以用 resource 模块来设置。用 subprocess 启动子进程时,要在子进程那一侧设定上限,父进程才不会一起崩溃。即使子进程崩溃,cap.py 自己也以 0 结束,只是以 JSON 告知子进程的退出码。
用拆分与排序、归并做外部排序
增加 sortmerge <파일> <칼럼> <청크행수>(占位符依次为文件、列名、数据块行数),让它对每个数据块排序并写入 chunks/,再合并后连同表头一起写入同一目录的 sorted.csv。响应为 rows、chunks、chunk_rows、output、peak_kib。
反复进行这样的操作:攒够数据块行数的数据就排序,导出为文件,然后清空缓冲区。归并时,把数据块文件全部打开,但每个文件每次只取一行,heapq.merge 接收排序键并替你完成这件事。开始时要把之前生成的数据块文件删掉:不删的话,下一次运行时结果会混在一起。
不靠中间产物来验证
增加 verify <파일> <칼럼>(占位符依次为文件、列名),让它给出 rows、ordered、sum_amount、digest、peak_kib,并对原始文件和排序后的文件进行核对,把 column、source、sorted、same_multiset 写入 /root/stream/verify.json。
把原始文件和结果都放进列表来比较,流式排序的意义就没有了。给每一行算出哈希并全部相加,顺序不同也会得到相同的值,而只要缺了一行或者重复了一行,就会不同。是否有序,只需要记住上一行的键就能知道。
把证明留成一份
在 /root/stream/stream_report.json 中写入 file_bytes、rows、sum_amount、chunks、limit_mib、stream_exit、slurp_exit、digest_match,并在 /root/stream/stream_report.md 中分 ## 무엇을 받았나 ## 어떻게 처리했나 ## 정말로 스트리밍인가 ## 남은 한계 四节书写(韩文标题,依次意为“收到了什么”“如何处理”“真的是流式处理吗”“剩余局限”)。
读取前面步骤留下的 limit.json、verify.json 和 distinct.json 并合并即可。报告里要用数字写明上限值和两个退出码:那就是证明;在剩余局限一节,要写明统计唯一值的内存与种类数成正比这一事实。