VE472复习站 / 项目1

项目1:百万歌曲基础数据与紧凑图BFS

本页先说明五个数据处理阶段,再说明里程碑2。四种实现使用同一张有向无权图,并遵守同一套广度优先搜索(BFS)语义。优化实现使用紧凑图和位图。

来源、证据等级与主要结论

主要代码来源固定为仓库 .research/p1team01 的以下Git提交(提交):

09582fbeac6e5dfaf3efc19680ee2ae09b706f36

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

仓库事实
可在上述Git提交的规格、实现、测试或性能记录中直接定位。
可验证推导
由图论、BFS或Spark执行模型直接推出;可用小图或执行计划复核。
历史记录
仓库保留的某次实验数值。只有口径完全相同时才能横向比较。
未给定
缺少PDF、外部原始证据(原始证据)或运行清单(清单),无法在当前工作区唯一确认。
复习要点
问题答案
图是什么?artist_id → similar_artists[i] 生成的有向、无权、去重边。
BFS保证什么?第一次发现顶点时的深度是从源点出发的最少边数。
如何公平比较四种组合?固定输入语义、源点、深度、边定义、停止条件和资源口径。只改变计算框架或文件格式。
正式MR与旧MR的核心差别是什么?正式MapReduce(MR)实现每层直接扫描Avro或Parquet基础数据。旧MR实现读取预先物化的制表符分隔值(TSV)边表。
cache() 是否立即计算?否。cache() 只登记持久化意图。Spark动作才触发物化。
D1 / D2是什么?D1每层启动一个Spark作业。D2在一个作业的执行器中完成整次位图BFS。

项目1的五阶段数据流

术语:HDF5是分层数据格式第五版。Avro是按行组织并携带模式的数据序列化格式。Parquet是面向分析的列式存储格式。Drill是可查询多种数据源的分布式结构化查询语言(SQL)查询引擎。基础数据(基础数据)是阶段0生成的统一逻辑数据集。

仓库事实:系统按阶段0至阶段4组织。阶段2包含两种召回方法:艺人图BFS和98维歌曲向量的近似最近邻(ANN)检索。两种方法解决不同问题。阶段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                          │
└──────────────────────────────────────────┘

阶段0:建立统一基础数据

仓库事实:songs_foundation 同时写为Avro和Parquet。两种物理格式必须使用同一个逻辑模式,并且记录语义必须相同。仓库README记录了1,000,000行数据和零次失败。主要字段包括歌曲标识、艺人标识、标量音频特征、timbre_90similar_artists 和标签数组。

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

阶段1:使用Drill查询

对Avro和Parquet执行同一组课程查询。验证两种格式返回相同的排序结果。仓库在后续分析中选择Parquet。阶段2仍然比较四种正式实现。

阶段2:本页重点

必做部分包括艺人距离BFS和四种实现的比较。优化部分使用紧凑图代替重复扫描基础数据或边表。位图核心算法按精确层数扩展当前边界。阶段2还包含歌曲向量ANN,但歌曲向量ANN不能代替艺人图BFS。

阶段3和阶段4

阶段3预测发行年份。阶段4合并图召回候选和向量召回候选,然后执行过滤和排序。仓库明确声明,发行年份预测单独报告,不计入当前阶段4的排名分数。

需要区分:艺人距离、歌曲余弦相似度和发行年份预测不是同一个指标。三个任务的输入、目标和验证指标不同。

里程碑2的共同约束:先固定问题,再比较实现

仓库约束:任务2只处理艺人距离BFS和格式与框架比较。任务2不负责建立新的歌曲推荐索引,也不负责产品演示。

输入

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有向路径的最少边数不可达时为 null
path重建的艺人ID路径可达时,第一个元素必须是起点,最后一个元素必须是目标。
visited_count已发现艺人数计数口径必须跨四路线一致。
iteration_count执行的BFS深度轮数不是Spark阶段数的同义词。

完成条件

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

1.边生成规则

