Skip to content

Storyline 三表 Lance 存储

StorylineLanceStore 是 pChronicle 的 Storyline-native 规范化物理表示。它与 events.lance 原始事件日志并列存在,不替代后者。

这是 pChronicle 唯一的规范化三表模型。旧的 ATIF sessions / steps / tool_callsNormalizedStore 和内存联表视图已经删除。ATIF 仍作为输入输出格式存在,但查询时先转换为 Storyline,再投影到本页定义的 runs / steps / tool_calls schema,不再维护第二套表结构。

表模型

粒度 逻辑主键 外键
runs.lance 每个 Storyline 一行 session_id
steps.lance 每个 turn 一行 (session_id, step_id) session_id → runs
tool_calls.lance 每个 tool call 一行 (session_id, tool_call_id) (session_id, step_id) → steps

run_id 是 Run 分组键;一个 Run 可以包含主 Story 和多个 subagent Story,因此 runs.lance 中可能有多行共享同一个 run_idsession_id 才是 Storyline 文档的 唯一键。

常用 JSON 值(message、arguments、metrics、extra 等)以 UTF-8 JSON 列保存;身份、 顺序、类型、时间和性能字段使用独立的 Arrow 标量列,便于过滤和分析。

steps.timestamp 是规范化到 UTC 的 Timestamp(Millisecond, "UTC")。写入端拒绝无法 解析为 RFC3339 的非空时间;原始时区偏移和亚毫秒精度不会写入物理表,读取为 Storyline / ATIF 时统一编码成带 Z 后缀、精度不超过毫秒的 UTC 字符串(整秒省略小数部分)。SQL 排序、范围过滤和时间聚合直接使用 timestamp。这是一次物理 schema 变更;旧三表投影 需要从 canonical events 或原始交换文件重建,不对既有 Lance generation 做原地类型猜测。

tool result 不再留在 step 的 observation JSON 中。写入时根据 observation.results[].source_call_id 关联到对应 tool call,并保存到该行的 results_json。缺失或错误的关联会拒绝整次写入。

大块内容层

设计目标与边界

Agent 轨迹中的长 reasoning、工具输出、源代码、日志和多模态载荷会让列式表出现少量超大 cell。若把它们与身份、顺序、类型和指标一起内联,常规过滤与聚合也要承受更大的 fragment、 page cache 和解码开销。pChronicle 因此在三表之外增加共享内容层,但保持三个约束:

  1. runs / steps / tool_calls 的 Arrow schema 和 SQL 结果不变;内容层是内部物理优化。
  2. 小值继续内联,只有达到阈值的 UTF-8/JSON cell 才外置,避免所有读取都退化成 KV lookup。
  3. 内容按原始字节寻址并跨 Storyline 复用;不把轨迹主键、生命周期或业务去重混入内容层。

当前实现没有定制 Lance 文件格式或私有索引类型,而是组合 Lance Blob v2、普通 BTree scalar index 和 DataFusion execution node。这样可以得到需要的延迟物化能力,同时把维护 面限制在 pChronicle 自己的协议与执行计划中。

内部描述符协议

超过默认 64 KiB 的内容列在三表中暂时编码为:

<RS>PCHRONICLE-CONTENT-1:<type>:<codec>:<blake3-256>:<raw_length>:<preview-base64url>
字段 当前编码 作用
magic/version PCHRONICLE-CONTENT-1 严格识别内部引用并允许协议演进
logical type u / j / b UTF-8、JSON;binary 标签已保留给后续二进制列
codec i / z identity 或 Zstd
content id 64 位十六进制 BLAKE3-256 对未压缩原始字节寻址、校验和跨轨迹复用
raw length u64 解压后长度校验,也允许无 payload 的代价判断
preview URL-safe Base64 默认最多 256 个 UTF-8 字节的安全前缀

描述符只允许存在于内部物理列。用户原文若恰好以 magic 开头,也会被强制外置,读取时再 恢复为原文,从而消除“用户字符串被误认为引用”的歧义。公开的读取、SQL、转换和导出 API 必须返回完整值或显式 preview,不能泄露描述符。

objects.lance 使用以下物理列:

