给数据湖加一层索引:让在线点查询不再扫描海量文件

2026-07-28 15 预计阅读时间: 1 分钟
来源: engineering.atspotify.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:9 分钟

数据湖擅长以较低成本保存海量历史数据,却不天然适合在线服务中的点查询:当请求只需要某个用户、歌曲或订单的一条记录时,扫描大量对象存储文件的延迟和成本都难以接受。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_idrow_groupupdated_at 和校验值,从而支持快照切换、压实迁移与一致性检查。

从原型走向在线服务

生产系统不应只比较“有索引”和“无索引”的平均延迟。更有价值的验收清单包括:

  • 统计索引命中率以及不存在键的查询成本。
  • 分别测量索引查找、对象存储读取、解压和反序列化延迟。
  • 观察 p95p99 和超时率,而不只是平均值。
  • 验证文件提交、索引发布和旧文件清理的顺序。
  • 对索引滞后、重复键、删除记录和 compaction 建立自动化测试。
  • 为索引重建设计离线流程,避免只能依赖增量事件恢复。
  • 根据可接受的新鲜度决定使用同步更新、异步更新还是周期性批处理。

索引数据湖的价值,在于把“数据存在哪里”从每次在线查询要解决的问题,变成写入流程提前维护的元数据。代价则是新增的一致性协议、索引存储和运维责任。只有当点查询规模、延迟目标和数据新鲜度被明确量化后,团队才能选择合适的索引粒度与更新机制,而不是把分析型数据湖直接当作在线数据库使用。


相关推荐