VE472 复习站 / Project 1

Project 1:从百万歌曲 Foundation 到 Compact BFS

先用五阶段数据流定位全局,再集中掌握 Milestone 2:同一张有向无权图、同一套 BFS 语义、四条正式课程路线,以及 compact/bitmap 优化。

来源、证据等级与最小结论

主要代码来源固定为仓库 .research/p1team01 的 commit:

09582fbeac6e5dfaf3efc19680ee2ae09b706f36

本页同时参考 projects/avro_parquet_sdd.htmlprojects/sdd_guide.html。请求中列出的 projects/p1.pdf 在当前工作区不存在,因此本页不会声称已从该 PDF 核对原题措辞。仓库 spec 中对 “p1.pdf requirement” 的转述会标为“仓库契约”,而不是 PDF 直接引文。

仓库事实
可在上述 commit 的 spec、实现、测试或性能记录中直接定位。
可验证推导
由图论、BFS 或 Spark 执行模型直接推出;可用小图或执行计划复核。
历史记录
仓库保留的某次实验数值。只有口径完全相同时才能横向比较。
未给定
缺少 PDF、外部 raw evidence 或 manifest,无法在当前工作区唯一确认。
考试最小记忆集
问题最短答案
图是什么?artist_id → similar_artists[i] 生成的有向、无权、去重边。
BFS 保证什么?第一次发现顶点时的深度是从源点出发的最短 hop 数。
四组合为什么公平?同一输入语义、源点、深度、边定义、停止条件与资源口径,只改变框架/格式。
正式 MR 与旧 MR 的核心差别?正式 MR 每 hop 原生扫描 Avro/Parquet foundation;旧 MR 读预物化 TSV 边。
cache() 是否立即计算?否。它只登记持久化意图;action 才物化。
D1 / D2 是什么?D1 每层跨一个 Spark job;D2 把完整 bitmap BFS 融合为一个 job 内的 executor 计算。

Project 1 五阶段数据流

仓库事实:系统以 Stage 0–4 组织。Stage 2 有两条 recall 路径:artist BFS 与 98D song-vector ANN;两者回答不同问题,只在 Stage 4 合并。

数据流图:阶段编号表示仓库中的执行阶段
1,000,000 个 MSD HDF5
          │
          ▼
┌──────────────────────────────────────────┐
│ Stage 0  typed foundation                │
│ 同一 logical schema → Avro + Parquet     │
└────────────┬───────────────┬─────────────┘
             │               │
       ┌─────▼─────┐   ┌────▼──────────────────────────────┐
       │ Stage 1   │   │ Stage 2                           │
       │ Drill SQL │   │ A. directed artist BFS            │
       │ 四个查询  │   │ B. 98D vectors → ANN/HNSW recall  │
       └─────┬─────┘   └────────────┬──────────────────────┘
             │                      │
             │               ┌──────▼──────┐
             │               │ Stage 3     │
             │               │ release-year│
             │               │ prediction  │
             │               └──────┬──────┘
             │                      │(独立报告,不进入当前排序分数)
             └──────────────┬───────┘
                            ▼
┌──────────────────────────────────────────┐
│ Stage 4 recommendation                   │
│ BFS candidates + ANN candidates          │
│ → merge/dedup → SQL policy filter        │
│ → weighted score → deterministic lottery │
│ → ranked Top 10                          │
└──────────────────────────────────────────┘

Stage 0:统一数据地基

仓库事实:songs_foundation 同时落为 Avro 与 Parquet。两种物理格式必须来自同一 logical schema、具有相同记录语义。仓库 README 记录 1,000,000 行、零失败;主要字段包含歌曲/艺人标识、标量音频特征、timbre_90similar_artists 与标签数组。

SDD 要点:不要把 Avro 与 Parquet 写成两套业务表。比较格式时,字段、空值规则和行集必须一致,否则测到的是数据差异,不是格式差异。

Stage 1:Drill 查询

对 Avro 与 Parquet 执行同一组课程查询,并验证排序结果一致。仓库后续分析路径选择 Parquet,但这不取消 Stage 2 的四路线正式比较。

Stage 2:本页主体