作用
content_id BLAKE3 内容地址;建立 BTree index
logical_type, media_type 逻辑类型和 MIME 提示
raw_length, stored_length, codec 完整性检查和存储代价
preview 无 Blob I/O 的安全预览
payload Lance Blob v2,保存 identity/Zstd 字节
created_at_ms 对象创建时间

写入、复用与发布

写入端对候选 cell 依次执行:

原始 UTF-8/JSON
  ├─ 小于阈值 ───────────────────────────────► 原值内联
  └─ 达到阈值 / 命中 magic
       ├─ BLAKE3(raw bytes) + UTF-8 preview
       ├─ Zstd;没有净收益则保留 identity
       ├─ batch 内按 content_id 合并并检查碰撞
       ├─ BTree 批量查询 objects.lance,跳过已存在对象
       └─ 先提交对象 version,再写三表 descriptor,最后发布 CURRENT

对象必须先于引用持久化;CURRENT 同时固定三张业务表和对象表的精确 Lance version。 任一步骤失败都不会发布新快照:允许留下不可达对象,但不会发布悬空引用或跨表半提交。 跨轨迹复用只依赖内容地址,不依赖 session 生命周期,因此同一长文本在不同 Run 中只保存 一次。同一写入批次内若 content id 相同但 codec、原始长度或存储字节不一致,会拒绝写入。

对象层在普通写入期间保持 append-only,GC 不进入写入热路径。显式 maintain 会只扫描三表 的内容引用列,计算当前快照的可达 content id,并清理不可达 payload。生产环境仍需要把对象 增长率、不可达字节和维护耗时纳入指标。

查询期延迟物化

StorylineDataSource 先让 Lance 完成业务表的 projection、可安全谓词、scalar index、limit 和并行扫描,再在计划中插入 ContentHydrationExec

  • 查询不引用内容列时,不打开 objects.lance payload;
  • 只收集投影中实际出现的 content id,以最多 512 个为一组走 BTree lookup;
  • 根据 row address 批量读取 Blob,解压后验证长度与 BLAKE3,再恢复原 Utf8 列;
  • 内容列谓词不能作用于描述符,必须保留在 hydration 之后由 DataFusion 计算;
  • Preview 模式只返回描述符中的 UTF-8 前缀,零 payload I/O,并拒绝内容列谓词,避免把 preview 错当成完整值。

因此大内容的成本只由真正读取这些内容的查询承担;身份过滤、计数、分组和指标分析仍沿用 紧凑的三表列式路径。

提交布局

root/
├── CURRENT
├── objects.lance/
└── generations/
    └── <table-generation>/
        ├── runs.lance/
        ├── steps.lance/
        └── tool_calls.lance/

首次导入创建三张规范化 Lance dataset、共享的 objects.lance 和标量索引。后续 replace_storyline 不再读取或重写全库,而是按 session_id 删除旧行、追加新行。每次替换 会产生一个新的逻辑 snapshot; CURRENT 是一段 JSON,记录逻辑 snapshot id、物理 table_generation、三张表以及对象表 各自精确的 Lance version id。对象先持久化,三张业务表随后写入,最后才更新 CURRENT; 因此失败最多留下不可达对象,不会发布悬空引用或跨表半提交。

阈值、preview 长度和 Zstd level 可通过 StorylineContentOptions 配置;三表 schema 不变。

Lance MVCC 的旧版本默认保留,便于已打开的 reader 固定快照及故障恢复。频繁增量更新 会积累 fragment、delete file 和未合并的索引增量。普通 replace 不执行 index refresh 或 compaction,避免某次写请求出现维护型长尾;生产环境通过 maintain 显式执行三表并行 compaction、补齐/刷新索引、内容 GC 和按保留期 vacuum。维护产生的三个新 version 仍先原子更新 CURRENT,之后才回收旧版本。CURRENT 必须是包含全部精确版本的 JSON 指针,不读取旧的 纯文本 generation 指针。

本地写入通过进程内锁和文件锁串行化;对象存储通过 CURRENT 的 ETag/version 条件更新 执行 optimistic CAS。stale commit 不能移动 CURRENTStorylineLanceStore 在 CAS 冲突后 直接返回错误,不会重新读取、merge 或自动重试。调用方若选择重试,必须从最新 snapshot 重新 开始完整 replace。上层 lease 可减少冲突,但不改变这一失败语义。

