数据湖擅长以较低成本保存海量历史数据,却不天然适合在线服务中的点查询:当请求只需要某个用户、歌曲或订单的一条记录时,扫描大量对象存储文件的延迟和成本都难以接受。Spotify 这类需要同时处理庞大数据规模与低延迟在线请求的公司,面临的核心问题就是:如何在保留数据湖存储优势的同时,快速定位一条具体记录。
一种实用思路是在数据湖与在线服务之间增加索引层。索引不复制全部业务数据,而是记录“键位于哪个文件、哪个分区或哪个行组”,把原本昂贵的全湖扫描缩小为一次索引查询和一次定向读取。
点查询为什么会拖慢数据湖
数据湖通常由对象存储上的 Parquet、Avro 或 ORC 文件组成。这些格式适合列裁剪、压缩和批量分析,但在线点查询具有不同的访问模式:
- 查询条件通常是单个精确键,例如
user_id = 38291。 - 每次只返回一条或少量记录。
- 请求量可能很高,并且尾延迟比平均吞吐量更重要。
- 新数据持续写入,索引必须跟上文件提交和删除的节奏。
即使 Parquet 文件内部带有统计信息,服务仍要先知道应该打开哪些文件。如果对象存储里存在数百万个文件,逐个读取元数据同样不可行。按日期、地区等字段分区只能减少一部分范围;当查询键与分区键不同,文件定位问题仍然存在。
因此,面向点查询的索引通常维护类似下面的映射:
record_key -> object_path + row_group + optional_offset/version
索引查找负责定位,Parquet 或其他数据文件仍是事实数据的载体。两者职责分离后,索引可以针对低延迟随机读取优化,数据湖则继续针对低成本批量存储优化。
索引层需要解决的不只是查找
真正困难的部分往往不是建立一个键值映射,而是让映射长期保持正确。
原子可见性。 数据文件应先成功写入,再发布对应索引。否则索引可能指向尚不存在或没有完整写入的对象。若数据湖使用快照或表格式,索引最好基于已提交快照更新。
更新与删除。 同一个键可能出现在多个文件中。索引需要记录版本、事件时间或快照序号,以确定哪条记录有效。删除操作也可能表现为 tombstone,而不是立即重写历史文件。
压实后的重建。 Compaction 会把小文件合并成大文件,并改变记录的物理位置。旧索引必须在新文件发布后切换,随后才能清理旧文件。
热点和尾延迟。 流行歌曲、热门用户或公共配置可能形成热点。索引分片不能只考虑总容量,还要考虑键的访问分布;必要时可在服务侧增加小型缓存。
降级路径。 索引不可用时,系统需要明确选择:返回错误、读取稍旧的副本,还是退回较慢的批查询。直接扫描整个数据湖通常不适合作为在线请求的自动降级方案。
可以这样实践:用 SQLite 给 Parquet 文件建立轻量索引
下面是一个可运行的最小示例。它生成两个 Parquet 文件,将 user_id 到文件路径和行号的映射写入 SQLite,然后执行一次定向点查询。这个示例用于说明架构,不代表来源文章采用了 SQLite;生产环境可以把 SQLite 替换为分布式键值存储、托管数据库或专用索引服务。
先安装依赖:
python -m pip install pyarrow
将下面内容保存为 lake_index_demo.py 并运行:
from pathlib import Path
import sqlite3
import pyarrow as pa
import pyarrow.parquet as pq
LAKE = Path("demo_lake")
INDEX = Path("lake_index.db")
def write_data() -> list[Path]:
LAKE.mkdir(exist_ok=True)
batches = [
[
{"user_id": "u-100", "plan": "free", "country": "SE"},
{"user_id": "u-101", "plan": "premium", "country": "DE"},
],
[
{"user_id": "u-200", "plan": "premium", "country": "US"},
{"user_id": "u-201", "plan": "free", "country": "BR"},
],
]
paths = []
for number, rows in enumerate(batches):
path = LAKE / f"users-{number}.parquet"
pq.write_table(pa.Table.from_pylist(rows), path)
paths.append(path)
return paths
def rebuild_index(paths: list[Path]) -> None:
with sqlite3.connect(INDEX) as db:
db.execute("DROP TABLE IF EXISTS record_index")
db.execute(
"""
CREATE TABLE record_index (
record_key TEXT PRIMARY KEY,
object_path TEXT NOT NULL,
row_number INTEGER NOT NULL
)
"""
)
for path in paths:
keys = pq.read_table(path, columns=["user_id"])["user_id"].to_pylist()
db.executemany(
"INSERT INTO record_index VALUES (?, ?, ?)",
[(key, str(path), row_number) for row_number, key in enumerate(keys)],
)
db.commit()
def get_user(user_id: str) -> dict | None:
with sqlite3.connect(INDEX) as db:
location = db.execute(
"SELECT object_path, row_number FROM record_index WHERE record_key = ?",
(user_id,),
).fetchone()
if location is None:
return None
object_path, row_number = location
table = pq.read_table(object_path)
return table.slice(row_number, 1).to_pylist()[0]
if __name__ == "__main__":
paths = write_data()
rebuild_index(paths)
print(get_user("u-200"))
运行命令:
python lake_index_demo.py
预期输出:
{'user_id': 'u-200', 'plan': 'premium', 'country': 'US'}
这个版本为了便于理解,读取了整个目标 Parquet 文件。在生产实现中,应将索引粒度细化到 row group,并利用过滤条件和列裁剪:
table = pq.read_table(
object_path,
columns=["user_id", "plan", "country"],
filters=[("user_id", "=", user_id)],
)
还可以为索引记录增加 snapshot_id、row_group、updated_at 和校验值,从而支持快照切换、压实迁移与一致性检查。
从原型走向在线服务
生产系统不应只比较“有索引”和“无索引”的平均延迟。更有价值的验收清单包括:
- 统计索引命中率以及不存在键的查询成本。
- 分别测量索引查找、对象存储读取、解压和反序列化延迟。
- 观察
p95、p99和超时率,而不只是平均值。 - 验证文件提交、索引发布和旧文件清理的顺序。
- 对索引滞后、重复键、删除记录和 compaction 建立自动化测试。
- 为索引重建设计离线流程,避免只能依赖增量事件恢复。
- 根据可接受的新鲜度决定使用同步更新、异步更新还是周期性批处理。
索引数据湖的价值,在于把“数据存在哪里”从每次在线查询要解决的问题,变成写入流程提前维护的元数据。代价则是新增的一致性协议、索引存储和运维责任。只有当点查询规模、延迟目标和数据新鲜度被明确量化后,团队才能选择合适的索引粒度与更新机制,而不是把分析型数据湖直接当作在线数据库使用。