必做部分是 artist-distance BFS 和四组合比较。优化部分把 repeated foundation/edge scan 替换为 compact graph,并用 bitmap kernel 执行 exact-hop 扩展。Stage 2 还包含 song-vector ANN,但它不是 artist BFS 的替代品。

Stage 3 与 Stage 4

Stage 3 预测发行年份。Stage 4 合并图召回与向量召回候选,再过滤和排序。仓库明确声明:year prediction 单独报告,不进入当前 Stage 4 排名分数。

考试陷阱:不要把 artist distance、song cosine similarity 和 release-year prediction 合并成一个“相似度”。三者的输入、目标与验证指标不同。

Milestone 2 总契约:先固定问题,再比较实现

仓库契约:Task 2 只解决 artist-distance BFS 与格式/框架比较,不负责另建歌曲推荐索引或产品演示。

输入

hdfs:///p1/full/avro/songs_foundation/
hdfs:///p1/full/parquet/songs_foundation/

所需字段是 artist_idartist_namesimilar_artists。显示字段可以包含 track_idsong_idtitlerelease,但这些字段不改变 BFS 图。

查询输出

字段语义边界
source_artist_id起点必须是清洗后的非空 ID。
target_artist_id目标与起点相同时距离为 0。
distance最少有向 hop 数不可达时为 null
path重建的 artist ID 路径可达时首尾必须分别为 source/target。
visited_count已发现艺人数计数口径必须跨四路线一致。
iteration_count执行的 BFS 深度轮数不是 Spark stage 数的同义词。

完成条件

有向、无权、去重图与 BFS 不变量

1. 边生成规则

仓库事实:每条 foundation 记录中的 artist_id 是源点;similar_artists 数组中的每个非空元素生成一条出边。

foundation row:
  artist_id = "A"
  similar_artists = ["B", "C", "B", null, " "]

cleaned directed edges:
  ("A", "B")
  ("A", "C")

仓库 src/2/foundation_graph.py 中的可执行 PySpark 边投影为:

from pyspark.sql import functions as F

edges = (
    foundation.select(
        F.trim(F.col("artist_id")).alias("src_artist_id"),
        F.explode_outer("similar_artists").alias("dst_artist_id"),
    )
    .select(
        "src_artist_id",
        F.trim(F.col("dst_artist_id")).alias("dst_artist_id"),
    )
    .where("src_artist_id IS NOT NULL AND src_artist_id <> ''")
    .where("dst_artist_id IS NOT NULL AND dst_artist_id <> ''")
    .dropDuplicates(["src_artist_id", "dst_artist_id"])
)

说明:以上是从仓库函数中抽出的可执行 PySpark(只省略读取与按源点重分区),不是 Scala 源码。仓库实现为 PySpark 与 Java;本页不会虚构 Scala 实现。

  • 有向:A → B 不推出 B → A。只有源数据也含反向关系时,反向边才存在。
  • 无权:每条边 hop cost 恒为 1。距离是边数,不是音频相似度。
  • 去重:相同 (src,dst) 只能保留一次;否则 frontier count、shuffle 量和性能都会失真。
  • 允许自环但不产生新发现:spec 没有显式要求删除 A → A。visited 集会挡住它。

考试陷阱:src_artist_id 单列去重会错误删除一个艺人的其他邻居;必须对二元组 (src,dst) 去重。

2. 队列版 BFS 伪代码

BFS(G, source, target, maxDepth):
    visited := {source}
    predecessor[source] := NONE
    frontier := [source]

    if source = target:
        return distance 0, path [source]

    for depth := 1 .. maxDepth:
        next := empty sequence
        for u in frontier:
            for each directed neighbor v in G[u]:
                if v not in visited:
                    visited.add(v)
                    predecessor[v] := u
                    next.append(v)
                    if v = target:
                        return depth, reconstruct(predecessor, target)
        if next is empty:
            break
        frontier := next

    return unreachable