仓库事实:每条基础数据记录中的 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。只有源数据也含反向关系时,反向边才存在。
  • 无权:每条边的代价恒为1。距离表示边数,不表示音频相似度。
  • 去重:相同 (src,dst) 只能保留一次。否则,当前边界计数、混洗数据量和性能结果都会失真。
  • 允许自环但不产生新发现:仓库规格没有明确要求删除 A → A。已访问集合会过滤该自环。

需要区分: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 中每个顶点与源点的最短距离恰为 d−1
  2. 已访问集合不变量:d 轮结束后,visited 恰好包含距离不超过 d 的已发现顶点。
  3. 互斥不变量:next_frontier ∩ visited_before = ∅。同一轮多个父节点发现同一目标时,目标也只能进入一次。
  4. 前驱不变量:首次发现 v 时记录的前驱位于上一层。因此,沿前驱回溯会得到长度为发现深度的路径。

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

可验证推导:BFS按0、1、2、… 跳的顺序扩展。如果顶点 v 首次在深度 d 出现,但存在长度小于 d 的路径,则该路径的倒数第二个顶点会在更早的层中扩展,并更早发现 v。这与首次在深度 d 发现 v 矛盾。因此,首次发现深度就是最少跳数。

分布式实现对应关系

抽象SparkMapReduce位图内核
当前层frontier DataFrame分布式缓存中的 frontier 文件long[] frontier
已访问visited DataFramevisited 文件long[] visited
扩展frontier 与边表连接扫描基础数据,由映射器匹配源点按64位字对邻接位图执行OR
过滤已访问顶点left_anti映射器过滤 visitednext & ~visited
轮内去重dropDuplicates归约器按艺人键去重位OR操作具有幂等性

4.精确层计数与累计层计数不能混用

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

范围内(D) = Σd=0D 精确[d],且 精确[0]=1

需要区分:“三层内的人数”是累计量。“第三层新发现的人数”是单层量。仓库的正式约束保存 exact_new_by_hop。不得用累计值与该字段比较。

MR/Spark × Avro/Parquet四组合

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

路线正式输入语义迭代边界主要成本
MapReduce + Avro 每一跳使用 AvroKeyInputFormat 直接扫描Avro基础数据。 每一跳启动一个独立的YARN应用和作业。 重复读取全表、启动作业、混洗数据和执行归约。
MapReduce + Parquet 每一跳使用 AvroParquetInputFormat 直接扫描Parquet基础数据。 每一跳启动一个独立的YARN应用和作业。 重复扫描数据,并承担适配器和输入格式成本。
Spark + Avro 读取Avro基础数据。执行一次边投影和去重,并缓存结果。 在同一个Spark应用中循环执行多跳。 首次投影、混洗、状态物化和每轮Spark动作。
Spark + Parquet 读取Parquet基础数据。执行一次边投影和去重,并缓存结果。 在同一个Spark应用中循环执行多跳。 列裁剪、首次投影、混洗和状态物化。

公平比较的固定项

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”只表示数据来源,不证明MR原生读取Avro的性能。

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

相同运行时间不证明算法相同。相同最终累计数也不能排除中间层错误。仓库的正式证据比较每一跳的新边界顶点数。当前性能记录说明,五名艺人、四种方法和十跳产生的200条“方法—跳数”记录全部一致。

正式直接读取基础数据的实现与旧脚本

1.正式MapReduce:每一跳直接扫描基础数据

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

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

每轮把 frontier.txtvisited.txt 作为缓存文件分发。映射器读取基础数据中的 artist_idsimilar_artists。仅当源点位于当前边界时,映射器才输出未访问邻居。归约器按目标艺人键去重。

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

结果元数据强制:

  • input_semantics = direct_foundation_scan_per_hop
  • tsv_fallback_used = false
  • Avro/Parquet的输入-格式类与路线一致;
  • 深度为 D 时必须存在 D 个不同应用程序ID;
  • exact_new_by_hop 必须等于冻结的正确性序列。

2.正式Spark:直接投影基础数据并缓存一次边表

仓库事实:Spark启动器的运行清单写入:

input_semantics = direct_foundation_projection_cached_then_dataframe_bfs
timing_scope = python_main_start_through_spark_stop