Rust API

let store = StorylineLanceStore::open(path).await?;
store.replace_storyline(&storyline).await?;
let runs = store.list_run_summaries(Default::default()).await?;
let steps = store.list_steps_page("session-id", Default::default()).await?;
let restored = store.get_storyline_full("session-id").await?;
let report = store.maintain(&LanceMaintenanceOptions::default()).await?;

replace_storylinesession_id 为边界替换三张表中的相关行,同时保留同一 store 内的其他 Storyline。

读取接口按成本分层:list_run_summaries 只扫描 runs 的标量摘要列,不打开 objects、steps 或 tool_calls;list_steps_pagelist_tool_calls_page 先读取排序键确定一页,再只为该页 读取全列并恢复外置内容。分页 cursor 固定 CURRENT generation;期间发生写入会返回 stale cursor 错误,调用方应从第一页重试。get_storyline_full 明确表示会读取三表并恢复该 Storyline 的全部内容。公共读取面只提供显式 full 或分页 API,不保留无成本提示的全量读取别名。

首次导入和替换都并行写三张表。Arrow 行按最多 8192 行一批懒编码并流入 Lance,避免 导入大型语料时同时保留整表的 Arrow 副本。CURRENT 只解析一次;DataSource 随后把每张 表直接打开到指针指定的 version,不再先验证、再重复打开同一 dataset。

生产环境通过 StorylineLanceStore::maintain Rust API 执行维护;公共 CLI 不提供维护命令。

DataFusion datasource

StorylineDataSource 在打开时固定 CURRENT 中的三个业务表 version 和对象表 version,并把三张 dataset 注册为 runsstepstool_calls。即使写入端随后切换 CURRENT,已经打开 的查询仍使用同一份三表快照。

let source = StorylineDataSource::open(path).await?;
let ctx = source.session_context()?;
let rows = ctx
    .sql("SELECT step_id, source FROM steps WHERE session_id = 's-1' ORDER BY step_id")
    .await?
    .collect()
    .await?;

Datasource 使用 Lance 原生 DataFusion execution plan,支持列裁剪、谓词和 limit 下推, 并采用 unordered physical scan 允许并行读取;有顺序要求的查询必须显式使用 ORDER BY step_id, call_index。未引用大内容列的查询不会打开 Blob;引用内容列时在 Lance 投影/安全谓词/limit 之后批量恢复。针对内容列的谓词不下推到内部引用,而是在恢复 后由 DataFusion 求值,确保 SQL 语义不变。内部引用不会由 pChronicle 的读取、查询、导出 API 返回;直接绕过 pChronicle 扫描底层 Lance 文件属于诊断接口,不在该保证内。

预览 UI 可把 StorylineDataSourceOptions::content_read_mode 设为 StorylineContentReadMode::Preview。该模式直接从描述符返回 UTF-8 安全的短 preview,零 Blob payload I/O;为避免把 preview 当成完整值产生错误结果,内容列谓词在 preview 模式 下会被明确拒绝。

首次创建 table generation 时建立以下标量索引:

BTree Bitmap
runs session_id, run_id
steps session_id, timestamp effective_kind, source
tool_calls session_id, tool_call_id function_name

这些索引针对按 Story/Run 定位、tool-call 查找和类型过滤。step_id 在每个 Storyline 内 从小值重新计数,全局选择性低,因此不单独建立 BTree;组合条件先用 session_id 定位到 单个 Storyline,再过滤很短的 step 范围。DataFusion 的索引谓词会下推为 Lance ScalarIndexQuery

StorylineDataSourceOptions 可显式控制 use_scalar_indexesscan_in_order;默认配置 面向在线分析查询启用索引、关闭物理顺序。关闭索引主要用于 benchmark、诊断或极小表 的全扫描对照。

统一查询引擎

ChronicleQueryEngine 是对外的只读 SQL 门面。Lance 与 ATIF 后端注册完全相同的 runsstepstool_calls 表,因此查询语句不需要随物理格式改变:

let lance = ChronicleQueryEngine::open_lance("./storyline-store").await?;
let batches = lance.query(
    "SELECT session_id, step_id, source FROM steps WHERE step_id >= 10"
).await?;