3. 四个核心不变量

  1. 层不变量:进入第 d 轮时,frontier 中每个顶点与 source 的最短距离恰为 d−1
  2. visited 不变量:d 轮结束后,visited 恰包含距离不超过 d 的已发现顶点。
  3. 互斥不变量:next_frontier ∩ visited_before = ∅。同一轮多个父节点发现同一目标时,目标也只能进入一次。
  4. 前驱不变量:首次发现 v 时记录的 predecessor 位于上一层,因此沿 predecessor 回溯必然得到长度为发现深度的路径。

为什么第一次发现就是最短路

可验证推导:BFS 按 0、1、2、… hop 的顺序扩展。若顶点 v 首次在深度 d 出现,却存在长度小于 d 的路径,则该路径的倒数第二个顶点应在更早层被扩展,并更早发现 v,矛盾。因此首次发现深度就是最短 hop 数。

分布式实现对应关系

抽象SparkMapReduceBitmap kernel
当前层frontier DataFrame分布式 cache 中的 frontier 文件long[] frontier
已访问visited DataFramevisited 文件long[] visited
扩展frontier join edges扫描 foundation,mapper 匹配 source按 64-bit word OR 邻接位图
去旧点left_antimapper 过滤 visitednext & ~visited
轮内去重dropDuplicatesreducer 对 artist key 唯一化bit OR 天然幂等

4. Exact-hop 与 within-hop 不能混用

exact[d] 表示恰在第 d 跳首次发现的节点数,则:

within(D) = Σd=0D exact[d],且 exact[0]=1

考试陷阱:“三跳内人数”是累计量;“第三跳新发现人数”是单层量。仓库 formal contract 保存 exact_new_by_hop,不能拿累计值与它比较。

MR/Spark × Avro/Parquet 四组合

仓库契约:四路线必须解决同一个 BFS,而不是四个“差不多”的任务。

路线正式输入语义迭代边界主要成本
MapReduce + Avro 每 hop 用 AvroKeyInputFormat 直接扫描 Avro foundation。 一 hop 一个独立 YARN application/job。 重复全表读、job 启动、shuffle/reduce。
MapReduce + Parquet 每 hop 用 AvroParquetInputFormat 直接扫描 Parquet foundation。 一 hop 一个独立 YARN application/job。 重复扫描与 adapter/input-format 成本。
Spark + Avro 读取 Avro foundation,投影/去重边一次并 cache。 同一 Spark application 内循环多 hop。 首次 projection、shuffle、状态物化与每轮 action。
Spark + Parquet 读取 Parquet foundation,投影/去重边一次并 cache。 同一 Spark application 内循环多 hop。 列裁剪、首次 projection、shuffle 与状态物化。

公平比较的固定项

fixed:
  foundation row set
  graph cleaning and directed-edge definition
  source artist(s)
  depth / target / stop condition
  exact-hop correctness contract
  cluster topology and resource envelope
  application lifecycle class

varied:
  framework ∈ {MapReduce, Spark}
  foundation format ∈ {Avro, Parquet}

考试陷阱:文件格式标签必须说明“实际被谁读取、在何时读取”。如果 MR 实际读的是从 Avro 派生的 TSV,称它为 “MR + Avro” 只表示 provenance,不证明 MR 的 Avro reader 性能。

为什么四条路线都要做正确性对齐

相同 runtime 不证明相同算法;相同最终累计数也不能排除中间层错误。仓库的正式证据比较每一 hop 的新 frontier 数。当前性能记录称五名艺人、四种方法、十跳的 200 个 method-hop 记录全部一致。

正式 direct-foundation 与旧脚本的差异

1. 正式 MapReduce:每 hop 直接扫描 foundation

仓库事实:FoundationMrJob 根据格式选择原生输入类:

if (format.equals("avro")) {
    job.setInputFormatClass(AvroKeyInputFormat.class);
} else if (format.equals("parquet")) {
    job.setInputFormatClass(AvroParquetInputFormat.class);
}

每轮把 frontier.txtvisited.txt作为 cache files 分发。mapper 读取 foundation 的 artist_idsimilar_artists;只有 source 位于当前 frontier 时才发出未访问邻居。reducer 以目标 artist 为 key 去重。