Spark在同一个应用中读取指定格式,并清洗和去重边。然后Spark调用 cache()。默认模式 eager_count 先物化边缓存,再开始逐跳循环。frontiervisited 默认通过HDFS写入后重读来物化。

3.旧基准测试脚本:格式标签只代表边的来源

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

  • 供Spark使用的Parquet边表;
  • 供Hadoop Streaming使用的TSV边文本。

旧MapReduce路线的真实输入是 avro_tsvparquet_tsv。Python映射任务/归约任务在Hadoop Streaming中处理这些TSV边,而不是原生读取基础数据。

问题正式基础数据直读旧脚本
MR实际读取Avro/Parquet基础数据预物化TSV边
MR实现含依赖打包Java JAR + 原生InputFormatHadoop Streaming + Python
边构建是否计入每跳数的基础数据扫描天然在查询内常作为独立预处理;须另行声明是否计时
格式结论可比较MR原生格式读取主要比较派生边来源与Streaming路径
证据强度清单、输入类、无TSV后备方案、应用程序IDs旧日志可保留,但不能替代正式契约

需要区分:旧脚本并非“错误代码”。它可以验证图语义和探索工程瓶颈;问题在于它不能回答正式的原生MR × physical格式比较。

Spark:延迟求值、缓存与血缘关系

术语:Spark转换只构建计算计划。Spark动作触发计算。cache() 请求持久化计算结果。血缘关系记录数据集的上游计算依赖。

1. cache() 是声明,不是Spark动作

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

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

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

可验证推导:若跳过 count(),第一个需要边的动作会一并承担首次计算成本。这样“BFS query时间”包含多少准备,取决于计时边界。仓库正式清单记录边物化模式,避免暗中改变口径。

需要区分:第二次Spark动作能否完全复用缓存,取决于分区是否已经物化以及分区是否仍在缓存中。调用 cache() 不表示数据会永久驻留。

2.多个Spark动作会重复触发未缓存的上游计算

简化版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

第一个动作会物化 next_frontier;后续动作可以复用它,但 visited 若只保留长血缘关系,仍可能重算上游。性能分析必须看动作、作业与阶段,而不能只数Python循环次数。

3.血缘关系为什么会增长

每一轮都构造:

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

若不截断血缘关系,第 d 轮的逻辑计划包含前 d−1 轮的并集/去重链。长血缘关系会增加计划、调度、序列化与故障重算成本。

仓库支持的物化策略

状态选项效果
eager_count / lazy_cache决定缓存是否在BFS前预热。
当前边界hdfs_roundtrip / spark_checkpoint切断当前层血缘;前者写后重读,后者交给Spark检查点。
已访问hdfs_roundtrip / lineage决定累计集合是否每轮物化。

历史记录:仓库记录Spark 4.1 DataFrame检查点在该作业中遇到血缘关系问题,最终脚本采用每跳数将当前边界与已访问写入HDFS再读回的方式。

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

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

历史记录:当前性能文档中G4 “已缓存混合边” 为40.630 s,G2为40.172 s,被标为 NO EFFECT;G5预热边缓存为41.033 s,也没有改善。

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

紧凑图:从字符串边表到内存映射数组

1.为什么需要紧凑图

基础数据/Parquet适合存储与分析,却不是重复低延迟BFS的理想查询结构。每跳数扫列式表、展开数组、连接当前边界和混洗,会让调度与数据移动成本远大于位运算本身。

仓库事实:优化路径从规范基础数据Parquet派生紧凑产物。README记录一个50,996,454-字节产物,包含354,307艺人、约4.441百万有向边与998,795艺人-到歌曲链接。

注意:性能文档另列格式-v1稀疏位图为344,151艺人。两个艺人计数在同一提交中不一致,表明它们不是可无条件合并的同一产物;详见性能冲突章节。

2.稠密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

仓库事实:索引记录是小端序 <IB,共5字节;载荷值是小端序 uint32。因此每行偏移量上限为232−1,计数上限为255。

紧凑邻接表示意图
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.复用相同的连续邻接片段

仓库事实:ContainmentArrayWriter 先收集每个来源的有序、无重复行,再按“长度降序、字典序”处理唯一行。若某行已作为载荷中连续子序列存在,则索引复用该偏移量;否则追加。

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

