题面
实现教学版 HLLLite。64 位哈希值的最高 p 位选择桶。剩余 q=64−p 位中,第一个 1 之前的零数加一为 rho。实现以下接口:
- 实现
offerHashed和add。 - 实现兼容性检查和原地
merge。 - 实现未经修正的 HLL 估计值。小基数时使用线性计数修正。
- 说明 MapReduce 合并器和归约器为什么可以逐寄存器取最大值。
哈希函数由训练题提供。只有当两个摘要的 p 和哈希标识同时相同时,才允许执行 merge。
概率型数据结构手册 / Scala 算法
本页包含三套 45 分自测训练。内容覆盖概率型流式结构、线性代数、图算法、Scala 集合和 Spark 弹性分布式数据集(RDD)。45 分是本页定义的训练分值,不是官方评分标准,也不表示官方考试承诺。
术语:HyperLogLog 基数估计器(HyperLogLog,HLL)使用寄存器估计不同元素数量。对称正定矩阵(SPD)满足对称性和严格正定性。广度优先搜索(BFS)按路径边数逐层访问图。主成分分析(PCA)用于线性降维。奇异值分解(SVD)把矩阵分解为两个正交因子和一个奇异值矩阵。Spark 的映射器(mapper)、归约器(reducer)、驱动进程(driver)、执行进程(executor)、广播变量(broadcast variable)、惰性求值(lazy evaluation)和混洗(shuffle)后文固定使用这些中文名称或对应代码标识符。
使用边界:本页根据课程材料安排训练顺序。本页不预测任何考试内容。当前材料不包含正式题面,也不包含官方 45 分评分规则。
| 等级 | 材料依据 | 可使用的结论 |
|---|---|---|
| A | 当前材料提供完整接口、算法或计算要求。 | 可以据此建立完整实现训练。 |
| B | 当前材料提供可执行流程,但没有提供完整接口。 | 可以据此建立推导或半实现训练。 |
| C | 历史训练材料提供同类代码。 | 只能称为历史训练证据。 |
| D | 扩展资料提供相邻主题。 | 只能作为补充训练,不表示官方要求。 |
| 顺序 | 训练内容 | 依据等级 | 安排理由 |
|---|---|---|---|
| 1 | HyperLogLog(HLL) | C | prob_struct.html 提供完整的短 Scala 实现。该实现适合训练状态、位运算、merge 和估计器。材料没有把它声明为官方考核要求。 |
| 2 | Cholesky | A+B | 当前材料提供 SPD 和 Cholesky 分解的公式、循环不变量和失败条件。 |
| 3 | 无权有向图 BFS | A | 当前材料提供最短距离、路径和分布式当前边界的语义。 |
| 4 | Scala Collection / Spark RDD 聚合 | A | 当前材料提供 Scala、WordCount 和 flatMap → map → reduceByKey 数据流。 |
| 5 | Spark 批量梯度下降 | B | 当前材料提供梯度下降的数据流和通信要求,适合半实现训练。 |
| 6 | PCA/SVD Gram | B | 当前材料提供 Gram 矩阵、SVD、PCA 和高而窄矩阵聚合。 |
| 扩展 | 蓄水池抽样(Reservoir Sampling)/计数最小摘要(Count-Min Sketch,CMS)/布隆过滤器(Bloom filter) | D | 可用于读码迁移;不应挤占前六项训练时间。 |
三套参考实现采用 Scala 2.13 语法。Spark 自测训练固定使用 Scala 2.12.18 和 Apache Spark 3.5.1,因为 Spark 构件带有 Scala 二进制版本后缀。不得混合编译两个版本的类文件。
// Array 与二维数组
val a = Array(1.0, 2.0)
val matrix = Array.ofDim[Double](3, 3)
val copy = a.clone()
// 现场算法常用 mutable 容器
import scala.collection.mutable
val q = mutable.Queue.empty[(String, Int)]
val seen = mutable.Set.empty[String]
val parent = mutable.Map.empty[String, String]
// Option
def lookup(k: String): Option[Int] =
if (k.nonEmpty) Some(k.length) else None
// 迭代器适合一次流过
val words = lines.iterator
.flatMap(_.toLowerCase.split("[^a-z0-9]+"))
.filter(_.nonEmpty)
Double。merge 和归约操作。确认操作满足算法要求的结合性。实现教学版 HLLLite。64 位哈希值的最高 p 位选择桶。剩余 q=64−p 位中,第一个 1 之前的零数加一为 rho。实现以下接口:
offerHashed 和 add。merge。哈希函数由训练题提供。只有当两个摘要的 p 和哈希标识同时相同时,才允许执行 merge。
final class HLLLite(
val p: Int,
val hashId: String,
hash64: String => Long
) {
def add(value: String): Unit
def offerHashed(hash: Long): Unit
def merge(that: HLLLite): Unit
def estimate(): Double
def registers: Vector[Int]
}
4 ≤ p ≤ 20,m=2p。offerHashed 每次只更新一个 register,并且 register 只能增大。merge 对不同 p/hashId 必须抛 IllegalArgumentException。| 部分 | 自测分值 | 可观察检查项 |
|---|---|---|
| 状态与构造 | 7 | p 检查、m、零寄存器、hash 标识。 |
| hash 拆分与 rho | 14 | 无符号右移取 bucket;rho 有全零 remainder 上界;max 更新。 |
| merge | 8 | 兼容性;逐项 max;交换/结合/幂等。 |
| estimate | 9 | 调和和、alpha、Double、小基数 linear counting。 |
| 复杂度与不变量 | 7 | add O(1),merge/estimate O(m),state O(m),MapReduce 合并理由。 |
| 总分 | 45 |
以下单文件不依赖第三方库,可保存为 MockHLL.scala,用 Scala 2.13 编译。
final class HLLLite(
val p: Int,
val hashId: String,
hash64: String => Long
) {
require(p >= 4 && p <= 20, "p must be in [4, 20]")
require(hashId.nonEmpty, "hashId must be non-empty")
private val m: Int = 1 << p
private val q: Int = 64 - p
private val reg: Array[Int] = Array.fill(m)(0)
def add(value: String): Unit = offerHashed(hash64(value))
def offerHashed(hash: Long): Unit = {
val bucket = (hash >>> q).toInt
val shiftedRemainder = hash << p
val rawRho = java.lang.Long.numberOfLeadingZeros(shiftedRemainder) + 1
val rho = math.min(rawRho, q + 1)
if (rho > reg(bucket)) reg(bucket) = rho
}
def merge(that: HLLLite): Unit = {
require(this.p == that.p, "precision mismatch")
require(this.hashId == that.hashId, "hash function mismatch")
val other = that.registers
var i = 0
while (i < m) {
reg(i) = math.max(reg(i), other(i))
i += 1
}
}
def estimate(): Double = {
var inversePowerSum = 0.0
var zeroRegisters = 0
var i = 0
while (i < m) {
inversePowerSum += math.pow(2.0, -reg(i).toDouble)
if (reg(i) == 0) zeroRegisters += 1
i += 1
}
val alpha = m match {
case 16 => 0.673
case 32 => 0.697
case 64 => 0.709
case _ => 0.7213 / (1.0 + 1.079 / m.toDouble)
}
val raw = alpha * m.toDouble * m.toDouble / inversePowerSum
if (raw <= 2.5 * m && zeroRegisters > 0)
m.toDouble * math.log(m.toDouble / zeroRegisters.toDouble)
else
raw
}
def registers: Vector[Int] = reg.toVector
}
object MockHLL {
private def hashed(p: Int, bucket: Int, rho: Int): Long = {
val q = 64 - p
require(bucket >= 0 && bucket < (1 << p))
require(rho >= 1 && rho <= q)
(bucket.toLong << q) | (1L << (q - rho))
}
def main(args: Array[String]): Unit = {
val noHash: String => Long = _ => 0L
val left = new HLLLite(4, "test-v1", noHash)
left.offerHashed(hashed(4, 0, 2))
left.offerHashed(hashed(4, 0, 4))
left.offerHashed(hashed(4, 3, 1))
assert(left.registers(0) == 4)
assert(left.registers(3) == 1)
val right = new HLLLite(4, "test-v1", noHash)
right.offerHashed(hashed(4, 1, 5))
right.offerHashed(hashed(4, 3, 3))
left.merge(right)
assert(left.registers(0) == 4)
assert(left.registers(1) == 5)
assert(left.registers(3) == 3)
val before = left.registers
left.offerHashed(hashed(4, 3, 3))
assert(left.registers == before)
val empty = new HLLLite(4, "test-v1", noHash)
assert(empty.estimate() == 0.0)
println("MockHLL tests passed")
}
}
| 输入 | 预期 | 检查目的 |
|---|---|---|
| 同 bucket 的 rho 2、4 | register 保留 4 | max 不变量。 |
| bucket 3 的 rho 1 | 只更新 register 3 | bucket 拆分。 |
| merge 两个 sketch | 逐项 max | union summary。 |
| 重复 offer 同一 hash | 状态不变 | 幂等。 |
| 空 sketch | estimate = 0 | linear counting。 |
| p/hashId 不同 | 拒绝 merge | 兼容性。 |
add 为 O(1);merge 与 estimate 为 O(m)。M[j] 是所有落入 bucket j 的观测中最大 rho。def cholesky(
a: Array[Array[Double]],
eps: Double = 1e-12
): Array[Array[Double]]
def shortestPath(
adjacency: Map[String, Seq[String]],
source: String,
target: String,
maxDepth: Int
): Option[List[String]]
source==target 返回单点路径;缺失邻接视为无出边;不可达为 None;负深度拒绝。| 部分 | 自测分值 | 检查项 |
|---|---|---|
| Cholesky 公式与循环 | 10 | 对角 residual、sqrt、下方 dot product/division。 |
| 验证与复杂度 | 5 | square/symmetry/SPD、Θ(n³)。 |
| 测试与重构 | 5 | 3×3 例、非 SPD 失败、LLᵀ。 |
| BFS 状态 | 14 | queue、入队即 seen、parent、first discovery。 |
| 边界与回溯 | 6 | 深度、同点、不可达、环。 |
| Spark frontier 语义 | 5 | join、dedup、left-anti、visited union。 |
以下单文件不依赖第三方库,可保存为 MockLinearGraph.scala。
import scala.collection.mutable
object MockLinearGraph {
def cholesky(
a: Array[Array[Double]],
eps: Double = 1e-12
): Array[Array[Double]] = {
require(eps >= 0.0, "eps must be nonnegative")
val n = a.length
require(a.forall(_.length == n), "square matrix required")
var i = 0
while (i < n) {
var j = 0
while (j < n) {
require(
math.abs(a(i)(j) - a(j)(i)) <= eps,
"symmetric matrix required"
)
j += 1
}
i += 1
}
val l = Array.ofDim[Double](n, n)
var column = 0
while (column < n) {
var diagonalSum = 0.0
var k = 0
while (k < column) {
diagonalSum += l(column)(k) * l(column)(k)
k += 1
}
val residual = a(column)(column) - diagonalSum
require(residual > eps, "matrix is not SPD")
l(column)(column) = math.sqrt(residual)
i = column + 1
while (i < n) {
var cross = 0.0
k = 0
while (k < column) {
cross += l(i)(k) * l(column)(k)
k += 1
}
l(i)(column) =
(a(i)(column) - cross) / l(column)(column)
i += 1
}
column += 1
}
l
}
def shortestPath(
adjacency: Map[String, Seq[String]],
source: String,
target: String,
maxDepth: Int
): Option[List[String]] = {
require(maxDepth >= 0, "maxDepth must be nonnegative")
if (source == target) return Some(List(source))
val queue = mutable.Queue.empty[(String, Int)]
val seen = mutable.Set.empty[String]
val parent = mutable.Map.empty[String, String]
queue.enqueue((source, 0))
seen += source
while (queue.nonEmpty) {
val (u, depth) = queue.dequeue()
if (depth < maxDepth) {
val neighbors = adjacency.getOrElse(u, Seq.empty).distinct.sorted
neighbors.foreach { v =>
if (!seen(v)) {
seen += v
parent(v) = u
if (v == target)
return Some(reconstruct(parent.toMap, source, target))
queue.enqueue((v, depth + 1))
}
}
}
}
None
}
private def reconstruct(
parent: Map[String, String],
source: String,
target: String
): List[String] = {
var path = List(target)
var current = target
while (current != source) {
current = parent(current)
path = current :: path
}
path
}
private def multiplyLLT(l: Array[Array[Double]]): Array[Array[Double]] = {
val n = l.length
val out = Array.ofDim[Double](n, n)
var i = 0
while (i < n) {
var j = 0
while (j < n) {
var k = 0
while (k < n) {
out(i)(j) += l(i)(k) * l(j)(k)
k += 1
}
j += 1
}
i += 1
}
out
}
private def close(x: Double, y: Double): Boolean =
math.abs(x - y) <= 1e-9
def main(args: Array[String]): Unit = {
val a = Array(
Array(25.0, 15.0, -5.0),
Array(15.0, 18.0, 0.0),
Array(-5.0, 0.0, 11.0)
)
val l = cholesky(a)
val expected = Array(
Array(5.0, 0.0, 0.0),
Array(3.0, 3.0, 0.0),
Array(-1.0, 1.0, 3.0)
)
for (i <- a.indices; j <- a.indices)
assert(close(l(i)(j), expected(i)(j)))
val rebuilt = multiplyLLT(l)
for (i <- a.indices; j <- a.indices)
assert(close(rebuilt(i)(j), a(i)(j)))
val graph = Map(
"A" -> Seq("B", "C"),
"B" -> Seq("D"),
"C" -> Seq("D"),
"D" -> Seq("E")
)
assert(shortestPath(graph, "A", "E", 3).contains(List("A", "B", "D", "E")))
assert(shortestPath(graph, "A", "E", 2).isEmpty)
assert(shortestPath(Map("A" -> Seq("A", "B")), "A", "B", 1)
.contains(List("A", "B")))
assert(shortestPath(graph, "A", "A", 0).contains(List("A")))
println("MockLinearGraph tests passed")
}
}
| 算法 | 输入 | 预期 |
|---|---|---|
| Cholesky | 3×3 SPD 示例矩阵 | [[5,0,0],[3,3,0],[-1,1,3]] |
| Cholesky | [[25,-50],[-50,101]] | [[5,0],[-10,1]] |
| Cholesky | [[1,2],[2,1]] | 第二 pivot 失败。 |
| Cholesky | [[1,2],[0,1]] | 不对称拒绝。 |
| BFS | A→B/C→D→E,depth 3 | A-B-D-E(排序固定 tie)。 |
| BFS | 同图,depth 2 | None。 |
| BFS | A→A、A→B | 终止并返回 A-B。 |
parent(v) 指向上一层,回溯路径无环。// 关系语义;列名按实际 schema 替换
val candidates = frontier
.join(edges, frontier("artist") === edges("src"))
.select(frontier("source"), edges("dst").as("artist"))
.dropDuplicates("source", "artist")
val next = candidates.join(
visited,
Seq("source", "artist"),
"left_anti"
)
val visitedNext = visited.union(next).dropDuplicates("source", "artist")
以上是 DataFrame 关系骨架,不是本地 BFS 的必要依赖。核心顺序是 join → 轮内去重 → left-anti visited → union。
maxDepth=0 只允许 source 自己。输入日志每行包含自然语言 token。实现大小写无关词频,去除非字母数字分隔产生的空 token,只保留次数至少 2 的词:
Seq[String] => Map[String,Int];RDD[String] => RDD[(String,Int)];reduceByKey 与 groupByKey 的区别;persist 位置与小型测试。def localCounts(lines: Seq[String]): Map[String, Int]
def counts(logs: RDD[String]): RDD[(String, Int)]
参考实现固定:
2.12.18;3.5.1;"org.apache.spark" %% "spark-core" % "3.5.1" % "provided";provided,或需由 spark-submit 提供 Spark runtime。// build.sbt
ThisBuild / scalaVersion := "2.12.18"
libraryDependencies +=
"org.apache.spark" %% "spark-core" % "3.5.1" % "provided"
| 部分 | 自测分值 | 检查项 |
|---|---|---|
| tokenize/normalize | 8 | 小写、正则、过滤空 token、规则一致。 |
| RDD 键值处理流程 | 9 | flatMap、(word, 1)、reduceByKey。 |
| filter/action | 10 | 阈值过滤;不在算法函数中 collect 全量。 |
| lazy/persist | 7 | action 触发;复用 base RDD 才 persist。 |
| Collection 版本 | 6 | foldLeft/mutable 计数正确。 |
| 复杂度/shuffle | 5 | 线性 token 化;map-side combine;避免 materialize values。 |
import org.apache.spark.{SparkConf, SparkContext}
import org.apache.spark.rdd.RDD
object MockWordCount {
private val SplitPattern = "[^a-z0-9]+"
private def tokens(line: String): Iterator[String] =
line.toLowerCase
.split(SplitPattern)
.iterator
.filter(_.nonEmpty)
def localCounts(lines: Seq[String]): Map[String, Int] =
lines.iterator
.flatMap(tokens)
.foldLeft(Map.empty[String, Int]) { (acc, word) =>
acc.updated(word, acc.getOrElse(word, 0) + 1)
}
.filter { case (_, count) => count >= 2 }
def counts(logs: RDD[String]): RDD[(String, Int)] =
logs
.flatMap(tokens)
.map(word => (word, 1))
.reduceByKey(_ + _)
.filter { case (_, count) => count >= 2 }
def main(args: Array[String]): Unit = {
val input = Seq("To be, or not to be", "to")
assert(localCounts(input) == Map("to" -> 3, "be" -> 2))
assert(localCounts(Seq.empty).isEmpty)
val conf = new SparkConf()
.setAppName("MockWordCount")
.setMaster("local[2]")
val sc = new SparkContext(conf)
try {
val logs = sc.parallelize(input, 2).persist()
val actual = counts(logs).collect().toMap
assert(actual == Map("to" -> 3, "be" -> 2))
logs.unpersist(blocking = false)
println("MockWordCount tests passed")
} finally {
sc.stop()
}
}
}
| 输入 | 阈值后输出 | 覆盖 |
|---|---|---|
"To be, or not to be", "to" | to→3, be→2 | 大小写、标点、阈值。 |
| 空 Seq/RDD | 空 | 空输入。 |
"a---a", "A" | a→3 | 连续分隔、小写。 |
reduceByKey 可先做 map-side combine,再按 key shuffle。(word,1);同 key 的和值等于全局出现次数。1,而求和只需可结合的局部部分和。collect 只用于小测试;生产输出应 saveAsTextFile 或交给下游。这两项适合训练半实现、推导、复杂度和 Spark 数据流。本页的安排不表示官方考核要求。
接口:
def batchGD(
data: RDD[(Array[Double], Double)],
initial: Array[Double],
learningRate: Double,
iterations: Int
): Array[Double]
平方误差梯度:g=(1/n)Σ(xTw−y)x。每一轮所有样本必须读取同一个 wt;driver 只更新一次得到 wt+1。
import org.apache.spark.rdd.RDD
import org.apache.spark.storage.StorageLevel
def batchGD(
data: RDD[(Array[Double], Double)],
initial: Array[Double],
learningRate: Double,
iterations: Int
): Array[Double] = {
require(learningRate > 0.0, "learningRate must be positive")
require(iterations >= 0, "iterations must be nonnegative")
val dimension = initial.length
val cached = data.persist(StorageLevel.MEMORY_AND_DISK)
val count = cached.count()
require(count > 0L, "data must be nonempty")
require(cached.filter(_._1.length != dimension).take(1).isEmpty,
"feature dimension mismatch")
var weights = initial.clone()
var iteration = 0
while (iteration < iterations) {
val broadcast = data.context.broadcast(weights)
val gradient = cached.treeAggregate(Array.fill(dimension)(0.0))(
seqOp = (sum, row) => {
val (x, y) = row
val w = broadcast.value
var prediction = 0.0
var j = 0
while (j < dimension) {
prediction += x(j) * w(j)
j += 1
}
val error = prediction - y
j = 0
while (j < dimension) {
sum(j) += error * x(j)
j += 1
}
sum
},
combOp = (left, right) => {
var j = 0
while (j < dimension) {
left(j) += right(j)
j += 1
}
left
}
)
broadcast.destroy(blocking = false)
var j = 0
while (j < dimension) {
weights(j) -= learningRate * gradient(j) / count.toDouble
j += 1
}
iteration += 1
}
cached.unpersist(blocking = false)
weights
}
复杂度:每轮 Θ(nd),总 Θ(Tnd);每轮 broadcast d 个参数并 reduce d 维梯度。学习率过大不保证 loss 下降。
测试:一维点 (1→1),(2→2)、w0=0、小正学习率时,第一轮 gradient 为负,w 应向 1 增大;全零 feature 的 gradient 为 0。
对 rows xi∈ℝd,Gram 为:
G=XTX=ΣixixiT
def gram(rows: Iterable[Array[Double]]): Array[Array[Double]] = {
val iterator = rows.iterator
if (!iterator.hasNext) return Array.empty[Array[Double]]
val first = iterator.next()
val dimension = first.length
val out = Array.ofDim[Double](dimension, dimension)
def add(row: Array[Double]): Unit = {
require(row.length == dimension, "row length mismatch")
var j = 0
while (j < dimension) {
var k = j
while (k < dimension) {
out(j)(k) += row(j) * row(k)
k += 1
}
j += 1
}
}
add(first)
iterator.foreach(add)
var j = 0
while (j < dimension) {
var k = j + 1
while (k < dimension) {
out(k)(j) = out(j)(k)
k += 1
}
j += 1
}
out
}
测试:rows [1,2]、[3,4] 应得 [[10,14],[14,20]]。
不变量:处理任意前缀后,G=ΣxixiT,因此 G 对称且半正定。只计算上三角时必须恰好镜像一次。
复杂度:dense 时间 Θ(nd²),状态 Θ(d²)。当 d² 无法放入单机内存时,这条基线不成立。Spark 版可发出上三角 key ((j,k), x(j)*x(k)) 并 reduceByKey,但会产生 Θ(nd²) pairs;稀疏数据应只枚举非零 pair。
rho 上界、逐项 max、Double 和空摘要。parent,不要复制整条路径。| 代码 | 接口 | 括号/结构 | 关键类型 | 依赖边界 |
|---|---|---|---|---|
| HLL | 完整 | 单文件 object/class | Long unsigned shift、Double | 无第三方依赖 |
| Cholesky/BFS | 完整 | 单文件 object | Array[Double]、Option[List] | Scala 标准库 |
| WordCount | 完整 | 单文件 object | RDD[(String,Int)] | Spark 3.5.1 / Scala 2.12.18 |
| Batch GD | 核心函数完整 | 可放入 Spark object | RDD、broadcast、Array[Double] | 同上 |
| Gram | 完整 | 标准库函数 | Iterable、二维 Array | 无第三方依赖 |
seen 更新时机、深度限制、parent 回溯与 Spark 前沿集合的关系。Seq 与 RDD 间转换,并说明惰性求值、persist 与混洗。