for depth = 1..D:
    write frontier and visited
    launch native Avro/Parquet MapReduce job
    scan every foundation record
    mapper: if record.artist_id in frontier:
                emit each nonempty unvisited similar_artist
    reducer: emit each destination artist once
    next := job output
    visited := visited ∪ next

结果 metadata 强制:

  • input_semantics = direct_foundation_scan_per_hop
  • tsv_fallback_used = false
  • Avro/Parquet 的 input-format class 与路线一致;
  • 深度为 D 时必须存在 D 个不同 application ID;
  • exact_new_by_hop 必须等于冻结的正确性序列。

2. 正式 Spark:直接投影 foundation,一次缓存边

仓库事实:Spark launcher 的 manifest 写入:

input_semantics = direct_foundation_projection_cached_then_dataframe_bfs
timing_scope = python_main_start_through_spark_stop

它在同一个 application 中读取指定格式、清洗并去重边,然后 cache()。默认使用 eager_count 先物化 edge cache,再开始 hop 循环;frontier 与 visited 默认经 HDFS roundtrip 物化。

3. 旧 benchmark 脚本:格式标签只代表边的来源

历史实现:run-task2-expansion-benchmark.sh 先调用 build_edges.py,从 Avro 或 Parquet foundation 生成:

  • 供 Spark 使用的 Parquet edge table;
  • 供 Hadoop Streaming 使用的 TSV edge text。

旧 MapReduce 路线的真实输入是 avro_tsvparquet_tsv。Python mapper/reducer 在 Hadoop Streaming 中处理这些 TSV 边,而不是原生读取 foundation。

问题正式 direct-foundation旧脚本
MR 实际读取Avro/Parquet foundation预物化 TSV edges
MR 实现shaded Java JAR + native InputFormatHadoop Streaming + Python
边构建是否计入每 hop 的 foundation scan 天然在查询内常作为独立预处理;须另行声明是否计时
格式结论可比较 MR 原生格式读取主要比较派生边 provenance 与 Streaming 路径
证据强度manifest、input class、无 TSV fallback、application IDs旧日志可保留,但不能替代正式契约

考试陷阱:旧脚本并非“错误代码”。它可以验证图语义和探索工程瓶颈;问题在于它不能回答正式的 native MR × physical format 比较。

Spark:lazy evaluation、cache 与 lineage

1. cache() 是声明,不是 action

edges = read_foundation_edges(...).cache()

# lazy_cache:到这里还没有读取/去重完整 foundation

edges.count()
# eager_count:触发 projection、dedup、shuffle,并填充 cache

可验证推导:若跳过 count(),第一个需要 edges 的 action 会一并承担首次计算成本。这样“BFS query time”包含多少 preparation,取决于计时边界。仓库正式 manifest 记录 edge materialization mode,避免暗中改变口径。

考试陷阱:第二次 action 是否完全命中 cache 取决于分区是否成功物化且未被驱逐。写了 cache() 不等于数据永久驻留。

2. 多个 action 会重复触发未缓存血缘

简化版 Spark BFS 常在同一轮调用:

next_frontier = expand(...).cache()
hit = next_frontier.filter(is_target).limit(1).collect()  # action
next_count = next_frontier.count()                        # action
visited_count = visited.count()                           # action

第一个 action 会物化 next_frontier;后续 action 可以复用它,但 visited 若只保留长 lineage,仍可能重算上游。性能分析必须看 action、job 与 stage,而不能只数 Python 循环次数。

3. 为什么 lineage 会增长

每一轮都构造:

visited_d = dedup(visited_(d-1) UNION frontier_d)

若不截断 lineage,第 d 轮的逻辑计划包含前 d−1 轮的 union/dedup 链。长 lineage 会增加计划、调度、序列化与故障重算成本。

仓库支持的物化策略

状态选项效果
edgeseager_count / lazy_cache决定缓存是否在 BFS 前预热。
frontierhdfs_roundtrip / spark_checkpoint切断当前层血缘;前者写后重读,后者交给 Spark checkpoint。
visitedhdfs_roundtrip / lineage决定累计集合是否每轮物化。

历史记录:仓库记录 Spark 4.1 DataFrame checkpoint 在该作业中遇到 lineage 问题,最终脚本采用每 hop 将 frontier 与 visited 写入 HDFS再读回的方式。