physical payload:
  (3, 1, 0)

logical_count = 5
physical_count = 3

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

需要区分:计数只有8位。单个艺人若有256个邻居,写入器必须抛出溢出;不能静默截断。

4. mmap与启动校验

Java CompactGraph.open 会:

  1. 读取元数据;
  2. 验证每个文件的实际字节与元数据相同;
  3. 读取艺人/歌曲ID列表并核对计数;
  4. 内存映射索引与载荷;
  5. 逐行验证偏移量/计数不越界,稠密ID不超目录。

这些检查是性能数字的前置条件。若图校验和、文件大小或艺人计数不匹配,不能把运行称为同一产物的复现。

5.位图核心算法

仓库事实:Java核心算法为艺人全集分配三个 long[] 位集:

  • visited:所有已发现艺人;
  • frontier:当前精确跳数集合;
  • 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个候选的集合并、差与计数可用一个机器字的OR、AND-NOT与 bitCount 完成;整数稠密ID避免大量字符串哈希/连接。

路径开销:计数-仅查询不分配距离/前驱;需要路径时才分配 byte[] distancesint[] predecessors。正式单次契约当前是 counts-only-v1,不能声称它输出了路径。

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

D1 / D2:相同的BFS,不同的Spark作业边界

仓库事实:正式紧凑单次同时测D1与D2,且两者必须返回相同的 exact_new_by_hop

语义实现深度D的作业数主要开销
D1 engine.sequential:先创建状态,再每跳数启动一个Spark作业推进。 D+1 每层调度、状态序列化、驱动程序↔执行器桥接。
D2 engine.fused:把完整请求放入一个分区,在执行器内一次运行位图核心算法。 1 一次调度与一次桥接,内核在单执行器完成。
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

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

资源契约:正式单次命令对D2把 SPARK_EXECUTORS 固定为1;否则多执行器配置会制造“分配了但未参与一个单分区查询”的歧义。

需要区分:D1/D2不是图深度,也不是数据格式。它们是同一紧凑/Java后端的执行边界语义。

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

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

仓库权威性能规则:

  • 输入、来源、跳数、集群布局、资源限制与应用程序生命周期必须相同。
  • 慢实验和失败实验也要记录为 REGRESSIONINVALID
  • 小于正常抖动的差异标为 NO EFFECT
  • 使用计时前,先核对当前边界计数、行计数、模式与输出哈希。
  • 单次Spark与持久化热点服务是不同工作负载类,不得当作等价延迟。
  • 每个数字必须能追溯到命令、提交、Spark应用程序ID、配置快照与原始资源日志。

2.当前权威四路线记录

历史记录:pre/performance.md 明确自称当前权威性能记录。五名艺人各跑十跳数,得到:

路线5运行均值标准差范围
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。正确解读是:“在此工作负载与资源口径下,它的均值最低”;不是“Parquet对所有BFS永远最快”。

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

主题README权威性能文档/其他记录复习结论
四路线10-跳数均值 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,且表中为五个艺人运行。 数字与重复次数/实验代际不同。正式报告应选带完整当前口径与原始证据路径的性能文档,不能求平均或拼表。
艺人计数 紧凑产物:354,307艺人。 稀疏位图格式-v1:344,151艺人。 视为不同产物生成/过滤口径,直到清单证明相同。不得把一个计数配上另一个产物的延迟。
基础数据大小 Avro 2.923 GiB,Parquet 2.509 GiB。 性能表四舍五入为2.9与2.5 GiB。 这是精度差异,不是实质矛盾;报告时保留来源精度。

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

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

类别示例生命周期
正式四条路线Spark+Parquet十跳数均值69.270 s完整应用程序/正式基础数据输入
匹配最终工作负载五艺人十跳数:Parquet 54.718 s,位图38.373 s匹配输入与结果,但属于优化代际比较
持久化服务Java位图3-跳数1.296 ms,10-跳数12.706 msJVM、执行器、mmap与JIT已热
README热点后端批量1,000请求,31.070 QPS批量热点服务吞吐,不是单次全新墙钟