let atif = ChronicleQueryEngine::open_atif("./trajectories.ndjson")?;
let jsonl = atif.query_jsonl(
    "SELECT source, COUNT(*) AS steps FROM steps GROUP BY source ORDER BY source"
).await?;

query 返回 Arrow RecordBatch,适合服务端继续处理;dataframe 返回 lazy DataFrame, 适合追加 DataFusion 变换或查看计划;query_jsonl 用于 CLI/API 边界。调用者也可通过 context() 取得 SessionContext 注册 UDF 或额外表。

AtifDataSource 接受单个 ATIF JSON 对象、JSON 数组、每行一个完整 trajectory 的 JSONL/NDJSON,以及包含这些 ATIF 文档的目录。文件路径默认注册为按文件 lazy 的 StreamingTable:manifest 在打开时冻结路径和文件身份,scan 才读取命中文件;目录按 稳定顺序发现,每个文件是独立 partition,并以固定大小 Arrow batch 提供背压。显式 from_json / from_trajectories 因为调用者已经持有完整输入,仍使用 MemTable

pChronicle + JSON 投影查询快路径

旧路径把每份 JSON 完整解析为格式对象,再执行 ATIF → Storyline → 三表行 → 全宽 Arrow 后交给 SQL。即使查询只需要 sourceCOUNT(*),也会构造 message、reasoning、metrics、 tool calls 等未使用的大字段。新路径把优化边界前移到 TableProvider::scan

SQL / DataFrame
  → DataFusion projection + filters
  → FileScanSpec
      ├─ _file_ = / IN / LIKE:manifest 文件裁剪
      ├─ session_id:trajectory 裁剪
      ├─ step_id / source:step 裁剪
      └─ projected column set
  → BufRead / serde streaming decoder
      └─ DeserializeSeed + Visitor + IgnoredAny
  → 只为命中行解码被引用字段
  → projected Arrow RecordBatch
  → DataFusion 保留 inexact filter 再次校验

当前 fast path 的适用范围是 ATIF 单对象 JSON、JSON 数组(包括 pretty JSON),以及每行一个对象的 JSONL/NDJSON, 目标表为 steps,并且物理计划存在严格列裁剪。它有意保持保守:

输入/查询 执行路径
ATIF object/pretty object + projected steps reader-backed seeded projected decoder
ATIF array/pretty array + projected steps fill_buf 结构扫描 + 有界 element buffer + seeded from_slice
ATIF JSONL/NDJSON + projected steps BufRead 逐记录、有界 record buffer
_file_session_idstep_idsource 的安全简单谓词 可提前裁剪,DataFusion 仍复核
SELECT * 完整规范化 fallback
runs / tool_calls 完整规范化 fallback
OpenAI-message / ACTF 完整规范化 fallback
无法证明安全的表达式、OR/函数/跨列条件 不预裁剪,由 DataFusion 求值

DeserializeSeed 把查询 projection 和安全谓词传入 Visitor;未引用字段交给 IgnoredAny 做语法扫描,不构造 Value/Storyline。JSONL/NDJSON 以 BufRead 逐记录读取; JSON array 的结构扫描器识别字符串和转义,在不构造 DOM 的情况下提取单个 trajectory, 再通过 slice decoder 执行投影解析。单条 JSONL 记录或 array element 由 max_record_bytes 限制;单对象直接从 reader 解码。三种路径都不先复制整文件。 Arrow encoder 也只创建投影列,COUNT(*) 使用合法的零列 batch。轻量路径 校验 JSON、必需字段、重复 session、命中文档内的重复 step 和当前表内约束;跨表引用 完整性仍由导入路径或完整 fallback 负责。这一边界使临时查询不承担导入语义,同时不降低 SQL 结果正确性。

查询指标额外报告 projected_filesstreamed_recordsstreaming_buffer_peak_bytes、 scanned/pruned documents、scanned/pruned/emitted rows 和 projected_arrow_bytes,用来区分 “源字节扫描”“输入缓冲”“JSON 字段物化”和“Arrow 输出”四个成本。仓库 benchmark 报告 median/P95、rows/s、独立进程峰值 RSS,以及计数 allocator 观测到的 allocation calls/bytes; 这些是指定 corpus、查询和机器的回归数据,不是跨环境 SLA。