权衡:HDFS roundtrip 增加 I/O,却建立明确迭代边界并避免长 lineage。是否更快必须测量;它首先是一项可复现性与稳定性决策。

4. 缓存优化为何可能没有效果

历史记录:当前性能文档中 G4 “cached hybrid edges” 为 40.630 s,G2 为 40.172 s,被标为 NO EFFECT;G5 warm edge cache 为 41.033 s,也没有改善。

可能原因包括:边读取已不是主瓶颈;额外 action/物化抵消收益;cache 已在先前路径中复用;测量差异落在正常噪声内。正确结论是“在该匹配 workload 下无测得效果”,不是“cache 永远无用”。

Compact graph:从字符串边表到内存映射数组

1. 为什么需要 compact graph

Foundation/Parquet 适合存储与分析,却不是重复低延迟 BFS 的理想查询结构。每 hop 扫列式表、explode 数组、join frontier 和 shuffle,会让调度与数据移动成本远大于位运算本身。

仓库事实:优化路径从 canonical foundation Parquet 派生 compact artifact。README 记录一个 50,996,454-byte artifact,包含 354,307 artists、约 4.441 million directed edges 与 998,795 artist-to-song links。

注意:性能文档另列 format-v1 sparse bitmap 为 344,151 artists。两个 artist count 在同一 commit 中不一致,表明它们不是可无条件合并的同一 artifact;详见性能冲突章节。

2. Dense ID 与六个数据文件

artist_ids.txt          dense artist index → original artist_id
song_ids.txt            dense song index   → original song_id

artist_index.bin        per artist: <uint32 offset, uint8 count>
artist_neighbors.u32    packed dense neighbor IDs

song_index.bin          per artist: <uint32 offset, uint8 count>
artist_songs.u32        packed dense song IDs

仓库事实:index record 是 little-endian <IB,共 5 bytes;payload value 是 little-endian uint32。因此每行 offset 上限为 232−1,count 上限为 255。

Compact 邻接表示意图
artist_index.bin
dense artist 0 ── (offset=0, count=3) ─────────┐
dense artist 1 ── (offset=1, count=2) ──────┐  │
                                             ▼  ▼
artist_neighbors.u32 payload:               [3, 1, 0]
artist 0 neighbors = payload[0:3] = [3,1,0]
artist 1 neighbors = payload[1:3] =   [1,0]

第二行复用第一行的后缀,不再复制 [1,0]。

3. Containment compaction

仓库事实:ContainmentArrayWriter 先收集每个 source 的有序、无重复 row,再按“长度降序、字典序”处理唯一 row。若某 row 已作为 payload 中连续子序列存在,则 index 复用该 offset;否则追加。

logical rows:
  row 0 = (3, 1, 0)
  row 1 = (1, 0)

physical payload:
  (3, 1, 0)

logical_count = 5
physical_count = 3

可验证推导:复用连续子序列保持 row 的顺序与内容完全不变,因此 BFS 邻居枚举语义不变。它不是集合近似压缩,也不允许漏边。

考试陷阱:count 只有 8 bit。单个 artist 若有 256 个邻居,writer 必须抛出 overflow;不能静默截断。

4. mmap 与启动校验

Java CompactGraph.open 会:

  1. 读取 metadata;
  2. 验证每个文件的实际 bytes 与 metadata 相同;
  3. 读取 artist/song ID 列表并核对 count;
  4. 内存映射 index 与 payload;
  5. 逐行验证 offset/count 不越界,dense ID 不超 catalog。

这些检查是性能数字的前置条件。若 graph checksum、文件大小或 artist count 不匹配,不能把运行称为同一 artifact 的复现。

5. Bitmap kernel

仓库事实:Java kernel 为 artist universe 分配三个 long[] 位集:

  • visited:所有已发现 artist;
  • frontier:当前 exact-hop 集合;
  • next:下一层候选。
frontier := {source}
visited  := {source}
exact[0] := 1

