Netflix 如何把 Cassandra 宽分区读延迟从秒级压到毫秒级

2026-07-06 31 预计阅读时间: 1 分钟
来源: infoq.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.

预计阅读时间:10 分钟

时间序列数据写进 Cassandra 很顺手,但一旦某个租户、设备或内容 ID 的数据量远超预期,原本看起来干净的分区键就会变成“宽分区”。Netflix 工程团队针对这类负载引入了动态分区拆分:用元数据识别过大的逻辑分区,把它拆成多个更小的子分区,并在读取时自动路由到这些子分区。结果是读延迟从秒级降到毫秒级,超时减少,集群稳定性也更好,同时对业务侧保持透明。

问题不在 Cassandra,而在分区形状

Cassandra 的性能很依赖分区大小。时间序列场景常见的表设计会把某个实体和时间桶拼成分区键,例如:

CREATE TABLE events_by_asset_day (
  asset_id text,
  day date,
  event_ts timestamp,
  payload text,
  PRIMARY KEY ((asset_id, day), event_ts)
) WITH CLUSTERING ORDER BY (event_ts DESC);

这个设计在多数 asset_id + day 数据量均匀时很有效。问题出现在热点实体上:某个 asset_id 一天内写入量特别大,单个分区会积累大量行。读取最近 N 条、读取某个时间范围,甚至后台 compaction、repair、缓存命中率都会被拖累。

Netflix 的思路不是要求业务方重新理解所有物理分区细节,而是在逻辑分区和物理分区之间加一层元数据。业务仍然按原来的逻辑键查询,系统根据分区大小和拆分记录,把请求分发到更小的子分区。

动态拆分的关键:元数据驱动,而不是硬编码规则

动态分区拆分有三个核心动作:

  1. 检测 oversized partition:持续观察某个逻辑分区是否过宽,例如行数、字节数、读延迟或超时率超过阈值。
  2. 写入拆分元数据:为一个逻辑分区记录多个 child partition,例如按时间范围、序号或哈希段拆开。
  3. 读路径自动 fan out:查询时先查元数据,再向相关子分区发起读取,最后合并排序并返回。

这类方案的价值在于透明性。调用方不需要知道“这个 asset 今天被拆成了 8 个子分区”,也不需要临时改查询参数。系统内部承担了多一次元数据查找、并发读取和结果归并的复杂度,换来单次底层读取更稳定、更短尾。

可以这样实践:用拆分元数据改造时间序列表

下面是一个可改造的最小方案。它不是 Netflix 内部实现细节,而是把“元数据驱动的动态分区拆分”落到 Cassandra 表结构和读路径上的一种工程化实践。

原始数据表可以把物理分片编号放进 partition key:

CREATE TABLE events_by_asset_day_shard (
  asset_id text,
  day date,
  shard int,
  event_ts timestamp,
  payload text,
  PRIMARY KEY ((asset_id, day, shard), event_ts)
) WITH CLUSTERING ORDER BY (event_ts DESC);

CREATE TABLE partition_split_metadata (
  asset_id text,
  day date,
  shard int,
  start_ts timestamp,
  end_ts timestamp,
  state text,
  PRIMARY KEY ((asset_id, day), shard)
);

假设某个逻辑分区 asset-42 / 2026-01-20 过大,可以写入 4 个子分区元数据:

INSERT INTO partition_split_metadata (asset_id, day, shard, start_ts, end_ts, state)
VALUES ('asset-42', '2026-01-20', 0, '2026-01-20T00:00:00Z', '2026-01-20T06:00:00Z', 'active');

INSERT INTO partition_split_metadata (asset_id, day, shard, start_ts, end_ts, state)
VALUES ('asset-42', '2026-01-20', 1, '2026-01-20T06:00:00Z', '2026-01-20T12:00:00Z', 'active');

INSERT INTO partition_split_metadata (asset_id, day, shard, start_ts, end_ts, state)
VALUES ('asset-42', '2026-01-20', 2, '2026-01-20T12:00:00Z', '2026-01-20T18:00:00Z', 'active');

INSERT INTO partition_split_metadata (asset_id, day, shard, start_ts, end_ts, state)
VALUES ('asset-42', '2026-01-20', 3, '2026-01-20T18:00:00Z', '2026-01-21T00:00:00Z', 'active');

