GreptimeDB 支持给 Prometheus 做 remote read 后端。这条路径的最后一步是把查询引擎产出的列式 RecordBatch,转成 Prometheus 协议里行式的 TimeSeries:按 label 组合把行分组,同一条时间线的采样点聚到一起。

图 1:查询引擎产出的列式 RecordBatch,转成 Prometheus 协议的行式 TimeSeries
这个函数(recordbatches_to_timeseries)在仓库里躺了很久没人动过,直到我们的 committer @lyang24(去年我们专访过他)提了 PR #8587 把它重写了一遍,只改了 src/servers/src/prom_store.rs 一个文件。
ℹ️ Note: GreptimeDB v1.2.0-beta.1 已经发布,带上了本次优化。
我是看到这个 PR 才去读的这段代码。PR 描述里配了一张作者画的优化示意图,右下角带着一份 CPU profile:300k series、4 并发读的场景下,RecordBatch → TimeSeries 这一步占了 36.2% 的 CPU,其中光 dictionary materialization 就占 8.9%。一个既不解压也不读盘、纯粹在内存里搬数据的 函数吃掉三分之一多的 CPU,这个比例显然不正常。
profile 是作者提供的,我没有复现那套环境。想知道改动本身带来多少优化,还得有个能跑的基准,可 bench 代码没有跟着 PR 进仓库,于是我照着入口重写了一个,在改动前后各跑一遍。优化前的数字也不 好看:一万行数据转一次要 3.7 毫秒,十万行要 44 毫秒,每秒处理不到 300 万行。优化后最好的情况是 现在的十几倍。
本文是我对这个 PR 的拆解,路子跟《Rust 解码 Protobuf 数据比 Go 慢五倍?记一次性能调优之旅》 一样:先写 benchmark 把问题复现出来,再逐段读代码。benchmark 代码放在 这个 gist 里。
Step1:复现
先把优化前的状态固定下来。bench 直接压测公共入口 recordbatches_to_timeseries,用 criterion 跑:
fn bench_prom_read_convert(c: &mut Criterion) {
let mut group = c.benchmark_group("prom_read_convert");
group.measurement_time(Duration::from_secs(5));
// (series count, samples per series)
let sizes = [(10, 1000), (100, 100), (1000, 100), (5000, 20)];
for encoding in [Encoding::Dictionary, Encoding::Utf8] {
for ordering in [Ordering::Adjacent, Ordering::Interleaved] {
for (series, samples) in sizes {
let rows = series * samples;
let (schema, batch) = build_recordbatch(series, samples, ordering, encoding);
group.throughput(Throughput::Elements(rows as u64));
group.bench_with_input(
BenchmarkId::new(
format!("{}/{}", encoding.name(), ordering.name()),
format!("{series}x{samples}"),
),
&(schema, batch),
|b, (schema, batch)| {
b.iter(|| {
let batches =
RecordBatches::try_new(schema.clone(), vec![batch.clone()])
.unwrap();
black_box(
recordbatches_to_timeseries("bench_metric", batches).unwrap(),
)
});
},
);
}
}
}
group.finish();
}三个 for 循环里的维度是照着真实场景挑的。编码上,Utf8 是普通字符串列,Dictionary<UInt32, Utf8> 是 PromQL 读路径实际返回的 label 布局。排列上,adjacent 表示同一条时间线的采样点在结果里连续 (查询输出一般按 primary key 排序,就是这个形状),interleaved 表示时间线逐行交错。规模上是 时间线数乘每条时间线的采样点数,remote read 的典型请求是时间线不多、每条点很多。label 按真实 metric 表的形状模拟:host 每条时间线都不一样,datacenter、env、job 只有几个取值。
cargo bench -p servers --bench prom_read_convert结果如下:
prom_read_convert/utf8/adjacent/100x100
time: [3.6421 ms 3.6581 ms 3.6759 ms]
thrpt: [2.7204 Melem/s 2.7337 Melem/s 2.7457 Melem/s]
prom_read_convert/dict/adjacent/100x100
time: [4.9186 ms 4.9330 ms 4.9482 ms]
thrpt: [2.0209 Melem/s 2.0272 Melem/s 2.0331 Melem/s]
prom_read_convert/dict/adjacent/1000x100
time: [55.435 ms 55.574 ms 55.719 ms]
thrpt: [1.7947 Melem/s 1.7994 Melem/s 1.8039 Melem/s]thrpt 那行是 criterion 报的吞吐,这里一个 element 就是一行,也就是每秒 200 万行上下。
除了慢,还有个地方不太对劲:字典编码那组比普通字符串那组慢了 35%(2.0272M vs 2.7337M)。字典 编码本来是用来省内存和省拷贝的,怎么反而更慢?我们先记下这个疑点。
Step2:这个函数每一行都在干什么
打开旧代码,第一步是 collect_timeseries_ids。函数头上原本就挂着一句自嘲的注释,写它的人显然知道它不好,只是当时没有更好的下手点:
/// Collect each row's timeseries id
/// This processing is ugly, hope <https://github.com/GreptimeTeam/greptimedb/issues/336> making some progress in future.
fn collect_timeseries_ids(table_name: &str, recordbatch: &RecordBatch) -> Vec<TimeSeriesId> {它给每一行构造一个 owned 的 label 向量:
for row in 0..row_count {
let mut labels = Vec::with_capacity(recordbatch.num_columns() - 1);
labels.push(new_label(METRIC_NAME_LABEL.to_string(), table_name.to_string()));
for (column_name, column_values) in columns.iter() {
if let Some(value) = &column_values[row] {
labels.push(new_label((*column_name).clone(), value.clone()));
}
}
timeseries_ids.push(TimeSeriesId { labels });
}而这里的 columns,是提前把每一列整列物化出来的:
.map(|(i, column_name)| {
(column_name, recordbatch.iter_column_as_string(i).collect::<Vec<_>>())
})第二步拿这些 id 当 key 做分组:
let mut timeseries_map: BTreeMap<&TimeSeriesId, TimeSeries> = BTreeMap::default();
for (row, timeseries_id) in timeseries_ids.iter().enumerate() {
let timeseries = timeseries_map
.entry(timeseries_id)
.or_insert_with(|| TimeSeries {
labels: timeseries_id.labels.clone(),
..Default::default()
});
// push sample
}因为 protobuf 生成的 Label 没有实现 Eq,TimeSeriesId 还得手写一整套 PartialEq / Eq / Hash / Ord 实现。
代码读起来是清楚的,问题在于开销的量级。用 R 表示行数,S 表示时间线数,L 表示平均 label 数,这一趟下来:整列物化是 R·L 次 to_string();每行一个 Vec<Label>,是 R 次堆分配;每个 label 的 name 和 value 各克隆一次,又是 R·L 次字符串拷贝;BTreeMap 每行查一次,每次 log(S) 量级的比较,虽然比较会在第一个不同的 label 上短路,但最坏情况下要把整个向量比完。

图 2:旧实现的开销都挂在行数 R 上,绝大部分花在反复重建同一份 label 组合
也就是说,所有开销都是按行数计算的,可结果里真正的信息量只有 S 条时间线。Remote read 的典型 请求是一条时间线拉几十上百个点,R 是 S 的几十倍,绝大部分分配和拷贝都花在反复重建同一份 label 组合上。
ℹ️ Note: 这里的
R和S都是单个RecordBatch内的数字。分组是逐个 batch 做的,recordbatches_to_timeseries只是把每个 batch 的结果 flatten 起来,同一条时间线如果跨了 batch,不会在这一步合并(新旧实现都是如此)。
Step3:字典列为什么反而更慢
回头看 Step1 留下的疑点。PromQL 读路径返回的 label 列是 Dictionary<UInt32, Utf8>, 1000 行里 host 可能只有 3 个不同取值,字典里就存 3 个字符串,每行存一个 u32 下标。这本来是 #8541 特意保留下来的编码。
但 iter_column_as_string 不认识字典列,走的是通用兜底路径,逐行 to_string()。3 个字符串被 展开成了 1000 个,上游省下来的分配在这里又都做了一遍,还多了一次字典解引用的开销。比普通字符串 列更慢就是这么来的——开头那份 profile 里单独列出的 8.9%,就是这一段。

图 3:上游特意保留的字典编码,被下游一句 iter_column_as_string 展开了
到这里 PR 的思路就清楚了:分组只需要读某一行的值拿去比较,那就别拷贝,直接从 Arrow 数组里“借”。 只有确认“这是一条从没见过的新时间线”的时候,才真正分配。
Step4:从 Arrow 数组借,不物化
先解决拷贝。PR 给 label 列包了一层借用视图:
enum LabelValues<'a> {
Utf8(&'a StringArray),
LargeUtf8(&'a LargeStringArray),
Utf8View(&'a StringViewArray),
DictionaryUtf8 {
dictionary: &'a DictionaryArray<UInt32Type>,
values: &'a StringArray,
},
Other(Vec<Option<String>>),
}
impl LabelValues<'_> {
fn value(&self, row: usize) -> Option<&str> {
match self {
Self::Utf8(values) => values.is_valid(row).then(|| values.value(row)),
Self::LargeUtf8(values) => values.is_valid(row).then(|| values.value(row)),
Self::Utf8View(values) => values.is_valid(row).then(|| values.value(row)),
Self::DictionaryUtf8 { dictionary, values } => dictionary
.key(row)
.and_then(|key| values.is_valid(key).then(|| values.value(key))),
Self::Other(values) => values.get(row).and_then(Option::as_deref),
}
}
}label 列实际会碰到的四种 Arrow 布局 (StringArray、LargeStringArray、StringViewArray,以及字典编码的 DictionaryArray<UInt32Type>)各占一支,value(row) 返回的是从数组里借来的 &str,不拷贝也不分配。LabelValues<'a> 持有的全是 &'a 引用,「这些 &str 活不过 RecordBatch」这件事就交给 Rust 编译器管了,不用任何约定。

图 4:四种 Arrow 布局各占一支,按行取值直接从数组里借 —— 跟图 3 对照着看
字典那一支就是 Step3 想要的:dictionary.key(row) 拿下标,values.value(key) 借字典里那几个 字符串之一。1000 行还是要做 1000 次下标访问,但复用的始终是字典里那 3 份字符串,不会展开成 1000 个 owned String。
非字符串类型的 label 列(比如 Int32)还是得物化,落到 Other,这条路上每行仍然会产生一个 String。PR 顺手还改掉了一处原来糟糕的错误处理:
ensure!(
values.len() == recordbatch.num_rows(),
error::InvalidPromRemoteReadQueryResultSnafu {
msg: format!(
"Cannot convert label column '{}' of datatype {:?} to string",
column_schema.name,
array.data_type()
),
}
);这里我跟 PR 描述有个分歧,作者写的是这个改动避免了 silently omitting that column,也就是认为旧 行为是静默忽略这一列。但照代码看不太像:iter_column_as_string 遇到转不成 GreptimeDB vector 的列时返回的是一个空迭代器,collect 出来是个空 Vec,旧代码紧接着按行下标去取 column_values[row],只要 batch 非空就会越界 panic,静默不了。不管是哪一种,换成一个带列名和 数据类型的 Err 都比原来强——出问题的时候至少知道是哪一列。
Step5:分组不用 BTreeMap
拷贝解决了,还剩每行一次的树查找和每行一个的 Vec 分配。新的循环是这样(下面的 LabelColumn 就是列名加上前面那个 LabelValues):
let columns = label_columns(&recordbatch)?;
let mut timeseries: Vec<TimeSeries> = Vec::new();
let mut timeseries_by_hash: HashMap<u64, Vec<usize>> = HashMap::new();
let mut previous_timeseries: Option<usize> = None;
for row in 0..recordbatch.num_rows() {
let timeseries_index = match previous_timeseries {
Some(index) if matches_timeseries(×eries[index].labels, &columns, row) => index,
_ => {
let hash = hash_timeseries(&columns, row);
let candidates = timeseries_by_hash.entry(hash).or_default();
match candidates
.iter()
.copied()
.find(|index| matches_timeseries(×eries[*index].labels, &columns, row))
{
Some(index) => index,
None => {
let index = timeseries.len();
timeseries.push(new_timeseries(table, &columns, row));
candidates.push(index);
index
}
}
}
};
previous_timeseries = Some(timeseries_index);
// push sample into timeseries[timeseries_index]
}最外层那个 Some(index) if matches_timeseries(...) 是冲着数据形状去的。Mito 的 SeqScan 在每个 PartitionRange 内按 primary key 和时间排序,同一条时间线的采样点大体是连着的,所以先拿 这行跟上一行落到的时间线比一下,命中就直接复用,hash 都不用算。这只是个针对常见局部顺序的优化, 跨 partition、跨 range 都可能断掉,正确性不依赖它。
快路径断掉的时候,以及每条时间线第一次出现的时候,都走 hash 兜底。这里的 map 值是候选列表而不是 单个下标,因为 hash 会碰撞,找到候选之后必须再做一次完整的 label 比对。为了快而把两条不同的 时间线合并到一起,这种错误是不能犯的。

图 5:新实现里一行数据的三条去向 —— 快路径的命中率决定了提速倍数
真的判断是新时间线时,才走到 new_timeseries。走上面那四种布局的 label 列,这里的 to_string() 就是整趟下来唯一一次拷贝:
fn new_timeseries(table: &str, columns: &[LabelColumn<'_>], row: usize) -> TimeSeries {
let mut labels = Vec::with_capacity(columns.len() + 1);
labels.push(new_label(METRIC_NAME_LABEL.to_string(), table.to_string()));
for (name, value) in row_labels(columns, row) {
labels.push(new_label(name.to_string(), value.to_string()));
}
TimeSeries { labels, ..Default::default() }
}R·L 次拷贝变成 S·L 次(注意:仍然是单个 batch 内的数字)。整个优化就是把分配从“每一次迭代” 挪到了“发现新东西的那次迭代”。
Step6:别把行为改了
这是协议层的代码,改动尽量别动外部用户可观察的行为,免得给下游埋兼容性问题。PR 里有三处特别注意 了这点。
旧实现用 BTreeMap,输出天然按 label 有序;新实现用 Vec,顺序是首次出现顺序。所以末尾补了 一次排序,比较规则照抄旧的 Ord for TimeSeriesId:
timeseries.sort_unstable_by(|left, right| compare_timeseries_labels(&left.labels, &right.labels));S·log(S) 的一次排序,跟省下来的开销比可以忽略。
NULL 的语义也得一样。旧代码 if let Some(value) 把 NULL 的 label 跳过,等于这个 label 不存在。 新代码用 filter_map 保持一致,字典列还多一种情况:key 是 NULL,或者 key 指向的 value 是 NULL,这就是上面那两层 and_then 的来历。
既然 NULL label 会被跳过,两条时间线的 label 序列就可能一个是另一个的前缀,比较时要小心:
fn matches_timeseries(labels: &[Label], columns: &[LabelColumn<'_>], row: usize) -> bool {
let mut labels = labels.iter().skip(1);
for (name, value) in row_labels(columns, row) {
let Some(label) = labels.next() else {
return false;
};
if label.name != name || label.value != value {
return false;
}
}
labels.next().is_none()
}最后那句 labels.next().is_none() 就是干这个的,少了它 {host=a} 和 {host=a, env=prod} 会被 当成同一条。开头的 skip(1) 跳过的是特殊 label __name__,它由表名决定,同一个 RecordBatch 里都一样,没必要比。
结果
Apple M4 Max、macOS 上跑的对比。before 是旧实现,after 是合并后的 32a6fc0915,两者只差 prom_store.rs 一个文件。criterion 每个 case 100 samples,所有 case 的 change 都是 p = 0.00 < 0.05。
字典编码 label(Dictionary<UInt32, Utf8>,PromQL 读路径的真实布局):
| 排列 | 时间线×采样点 | before | after | 提速 | change |
|---|---|---|---|---|---|
| adjacent | 10×1000 | 4.29 ms | 292 µs | 14.7× | −93.2% |
| adjacent | 100×100 | 4.93 ms | 335 µs | 14.7× | −93.2% |
| adjacent | 1000×100 | 55.6 ms | 3.38 ms | 16.5× | −93.9% |
| adjacent | 5000×20 | 61.8 ms | 4.98 ms | 12.4× | −91.9% |
| interleaved | 10×1000 | 4.31 ms | 820 µs | 5.3× | −81.0% |
| interleaved | 100×100 | 4.67 ms | 870 µs | 5.4× | −81.3% |
| interleaved | 1000×100 | 49.7 ms | 9.03 ms | 5.5× | −81.8% |
| interleaved | 5000×20 | 54.1 ms | 10.4 ms | 5.2× | −80.8% |
普通 Utf8 label:
| 排列 | 时间线×采样点 | before | after | 提速 | change |
|---|---|---|---|---|---|
| adjacent | 10×1000 | 3.08 ms | 288 µs | 10.7× | −90.8% |
| adjacent | 100×100 | 3.66 ms | 319 µs | 11.5× | −91.3% |
| adjacent | 1000×100 | 43.6 ms | 3.19 ms | 13.7× | −92.7% |
| adjacent | 5000×20 | 49.8 ms | 4.65 ms | 10.7× | −90.7% |
| interleaved | 10×1000 | 3.12 ms | 770 µs | 4.1× | −75.4% |
| interleaved | 100×100 | 3.44 ms | 832 µs | 4.1× | −75.6% |
| interleaved | 1000×100 | 37.9 ms | 8.46 ms | 4.5× | −77.7% |
| interleaved | 5000×20 | 42.3 ms | 10.1 ms | 4.2× | −76.2% |
吞吐从每秒 180 万~320 万行提到 960 万~3480 万行,没有一个 case 变慢。
adjacent 连续的那几组快得多(10~16 倍),来自 Step5 那条快路径:只有每条时间线的第一行需要算 hash,占比大概是 S/R。interleaved 每行都得算 hash、查表、比对,也还有 4~5 倍。实际的查询 结果大体是 adjacent 这个形状。
字典列提升最大,因为它原先的额外开销最大。同样是 adjacent/1000×100,字典列从 55.6 ms 降到 3.38 ms,跟普通字符串列的 3.19 ms 基本拉平,Step1 里那个“字典比字符串还慢”的反常总算没了。
小结
数据结构确实换了(BTreeMap 换成 Vec 加 hash 索引,多了一层 previous 快路径,有序输出改成 末尾统一排序),但收益主要不来自这些,而是两件很朴素的事:能“借”的别拷贝,能只做一次的别每行做 一遍。
当前 hash_timeseries 用的是 DefaultHasher,该 toolchain 下是固定 key 的 SipHash-1-3,算法 本身不在稳定 API 的承诺范围内。换 ahash 之类可以作为一个性能实验。末尾那次全量排序纯粹是 为了跟以前的 BTreeMap 的顺序保持一致,如果我们能确认下游不依赖它,也能省掉。
倒是 Step3 那个字典列的坑值得留意:上游为了省内存特意保留了字典编码,下游一句 iter_column_as_string 就把它展开了,还展开得比不用字典更慢。这种事情不写 benchmark 是很难发现的。
自己复现
bench 代码在这个 gist。 PR 只动了一个文件,对比不用重编整个仓库:
# 0) 取 bench 文件放进 src/servers/benches/,并在 src/servers/Cargo.toml 末尾加上
# [[bench]]
# name = "prom_read_convert"
# harness = false
# 1) 在改动前的父提交上开一个临时 worktree
git worktree add /tmp/gdb-8587-before a0e6f0b4
cp src/servers/benches/prom_read_convert.rs /tmp/gdb-8587-before/src/servers/benches/
# 2) 跑 before 基线
cd /tmp/gdb-8587-before
cargo bench -p servers --bench prom_read_convert -- \
--warm-up-time 2 --measurement-time 4 --save-baseline before
# 3) 回到含 PR 的代码跑 after,criterion 自动给出 change%
cd -
cargo bench -p servers --bench prom_read_convert -- \
--warm-up-time 2 --measurement-time 4 --baseline before本文的数据是在同一个工作树里跑的:先跑完 before,再把 prom_store.rs 换成 32a6fc0915 的重新 编译跑 after,criterion 数据目录共用。
参考链接
- 本文拆解的 PR(@lyang24 贡献):greptimedb#8587 perf(servers): optimize PromQL read conversion
- benchmark 代码和完整数据:gist: prom_read_convert.rs
- 优化后的源码:
src/servers/src/prom_store.rs - 上游保留字典编码的 PR:#8541
- 旧代码注释里挂了多年的 issue:#336
- 上一篇性能调优:《Rust 解码 Protobuf 数据比 Go 慢五倍?记一次性能调优之旅》
- PR 作者的专访:《“我不是为了头衔,而是不想忘了 Rust”|专访 GreptimeDB Committer @lyang24》
- 带上本次优化的版本:GreptimeDB v1.2.0-beta.1:JSON2 首次发布,支持 Prometheus Remote Write v2
- Prometheus remote read API:https://prometheus.io/docs/prometheus/latest/querying/remote_read_api/
- Arrow 数组类型:
StringArray、StringViewArray、DictionaryArray - criterion.rs 手册:https://bheisler.github.io/criterion.rs/book/index.html