for depth = 1..maxDepth:
    clear(next)
    for each set artist bit in frontier:
        for each 64-bit adjacency word (block, word):
            discovered := word AND NOT visited[block] AND NOT next[block]
            next[block] := next[block] OR word
            optionally record distance/predecessor for discovered bits

    for each block:
        next[block]    := next[block] AND NOT visited[block]
        visited[block] := visited[block] OR next[block]
        exact[depth]  += popcount(next[block])

    swap(frontier, next)

为什么快:64 个候选的集合并、差与计数可用一个 machine word 的 OR、AND-NOT 与 bitCount 完成;整数 dense ID 避免大量字符串 hash/join。

路径开销:counts-only 查询不分配 distances/predecessors;需要路径时才分配 byte[] distancesint[] predecessors。正式 one-shot contract 当前是 counts-only-v1,不能声称它输出了路径。

边界:底层 kernel 接受深度 0..15;正式 one-shot CLI 只允许 1..10。接口契约比内核能力更窄,调用者必须遵守 CLI 契约。

D1 / D2:相同 BFS,不同 Spark job 边界

仓库事实:正式 compact one-shot 同时测 D1 与 D2,且两者必须返回相同的 exact_new_by_hop

语义实现深度 D 的 job 数主要开销
D1 engine.sequential:先创建状态,再每 hop 启动一个 Spark job推进。 D+1 每层调度、状态序列化、driver↔executor bridge。
D2 engine.fused:把完整请求放入一个 partition,在 executor 内一次运行 bitmap kernel。 1 一次调度与一次 bridge,内核在单 executor 完成。
D1 与 D2 执行边界
D1, depth = 3
Job 0: initialize state
   │ serialize/collect
Job 1: hop 1
   │ serialize/collect
Job 2: hop 2
   │ serialize/collect
Job 3: hop 3

D2, depth = 3
Job 0: executor opens resident compact graph
       hop 1 → hop 2 → hop 3
       return exact counts once

可验证推导:compact graph 约几十 MB,单次 BFS 是内存位运算。将一个小查询拆成多个 Spark job,调度和序列化可能比计算更贵。D2 不代表“不使用 Spark”:它仍由 Spark/YARN 启动 application、分配 executor 和执行 task,只是把查询内核融合在一个 task 内。

资源契约:formal one-shot command 对 D2 把 SPARK_EXECUTORS 固定为 1;否则多 executor 配置会制造“分配了但未参与一个单 partition 查询”的歧义。

考试陷阱:D1/D2 不是图深度,也不是数据格式。它们是同一 compact/Java backend 的执行边界语义。

性能口径:先证明可比,再读数字

1. 正式比较必须固定的测量合同

仓库权威性能规则:

  • 输入、source、hop、集群布局、资源限制与 application lifecycle 必须相同。
  • 慢实验和失败实验也要记录为 REGRESSIONINVALID
  • 小于正常抖动的差异标为 NO EFFECT
  • 使用 timing 前,先核对 frontier counts、row counts、schema 与 output hash。
  • one-shot Spark 与 persistent hot service 是不同 workload class,不得当作等价延迟。
  • 每个数字必须能追溯到 command、commit、Spark application ID、配置快照与 raw resource log。

2. 当前权威四路线记录

历史记录:pre/performance.md 明确自称当前 authoritative performance record。五名艺人各跑十 hop,得到:

路线5 runs meanstddev范围
Spark + Avro74.530 s5.832 s69.141–83.384 s
Spark + Parquet69.270 s5.094 s62.608–74.036 s
MapReduce + Avro361.465 s7.826 s349.455–368.697 s
MapReduce + Parquet357.052 s9.295 s348.700–371.936 s

在该受控实验中选择 Spark + Parquet。正确解读是:“在此 workload 与资源口径下,它的均值最低”;不是“Parquet 对所有 BFS 永远最快”。

3. 同一 commit 中的历史数字冲突