毫秒级核心算法/服务延迟与几十秒全新应用程序墙钟同时成立,因为它们计时边界不同。核心算法只测图遍历;全新墙钟还含提交、资源分配、JVM、Spark上下文、产物分发与结果写出。

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墙钟更低但核心算法相同,收益主要来自作业融合与桥接减少。若精确计数不同,则该计时直接失效。

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

本页生成时实测:在提交 09582fbeac6e5dfaf3efc19680ee2ae09b706f36 的仓库虚拟环境中执行五个任务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工作负载、三跳/十跳、精确跳数数组形状与非负性。
test_task2_bfs_validation.py2四路线完整性、一致距离、缺失路线检测。
test_task2_compact_writer.py3后缀复用、uint8溢出、艺人→歌曲来源。
test_task2_expansion_benchmark.py2四路线扩展结果接受/拒绝规则。
合计18全部通过。

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

高频陷阱清单

  1. 把有向的相似艺人边自动补成双向边。
  2. 按来源单列去重,误删合法邻居。
  3. 用加权最短路径解释无权BFS。
  4. 把精确跳数与范围内跳数混为一个计数。
  5. 只比较最终计数,不核对每层当前边界。
  6. 把 “由Avro派生TSV” 写成 “MR原生Avro”。
  7. 认为 cache() 立即执行,忽略首次动作成本。
  8. 认为缓存会截断血缘关系;实际检查点/写后重读才建立新的可靠边界。
  9. 把Spark循环轮数当作作业/阶段数。
  10. 把D1/D2当作数据集或图深度。
  11. 把单次秒级墙钟与热点核心算法毫秒级延迟直接作倍率比较。
  12. 从README、性能文档和旧日志中挑最快数字拼成一张“最佳结果”表。
  13. 声称紧凑计数-仅单次返回完整路径。
  14. 把本页伪代码称为仓库Scala源码;该提交中可核验的实现语言是Python与Java。

自测题与答案

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

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

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

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

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

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

答案:X只能进入 next_frontier 一次。Spark用目标列去重,MapReduce归约任务按键唯一化,位图OR天然幂等。任取一个首次层父节点都能构成最短路径。

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

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

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

答案:旧路线先从Avro基础数据派生TSV边,Hadoop Streaming实际读取TSV。Avro只表示数据来源。正式路线必须让作业直接用 AvroKeyInputFormat 扫基础数据,并证明未使用TSV后备方案。

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

答案:取决于声明的计时契约。cache() 是延迟;若未先动作,首个BFS动作会承担投影与缓存物化。正式比较应记录 lazy_cacheeager_count,并固定准备是否计入。

7.为什么已访问血缘关系会随深度增长?

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

8.紧凑索引的 <uint32 offset,uint8 count> 有什么限制?

答案:载荷偏移量不能超过232−1;单行计数不能超过255。写入器对256个邻居必须报溢出,不能截断。

9.包含关系复用为何不改变BFS结果?

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

10.位图核心算法如何实现 next \ visited

答案:对每个64-位块执行 next[word] &= ~visited[word],再用OR加入已访问,并以 bitCount 统计精确跳数新节点。

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

答案:D1为11:一个初始化作业加十个逐跳数作业。D2为1:完整十跳在单个执行器任务内融合执行。

12. D2是否绕过Spark/YARN?

答案:否。Spark/YARN仍负责应用程序、执行器与任务。D2只是不把单个内存BFS的每一层拆成独立Spark作业。

13. 1.296 ms位图核心算法与69.270 s Spark+Parquet能直接算“快多少倍”吗?

答案:不能。前者来自持久化、JIT/mmap热状态的核心算法/服务;后者是正式基础数据路线的完整应用程序工作负载。输入结构、生命周期与计时边界不同。

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

答案:当前仓库将 pre/performance.md 声明为权威记录,且它给出工作负载、运行、方差、范围与原始证据路径。README的另一组数应保留为历史摘要,不能混合使用。

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

答案:不能。这18项测试证明局部BFS、契约、验证器和紧凑写入器行为。在线性能还需要HDFS/YARN、真实InputFormat、应用程序ID、资源快照与原始剖析证据。

考前最终检查