# 大数据工程：真实组件最小实验

配套专题：https://blog.traceof.me/topics/big-data

完整源码：https://blog.traceof.me/examples/big-data-lab.zip

这是本地教学实验，数据由固定种子生成。博客只分发源码、图解和结果记录，线上服务器不运行 Kafka、Spark 或 Flink。不要将工作目录指向业务数据、写作数据库或已有重要文件。

## 环境与版本

已实跑环境：macOS arm64，Amazon Corretto JDK 17.0.15，Python 3.12，Maven 3.9.9；机器物理内存 64 GiB。实验使用 Kafka 堆 512 MiB、Spark driver 2 GiB、Flink JobManager 768 MiB、TaskManager 1536 MiB，顺序运行阶段；建议为本地实验准备至少 8 GiB 可用内存及 5 GiB 可用磁盘。资源建议不是已验证的最低要求。

- Kafka 3.9.1（Scala 2.13 发行包），单节点 KRaft，3 个输入分区，副本数 1。
- Spark 3.5.7（Scala 2.12/Hadoop 3 发行包），`local[4]`。
- Iceberg Spark 3.5/Scala 2.12 runtime 1.10.0，本地 Hadoop catalog。
- Flink 1.20.3（Scala 2.12 发行包），并行度 2，4 个 slot。
- Flink Kafka Connector 3.4.0-1.20，Kafka client 3.9.1。
- Jackson core/annotations/databind 2.18.3，打包时重定位到 `lab.shaded.jackson`，避免运行时依赖混用。

官方下载 URL、官方摘要与实测 SHA-256 位于 `versions.json`。bootstrap 校验固定 SHA-256 和官方摘要后解压；依赖仅下载到指定工作目录，Maven 使用标准缓存和本包独立的 Central 镜像配置，不修改用户 Maven 设置。压缩包不携带大型发行版、JAR 或原始运行日志。

仅 macOS arm64 已完成实际运行。Linux 路径通过 `LAB_JAVA_HOME` 指定，但本次没有验证 Linux；Windows 原生不在本次覆盖范围内。脚本使用 Bash 与 Python 3.12+。脚本显式指定本次下载的 Spark/Flink 路径，不继承全局 SPARK_HOME、CLASSPATH 和 Java 注入选项；Kafka 不启用额外 JMX 监听。`run.py` 保留 Flink 官方 Java 17 `--add-opens`/`--add-exports` 设置。

## 完整运行

解压后进入 `big-data-lab` 目录，确认 `python3 --version` 为 3.12 或更高，`mvn` 可用。若系统默认 Python 较旧，通过环境变量指定：

```bash
export LAB_PYTHON=/absolute/path/to/python3.12
# macOS 默认自动选择已安装 JDK 17；也可以显式指定。
export LAB_JAVA_HOME=/absolute/path/to/jdk-17
bash all.sh /absolute/path/to/fresh-lab-work
```

`all.sh` 拒绝复用已经有实验数据的工作目录，不自动清空历史结果。它下载并校验依赖、编译 Java、执行 8 项参考语义测试、启动组件、按章节验证，最后停止本次组件。默认端口为 Kafka 19092/19093、Flink RPC 16123、REST 18081，启动前检查占用；不要停止不属于本实验的服务来抢占端口。

## 按章节运行

以下命令假定已执行 bootstrap 和 Maven package；用同一个 Python 3.12+ 解释器执行所有步骤。正常步骤只运行一次；重复接入会追加事件，重复建表不会自动清空旧表。需要完整复现时选择新的工作目录。

```bash
python3 bootstrap.py --work /absolute/path/to/lab-work
mvn -B -s java/settings.xml -f java/pom.xml -Dlab.build.directory=/absolute/path/to/lab-work/java-target package
python3 check_reference.py
python3 run.py setup --work /absolute/path/to/lab-work
python3 run.py kafka --work /absolute/path/to/lab-work
python3 run.py spark --work /absolute/path/to/lab-work
python3 run.py streaming --work /absolute/path/to/lab-work
python3 advanced.py /absolute/path/to/lab-work
python3 run.py stop --work /absolute/path/to/lab-work
```

`setup` 启动独立回环服务；`kafka` 生成输入、生产并捕获分区结束位置导出；`spark` 写入 Iceberg、对照 CSV/Parquet、验证字段与分区演进及历史快照；`streaming` 对照 100 个窗口，停止并重启 TaskManager，验证迟到修订；`advanced` 执行百万条倾斜 Join 对照，并删除、恢复实验日期的日报，连续补数两次。

