Storyline 三表 Lance 存储¶
StorylineLanceStore 是 pChronicle 的 Storyline-native 规范化物理表示。它与
events.lance 原始事件日志并列存在,不替代后者。
这是 pChronicle 唯一的规范化三表模型。旧的 ATIF
sessions / steps / tool_calls、NormalizedStore 和内存联表视图已经删除。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_id。session_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 因此在三表之外增加共享内容层,但保持三个约束:
runs/steps/tool_calls的 Arrow schema 和 SQL 结果不变;内容层是内部物理优化。- 小值继续内联,只有达到阈值的 UTF-8/JSON cell 才外置,避免所有读取都退化成 KV lookup。
- 内容按原始字节寻址并跨 Storyline 复用;不把轨迹主键、生命周期或业务去重混入内容层。
当前实现没有定制 Lance 文件格式或私有索引类型,而是组合 Lance Blob v2、普通 BTree scalar index 和 DataFusion execution node。这样可以得到需要的延迟物化能力,同时把维护 面限制在 pChronicle 自己的协议与执行计划中。
内部描述符协议¶
超过默认 64 KiB 的内容列在三表中暂时编码为:
| 字段 | 当前编码 | 作用 |
|---|---|---|
| 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.lancepayload; - 只收集投影中实际出现的 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 不能移动 CURRENT;StorylineLanceStore 在 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_storyline 以 session_id 为边界替换三张表中的相关行,同时保留同一 store
内的其他 Storyline。
读取接口按成本分层:list_run_summaries 只扫描 runs 的标量摘要列,不打开 objects、steps
或 tool_calls;list_steps_page 和 list_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 注册为 runs、steps、tool_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_indexes 与 scan_in_order;默认配置
面向在线分析查询启用索引、关闭物理顺序。关闭索引主要用于 benchmark、诊断或极小表
的全扫描对照。
统一查询引擎¶
ChronicleQueryEngine 是对外的只读 SQL 门面。Lance 与 ATIF 后端注册完全相同的
runs、steps、tool_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。即使查询只需要 source 或 COUNT(*),也会构造 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_id、step_id、source 的安全简单谓词 |
可提前裁剪,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_files、streamed_records、streaming_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_id、step_id、source 谓词传给
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 打开时间,以及列
裁剪、并行扫描和选择性索引收益。