主题README权威性能文档/其他记录复习结论
四路线 10-hop 均值 MR+Avro 354.608 s;MR+Parquet 327.892 s;Spark+Avro 113.335 s;Spark+Parquet 112.942 s。README 还称每艺人三次独立运行。 pre/performance.md 为 361.465、357.052、74.530、69.270 s,且表中为五个 artist run。 数字与重复次数/实验代际不同。正式报告应选带完整当前口径与 raw evidence 路径的性能文档,不能求平均或拼表。
artist count compact artifact:354,307 artists。 sparse bitmap format-v1:344,151 artists。 视为不同 artifact generation/过滤口径,直到 manifest 证明相同。不得把一个 count 配上另一个 artifact 的 latency。
Foundation size Avro 2.923 GiB,Parquet 2.509 GiB。 性能表四舍五入为 2.9 与 2.5 GiB。 这是精度差异,不是实质矛盾;报告时保留来源精度。

未给定:外部 raw evidence 位于 /4721/project/perf_runs/,不在 Git 中;本页无法仅凭仓库重建每个历史数字的完整运行链。因此只报告仓库声明,不把它们重新认证为本机实测。

4. 这些数字为什么不能放在同一柱状图

类别示例生命周期
正式 four-routeSpark+Parquet 十 hop均值 69.270 s完整 application/正式 foundation 输入
matched final workload五艺人十 hop:Parquet 54.718 s,bitmap 38.373 s匹配输入与结果,但属于优化代际比较
persistent serviceJava bitmap 3-hop 1.296 ms,10-hop 12.706 msJVM、executor、mmap 与 JIT 已热
README hot backend batch1,000 requests,31.070 QPS批量 hot 服务吞吐,不是单次 fresh wall

毫秒级 kernel/service latency 与几十秒 fresh application wall 同时成立,因为它们计时边界不同。kernel 只测图遍历;fresh wall 还含提交、资源分配、JVM、Spark context、artifact 分发与结果写出。

5. D1/D2 与性能归因

比较 D1/D2 时应至少同时记录:

launcher_process_wall_nanos
query_wall_nanos
kernel_nanos
executor_bridge_nanos
job_ids
execution_steps
resource_envelope
exact_new_by_hop

若 D2 wall 更低但 kernel 相同,收益主要来自 job 融合与 bridge 减少。若 exact counts 不同,则该 timing 直接失效。

验证状态:所列五个测试文件的 18 项测试

本页生成时实测:在 commit 09582fbeac6e5dfaf3efc19680ee2ae09b706f36 的仓库虚拟环境中执行五个 Task 2 测试文件:

.venv/bin/pytest -q \
  tests/test_task2_bfs.py \
  tests/test_task2_bfs_contract.py \
  tests/test_task2_bfs_validation.py \
  tests/test_task2_compact_writer.py \
  tests/test_task2_expansion_benchmark.py
..................  [100%]
18 passed in 0.03s
所列五个测试文件的 18 项测试覆盖面
测试文件数量覆盖重点
test_task2_bfs.py2最短路径与不可达结果。
test_task2_bfs_contract.py9W1/W2 workload、三跳/十跳、exact-hop 数组形状与非负性。
test_task2_bfs_validation.py2四路线完整性、一致 distance、缺失路线检测。
test_task2_compact_writer.py3后缀复用、uint8 overflow、artist→song provenance。
test_task2_expansion_benchmark.py2四路线 expansion 结果接受/拒绝规则。
合计18全部通过。

边界:这些是快速单元/契约测试,不等于在线四节点集群、HDFS、YARN、native InputFormat 与性能 raw evidence 的全量重跑。

高频陷阱清单

  1. 把 directed similar-artist edge 自动补成双向边。
  2. 按 source 单列去重,误删合法邻居。
  3. 用 weighted shortest path 解释无权 BFS。
  4. 把 exact hop 与 within hop 混为一个计数。
  5. 只比较最终 count,不核对每层 frontier。
  6. 把 “Avro-derived TSV” 写成 “MR native Avro”。
  7. 认为 cache() 立即执行,忽略首次 action 成本。
  8. 认为 cache 会截断 lineage;实际 checkpoint/写后重读才建立新的可靠边界。
  9. 把 Spark 循环轮数当作 job/stage 数。
  10. 把 D1/D2 当作数据集或图深度。
  11. 把 one-shot 秒级 wall 与 hot kernel 毫秒级 latency 直接作倍率比较。
  12. 从 README、性能文档和旧日志中挑最快数字拼成一张“最佳结果”表。
  13. 声称 compact counts-only one-shot 返回完整 path。
  14. 把本页伪代码称为仓库 Scala 源码;该 commit 的相关实现是 Python 与 Java。