该路径仍需顺序扫描命中文件的全部 JSON 字节,不是文件内索引。一次性或受控批次查询可 直接使用 JSON;超大、远端或反复查询的数据应先转换为 Lance,利用 snapshot、列裁剪、 并行 fragment scan 和 scalar index。

ATIF 导入同样默认走 AtifReader。空 store 使用一个 producer 单遍完成 校验、Storyline 规范化和三表拆分,再经三条有界 Arrow channel 并行创建三个 Lance dataset;已有 store 则以最多 256 个 Storyline 为一个增量替换批次。两种路径都在所有 输入和三表写入成功后才原子切换一次 CURRENT

CLI 使用相同引擎,输出稳定的 JSONL:

pchronicle query ./trajectories.ndjson \
  'SELECT source, COUNT(*) AS steps FROM dataset.steps GROUP BY source ORDER BY source'

# 含 CURRENT 的三表 store 根目录会被 auto 识别为 Lance
pchronicle query ./storyline-store \
  'SELECT step_id, source FROM dataset.steps WHERE session_id = '\''s-1'\'' ORDER BY step_id'

# OpenAI/ACTF 目录直接查询;_file_ 为查询期相对路径列,不写入 Lance
pchronicle query ./openai-data \
  "SELECT _file_, COUNT(*) FROM dataset.steps WHERE _file_ LIKE 'batch/%' GROUP BY _file_"

查询是只读的;SQL 可以使用 SELECT、CTE、JOIN、聚合和 DataFusion 内置函数,但不通过 这个门面执行 DDL/DML。Lance 引擎打开时固定 CURRENT 指向的三个版本,从而保证一次 查询会话内三张表来自同一快照。

仓库提供两组可执行 benchmark:

# scalar index 与 full scan A/B
cargo bench -p persisting-pchronicle --bench atif_storyline_lance

# 导入、冷查询、点查、增量替换,以及 warm SQL 对比
PCHRONICLE_BENCH_SCALE=128 PCHRONICLE_BENCH_ITERS=30 \
  cargo bench -p persisting-pchronicle --bench lance_vs_json

# 冷 datasource open + projected JSON SQL 的 allocation、rows/s、P95、进程 RSS 与输入缓冲上界
PCHRONICLE_BENCH_SCALE=128 PCHRONICLE_BENCH_ITERS=30 \
  cargo bench -p persisting-pchronicle --bench json_streaming

# 同一基准切换为 pretty array,单独观察 element scanner + slice decoder
PCHRONICLE_BENCH_JSON_SHAPE=array PCHRONICLE_BENCH_SCALE=128 \
  PCHRONICLE_BENCH_ITERS=30 cargo bench -p persisting-pchronicle --bench json_streaming

JSON 对照使用单个 NDJSON 文件,避免大量小文件打开开销。ATIF steps 直接查询会把 DataFusion projection 和可安全预裁剪的 session_idstep_idsource 谓词传给 projected decoder:未引用 JSON 字段只做语法扫描,不构造 Storyline/三表对象,Arrow batch 也只包含执行计划需要的列。object、array、pretty JSON 和 JSONL/NDJSON 共用流式 projection decoder;SELECT * 和其他格式仍走完整规范化 fallback。轻量路径执行 JSON、必需字段和表内约束校验,跨表引用完整性由导入或完整 fallback 校验。预解析内存 JSON 对照只计算查询逻辑,用来区分产品工作流与纯内存遍历成本。 benchmark 还单独输出 DataSource 冷打开并执行 SQL、get_storyline_full 点查和单 Storyline 替换的延迟,避免 warm SQL 吞吐掩盖在线读写路径的写放大。

性能结论不应写成“Lance 在所有规模和查询上必然更快”:显式构造的 MemTable 或预解析 内存 JSON 在小数据下仍可能更快。默认 ATIF streaming 解决的是内存上界,不提供物理 索引;Lance 的主要优势仍是更小的物理体积、近乎常数的 datasource 打开时间,以及列 裁剪、并行扫描和选择性索引收益。