读路径可以先查元数据,再并发查询命中的子分区。下面示例使用 Python Cassandra driver,运行前需要修改 CONTACT_POINTSKEYSPACE,并安装依赖:

python -m pip install cassandra-driver
from concurrent.futures import ThreadPoolExecutor, as_completed
from datetime import datetime, timezone
from cassandra.cluster import Cluster

CONTACT_POINTS = ["127.0.0.1"]
KEYSPACE = "demo"

cluster = Cluster(CONTACT_POINTS)
session = cluster.connect(KEYSPACE)

meta_stmt = session.prepare("""
SELECT shard, start_ts, end_ts
FROM partition_split_metadata
WHERE asset_id = ? AND day = ?
""")

read_stmt = session.prepare("""
SELECT event_ts, payload
FROM events_by_asset_day_shard
WHERE asset_id = ? AND day = ? AND shard = ?
  AND event_ts >= ? AND event_ts < ?
LIMIT ?
""")


def intersect(a_start, a_end, b_start, b_end):
    return max(a_start, b_start), min(a_end, b_end)


def read_events(asset_id, day, start_ts, end_ts, limit=100):
    shards = list(session.execute(meta_stmt, (asset_id, day)))

    # 未拆分时可以约定 shard=0,避免业务侧感知物理分片。
    if not shards:
        shards = [{"shard": 0, "start_ts": start_ts, "end_ts": end_ts}]

    tasks = []
    with ThreadPoolExecutor(max_workers=min(len(shards), 8)) as pool:
        for shard in shards:
            s, e = intersect(start_ts, end_ts, shard.start_ts, shard.end_ts)
            if s >= e:
                continue
            tasks.append(pool.submit(
                session.execute,
                read_stmt,
                (asset_id, day, shard.shard, s, e, limit)
            ))

        rows = []
        for task in as_completed(tasks):
            rows.extend(task.result())

    rows.sort(key=lambda r: r.event_ts, reverse=True)
    return rows[:limit]


if __name__ == "__main__":
    start = datetime(2026, 1, 20, 10, 0, tzinfo=timezone.utc)
    end = datetime(2026, 1, 20, 14, 0, tzinfo=timezone.utc)
    for row in read_events("asset-42", "2026-01-20", start, end, limit=20):
        print(row.event_ts, row.payload)

这个例子故意把复杂度放在读路径中:业务传入的仍然是 asset_id + day + time range,读服务负责决定要打到哪些 child partitions。

写路径和迁移不能含糊

分区拆分真正难的地方通常不在查询代码,而在状态转换。

拆分发生时,你需要明确这些问题:

  • 新写入应该落到哪个 child partition?按时间范围拆分时比较直观;按哈希拆分时要保证路由函数稳定。
  • 旧数据是否回填?如果不回填,读路径要同时读旧 shard 和新 shard;如果回填,要处理双写、幂等和校验。
  • 元数据更新是否原子?读路径看到半完成状态时,不能漏数据。
  • fan out 上限是多少?拆得太细会把一个慢查询变成很多小查询,压垮协调节点或线程池。
  • 监控是否按逻辑分区和物理子分区同时看?只看整体 P99 会掩盖单个热点。

Netflix 报告的收益来自把宽分区拆小之后,底层读请求变得更可预测。这个方向适合读超时、尾延迟和热点分区已经成为稳定性问题的系统;如果你的分区本来就小,强行加一层拆分元数据只会增加路径复杂度。

落地前的检查清单

引入动态分区拆分前,可以用这张清单压一遍设计:

  • 找出最坏的 1% 分区,而不是只看平均分区大小。
  • 给逻辑分区定义明确阈值:行数、字节数、读 P99、超时率至少选两个指标交叉验证。
  • 让拆分元数据可审计、可回滚、可灰度。
  • 读路径必须限制并发和总返回行数,避免 fan out 失控。
  • 写路径要有稳定的 shard 选择规则,并覆盖拆分过程中的新旧状态。
  • 压测要包含热点实体,不要只用均匀分布数据。

动态分区拆分的本质,是承认时间序列负载不会永远均匀,然后把“异常大的逻辑分区”变成系统可以管理的多个小物理分区。它不是免费优化,但当 Cassandra 宽分区已经把读延迟拖到秒级时,这类元数据驱动的拆分往往比反复扩容更接近问题根部。


相关推荐