自测题与答案

先口述不变量或数据边界,再展开答案。

1. A.similar_artists=[B] 能否推出边 B→A

答案:不能。正式图是有向图,只生成 A→B。只有另一条 foundation 关系显式给出 B→A 时才存在反向边。

2. 为什么第一次发现 target 时可以立即返回最短距离?

答案:BFS 按 hop 层递增扩展。首次在第 d 层发现 target,所有小于 d 的层已完整扩展;若存在更短路径,target 应已更早出现,矛盾。

3. 同一轮两个 parent 都发现顶点 X,应如何处理?

答案:X 只能进入 next_frontier 一次。Spark 用目标列去重,MapReduce reducer 按 key 唯一化,bitmap OR 天然幂等。任取一个首次层 parent 都能构成最短路径。

4. “第三跳 100 人”和“三跳内 100 人”是否等价?

答案:不等价。前者是 exact[3];后者是 exact[0]+exact[1]+exact[2]+exact[3]

5. 为什么旧 MR+Avro 路线不能证明 native Avro InputFormat 性能?

答案:旧路线先从 Avro foundation 派生 TSV edges,Hadoop Streaming 实际读取 TSV。Avro 只表示数据 provenance。正式路线必须让 job 直接用 AvroKeyInputFormat 扫 foundation,并证明未使用 TSV fallback。

6. Spark 对 edges 调用 cache() 后立刻开始计时 BFS,公平吗?

答案:取决于声明的 timing contract。cache() 是 lazy;若未先 action,首个 BFS action 会承担 projection 与 cache 物化。正式比较应记录 lazy_cacheeager_count,并固定 preparation 是否计入。

7. 为什么 visited lineage 会随深度增长?

答案:每轮的 visited_d 都由 visited_(d−1) UNION frontier_d 再去重生成。若不 checkpoint 或写后重读,第 d 轮计划递归包含此前所有轮的 union/dedup。

8. Compact index 的 <uint32 offset,uint8 count> 有什么限制?

答案:payload offset 不能超过 232−1;单行 count 不能超过 255。writer 对 256 个邻居必须报 overflow,不能截断。

9. Containment 复用为何不改变 BFS 结果?

答案:它只让多个逻辑 row 指向同一段完全相同的连续 payload 子序列。每个 artist 枚举到的邻居序列没有改变,因此边集合和 BFS 层计数不变。

10. Bitmap kernel 如何实现 next \ visited

答案:对每个 64-bit block 执行 next[word] &= ~visited[word],再用 OR 加入 visited,并以 bitCount 统计 exact-hop 新节点。

11. 深度为 10 时,D1 与 D2 分别应产生多少 Spark jobs?

答案:D1 为 11:一个初始化 job 加十个逐 hop job。D2 为 1:完整十跳在单个 executor task 内融合执行。

12. D2 是否绕过 Spark/YARN?

答案:否。Spark/YARN 仍负责 application、executor 与 task。D2 只是不把单个内存 BFS 的每一层拆成独立 Spark job。

13. 1.296 ms bitmap kernel 与 69.270 s Spark+Parquet 能直接算“快多少倍”吗?

答案:不能。前者来自 persistent、JIT/mmap 热状态的 kernel/service;后者是正式 foundation 路线的完整 application workload。输入结构、生命周期与计时边界不同。

14. README 与 pre/performance.md 的四路线均值不同,应引用哪个?

答案:当前仓库将 pre/performance.md 声明为 authoritative record,且它给出 workload、runs、方差、范围与 raw evidence 路径。README 的另一组数应保留为历史摘要,不能混合使用。

15. 所列五个测试文件的 18 项测试全过,是否证明线上四节点性能数字正确?

答案:不能。这 18 项测试证明局部 BFS、contract、验证器和 compact writer 行为。在线性能还需要 HDFS/YARN、真实 InputFormat、application ID、资源快照与 raw profile 证据。

考前最终检查