SQL 模板以 `${LAB_WORK}` 标记工作目录，主要用于阅读。实际执行入口会将路径和当前快照 ID 填入生成的 SQL；不要直接复制本次运行的快照 ID 到新表。为避免共享数据互相影响，所有入口应顺序执行。

## 输入与输出约定

一条 `play` 表示一次教学播放事件；`event_id` 稳定标识身份，重复传输使用相同 ID。`user_id`、`video_id` 为合成标识；`event_time`、`ingest_time` 为 UTC 毫秒；`watch_ms` 为整数毫秒且不能为负。参考计算器遇到同 ID 异内容直接失败；在线示例只演示同 ID 同内容去重，不承诺处理更正事件。

默认输入为 10,000 条唯一合法事件、100 条原样重复、1 条非法时长。事件在五秒小组内打乱，集中在短时间范围，供快速正确性验证。边界测试另覆盖上海午夜、窗口边界、身份冲突、乱序与并列排序。

日报按 Asia/Shanghai 自然日；实时采用不重叠的 5 分钟事件时间窗口、5 秒乱序余量、30 秒允许迟到、3 秒空闲检测。去重 ID 在事件时间一天后清理。`tick` 是推进实验时间的控制记录，广播至各输入分区，不参与播放量；真实业务不能随意用此方式宣称数据已完整。

Flink 将窗口绝对值与诊断记录写入 Kafka，启用 Checkpoint 与事务 sink；读取端使用 `read_committed`。窗口键为时间与视频，同键以较新记录覆盖，不累加修订值。Top 10 按次数降序、视频 ID 升序从最终快照生成，不是额外部署的实时榜单服务。

工作目录保留 `data/`、`warehouse/`、`checkpoint-*`、`spark-events/`、`logs/` 和 `results.json`。日志与 Checkpoint 可包含机器绝对路径，公开文章只使用整理后的结果摘要。

## 本次实际结果

- Kafka 显式暂停后 poll 为 0、积压为 10,101，恢复后导出 10,101 行；另一次普通导出 10,101 行，从已提交位置恢复为 0 行，主动回放为 10,101 行；去重合法数 10,000。
- Spark 日报与 5 分钟窗口逐组匹配独立参考；CSV 数据文件 550,999 字节，Parquet 114,085 字节。
- 新增 device 字段后 10,000 行保持存在且新字段为空；演进前快照仍可读取 10,000 行。
- Flink 100 个视频窗口匹配离线参考，输出 100 条重复诊断、1 条非法诊断。
- TaskManager 停止、重启后记录一次 Checkpoint 恢复；窗口依次输出 2/200 ms 与 3/300 ms，过迟事件 d 被分流；完整离线为 4/400 ms。
- 百万条热点样例两种 Join 结果一致。单次根 SQL 耗时为 896/308 ms，Shuffle 写出 5,550,118/430 字节；不能当作跨机器、多轮统计基准。
- 删除实验日期日报后为 0 个分组，两次补数均恢复相同 100 个分组，与参考一致。

`sample-events.jsonl` 与 `sample-expected.json` 提供八条可读小样例及独立答案。`results.json` 是本次记录；每次复现的时间、快照 ID 和运行身份可能不同。正确性按结果内容核对，性能没有硬编码通过阈值。单机、副本 1、进程重启实验没有证明多副本容灾、网络分区、磁盘损坏或 JobManager 高可用。

## 主要文件

- `generate.py`：确定性数据与独立参考聚合。
- `java/src/main/java/lab/KafkaTool.java`：生产、有界导出和恢复读取；显式分区分配，不测试消费组再均衡。
- `java/src/main/java/lab/StreamJob.java`：解析、去重、窗口、迟到侧输出、Checkpoint 与 Kafka 事务输出。
- `run.py`、`advanced.py`：真实组件实验与断言。
- `check_reference.py`：8 个语义边界测试；`spark-*.sql`：本次 SQL 模板。
- `01-platform.svg` 至 `06-reconciliation.svg`：六篇技术图。

示例代码按 MIT 许可证提供，见 LICENSE。下载的 Apache 项目及 Maven 依赖遵循各自许可证，本包不重新分发它们的二进制。
