VictoriaLogs 部署和使用都比较直接,适合中小团队,也适用于日志量较大、对资源成本敏感的场景。官方资料称,与 Elasticsearch、Grafana Loki 等方案相比,其 RAM 占用最低可降至 1/30,磁盘占用最低可降至 1/15,并可在 Raspberry Pi 上运行。这些是官方给出的上限数据,实际收益取决于日志结构、查询负载和保留周期。
官方建议:如果可以接受在单节点上垂直扩展来满足业务需求,那么就不必使用集群模式。
本文包含大量源码片段。若不关注代码细节,可依次阅读 从数据模型看整体设计、架构概览、可查询、落盘与合并、从 HTTP 参数到查询计划 和 设计取舍与适用场景。
源码基线:本文以 VictoriaLogs commit
995d1b3(2025-12-18,晚于 v1.41.1 发布)为准;其go.mod固定使用 VictoriaMetrics commitb6bc186。代码片段为突出主流程有所删减,不保证可独立编译。
从数据模型看整体设计#
VictoriaLogs 把查询加速分成两个层次:先在 Stream 维度确定可能相关的数据,再在 Part 和 Block 内利用时间范围、稀疏索引、Bloom filter 与列编码继续缩小扫描范围。写入端按相反方向组织数据,让同一 Stream 的日志在物理上相邻,并通过批量写入和后台合并摊薄小写入的成本。
单节点与集群模式共用同一套本地存储引擎。单节点进程直接调用 logstorage.Storage;集群模式只是在它前面增加无状态的 vlinsert 和 vlselect,分别承担写入分片和查询扇出。后面的源码可以沿两条主线阅读:
- 一条日志怎样从协议输入变成 Partition、Part 和 Block;
- 一条 LogsQL 查询怎样从 Stream、时间和 Block 索引逐层排除无关数据。
设计目标与边界#
VictoriaLogs 的设计建立在几个日志场景中常见的前提上:日志持续追加,很少原地修改;大部分查询带有时间范围;用于定位应用实例的字段数量有限且相对稳定;消息、trace_id、用户 ID 等普通字段的基数可能很高。基于这些前提,它采用了以下设计:
- 写入先进入内存缓冲和不可变 Part,再由后台合并生成更大的 Part,避免每条日志都触发随机磁盘写入;
- 只为 Stream 标签维护从标签到
streamID的索引,普通字段主要依赖 Block 级 Bloom filter 和列扫描; - 查询优先按日期、Stream 和 Block 时间范围剪枝,只有剩余候选 Block 才会读取具体列值;
- 集群按
vlstorage节点分片数据,不在存储节点之间执行共识、自动复制或自动重平衡。
这些选择同时构成了使用边界。时间范围过宽、没有 Stream filter 的查询会扫描更多 Block;Stream fields 选得过细会制造高基数 Stream;集群中的单个 vlstorage 丢失,也意味着该节点独有的历史数据无法从其他节点自动恢复。
数据模型:Log、Tenant 与 Stream#
VictoriaLogs 在写入阶段把不同协议的输入统一为扁平字段集合。嵌套 JSON 的键会用 . 连接,数组、数字和布尔值会转换成字符串;空字符串按缺失字段处理。字段中有三个需要先区分的概念:
_time是日志时间。未提供有效值时,系统使用接收时间;它决定日志进入哪个 UTC 日分区;_msg是主要消息字段。它仍是一列,只是针对全文检索和编码做了专门处理;- Stream fields 用于标识产生日志的应用实例,例如
service、instance或一组 Kubernetes 元数据。其他字段仍保存在日志中,但不参与 Stream 身份计算。
下面这条日志可以贯穿后文的写入和查询过程:
{
"_time": "2025-12-18T10:15:30Z",
"_msg": "request failed",
"service": "checkout",
"instance": "10.0.2.17:8080",
"level": "error",
"trace_id": "4bf92f3577b34da6"
}假设租户为 7:3,并通过 _stream_fields=service,instance 指定 Stream fields。LogRows.MustAdd 会规范字段名,将 Stream 标签排序并编码为 canonical bytes,得到与 {instance="10.0.2.17:8080",service="checkout"} 对应的字节序列,再计算 128 位哈希。内部 streamID 由 TenantID 和这个哈希共同组成:TenantID 隔离租户,哈希标识租户内的标签集合。因此,标签相同但租户不同的日志不会属于同一 Stream。
level 和 trace_id 仍是普通字段。尤其不能把 trace_id 这类几乎每条日志都变化的字段放进 Stream;否则每个请求都可能创建一个新 Stream,indexdb、缓存和查询需要处理的 Stream 数量会快速增长。反过来,如果没有配置任何 Stream fields,所有日志都会落入空 Stream {}。功能仍然可用,但按应用实例过滤和压缩的效果都会变差。
两类索引各管一层#
文章后面会反复出现 indexdb 和 datadb。两者都包含“索引”,但作用不同:
官方文档所说的“自动索引所有字段”容易让人联想到所有字段值都进入同一张全局倒排表。VictoriaLogs 不是这样实现的:Stream fields 在 indexdb 中建立标签倒排关系,普通字段的索引信息随 Block 保存在 datadb,粒度更粗。
| 结构 | 负责回答的问题 | 主要内容 |
|---|---|---|
indexdb |
哪些 Stream 符合 Stream filter? | Stream 是否存在、streamID 到标签集合的映射、标签到 streamID 的倒排索引 |
datadb |
这些 Stream 在哪些时间段和 Block 中可能有匹配行? | Part 时间范围、indexBlockHeader、blockHeader、列头、Bloom filter 和列值 |
例如查询 _stream:{service="checkout"} AND level:=error 时,service="checkout" 可以先通过 indexdb 得到候选 streamID;level:=error 不会在 indexdb 中查找,而是在候选 Block 上先检查 level 列及其 Bloom filter,可能命中时才读取该列的 values。这个分工避免为每个高基数字段值维护全局倒排表,但也意味着普通字段过滤的效率更依赖时间范围、Stream filter 和 Block 的组织方式。
架构概览#
VictoriaLogs 的基础集群由 vlinsert、vlstorage 和 vlselect 组成,本身不在 vlstorage 节点之间复制数据。若要跨故障域保留副本,可以让 vlagent 将同一批日志分别写入两套独立集群。下图展示的是这种跨集群高可用拓扑,而不是单个集群的内部副本关系:
vlagent将日志扇出到两套独立集群的vlinsert,副本边界是集群,而不是单个vlstorage节点。- 每套集群内,
vlinsert根据路由策略为每行日志选择一个vlstorage节点;vlselect查询时再向本集群的全部vlstorage节点扇出请求并合并结果。 - 查询入口可通过负载均衡或故障切换机制选择健康集群。若同时查询两套副本,需要由上层处理重复结果。
- 扩容
vlstorage时不必迁移存量数据;新写入会分布到更新后的节点集合,查询仍覆盖全部节点。
组件拆分与故障语义#
三个集群组件的边界很明确:vlinsert 和 vlselect 不保存日志数据,可以按写入量和查询量分别扩容;vlstorage 持有本地 Partition 与 Part,是需要稳定磁盘和备份的状态节点。每个 vlstorage 同时具备单节点 VictoriaLogs 的本地写入与查询能力,所以已有单节点实例也可以加入 -storageNode 列表,被集群统一写入或查询。
系统没有负责副本放置和成员变更的控制面。vlinsert、vlselect 都从配置的 -storageNode 列表认识后端,数据属于最初接收它的 vlstorage。增加节点只改变后续写入的分布;已有 Part 不迁移。这样省去了重平衡期间的额外 IO,也让集群容量可以近似随节点线性增加,代价是每次完整查询都必须覆盖保存过数据的全部节点。
写入可用性与查询完整性采用了不同策略。某个 vlstorage 写入失败时,只要还有其他可用节点,vlinsert 就会继续尝试,因此新日志仍可进入集群;已经写在故障节点上的历史日志不会自动产生副本。查询默认要求所有 vlstorage 成功返回,只要一个节点不可用,vlselect 就返回 502 Bad Gateway。显式启用 allow_partial_response 后可以返回其余节点的结果,但调用方必须接受结果不完整。
因此,基础集群解决的是容量和吞吐扩展,不等于数据复制。跨故障域的完整高可用需要 vlagent 向独立集群分别写入,并由查询入口只选择一套健康副本。节点级备份仍然必要。
vlstorage 缩容需要额外处理存量数据。一种做法是:
- 先从
vlinsert的storageNode列表中移除待缩容节点,但仍将其保留在vlselect的节点列表中; - 利用 VictoriaLogs 的数据保留策略,让该节点上的数据逐步过期;
- 数据全部过期后,再从
vlselect中移除并释放该节点。
也可以参考 Partitions lifecycle,手动执行数据迁移。
vlagent#
vlagent 是部署在日志源附近的代理和缓冲层,负责接收日志并发送到 vlinsert。
它可以把同一批数据写入多个独立的 VictoriaLogs 实例或集群,以实现跨集群复制。
vlinsert#
vlinsert 是 VictoriaLogs 的无状态写入组件。它在局部性与均匀分布之间折中,为每行日志选择目标 vlstorage 节点。
代码片段
// app/vlstorage/netinsert/netinsert.go
func (s *Storage) AddRow(streamHash uint64, r *logstorage.InsertRow) {
idx := s.srt.getNodeIdx(streamHash)
// sns 保存 vlstorage 节点列表
sn := s.sns[idx]
// sn.addRow 最终通过网络写入该 vlstorage 节点
sn.addRow(r)
}
// app/vlstorage/netinsert/netinsert.go
func (srt *streamRowsTracker) getNodeIdx(streamHash uint64) uint64 {
if srt.nodesCount == 1 {
// 如果只有一个 vlstorage 节点,那么就直接写入到该节点
return 0
}
// streamHash := sid.id.lo ^ sid.id.hi,sid 即 streamID
// sid.id = hash128(bb.B) 就是 stream 的标签和值 hash 得来的
// 可以参见 LogRows.MustAdd 中的过程
streamRows := srt.rowsPerStream[streamHash] + 1
srt.rowsPerStream[streamHash] = streamRows
// Stream 的前 1000 条日志写入同一个 vlstorage 组件
// 对于只包含少量日志的 Stream,可提高局部性;当某个 Stream 量较大时则
// 分散到不同的 vlstorage 节点,以提高查询性能。
if streamRows <= 1000 {
return streamHash % uint64(srt.nodesCount)
}
return uint64(fastrand.Uint32n(uint32(srt.nodesCount)))
}
// app/vlstorage/netinsert/netinsert.go
func (sn *storageNode) mustSendInsertRequest(pendingData *bytesutil.ByteBuffer) {
// sn 代表当前选中的写入节点
if err := sn.sendInsertRequest(pendingData); err == nil {
return
}
// 初始节点写入失败后,从随机位置开始遍历节点列表重试
for !sn.s.sendInsertRequestToAnyNode(pendingData) {
// 省略
}
}
// app/vlstorage/netinsert/netinsert.go
func (sn *storageNode) sendInsertRequest(pendingData *bytesutil.ByteBuffer) error {
// 如果没有禁用压缩,就使用 zstd 压缩数据
var body io.Reader
if !sn.s.disableCompression {
bb := zstdBufPool.Get()
defer zstdBufPool.Put(bb)
bb.B = zstd.CompressLevel(bb.B[:0], pendingData.B, 1)
body = bb.NewReader()
} else {
body = pendingData.NewReader()
}
// 对 vlstorage 节点请求 internal/insert 接口
reqURL := sn.getRequestURL("/internal/insert")
req, err := http.NewRequestWithContext(ctx, "POST", reqURL, body)
// ...
resp, err := sn.c.Do(req)
// ...
}小结
在集群模式下,vlinsert 负责从 vlstorage 节点列表中选择写入目标,路由策略如下:
- 同一 Stream 的前 1000 条日志固定写入由
streamHash选定的节点,以保持局部性。 - 从第 1001 条开始,每行日志都会随机选择一个
vlstorage节点,以便在多个节点上并行查询大流量 Stream。 - 如果写入失败,会尝试写入到其他节点。
vlselect#
vlselect 是 VictoriaLogs 的查询组件,负责接收请求并汇总查询结果。下面以 /select/logsql/query 接口为例梳理调用链。
入口函数为 app/vlselect/logsql/logsql.go#ProcessQueryRequest。
代码片段
// app/vlselect/logsql/logsql.go
// ProcessQueryRequest handles /select/logsql/query request.
//
// See https://docs.victoriametrics.com/victorialogs/querying/#querying-logs
func ProcessQueryRequest(ctx context.Context, w http.ResponseWriter, r *http.Request) {
// 省略参数解析,metrics 采集等等
// Execute the query
if err := vlstorage.RunQuery(qctx, writeBlock); err != nil {
// 省略错误处理
return
}
}
// app/vlstorage/main.go
func RunQuery(qctx *logstorage.QueryContext, writeBlock logstorage.WriteDataBlockFunc) error {
// 这部分的目的是用来判断查询是否可以优化为:直接返回最后的 N 条日志,而不用先执行过滤再排序
// 比如:
// - 'sort by (_time desc) offset <offset> limit <limit>'
// - 'first <limit> by (_time desc)'
// - 'last <limit> by (_time)'
qOpt, offset, limit := qctx.Query.GetLastNResultsQuery()
if qOpt != nil {
qctxOpt := qctx.WithQuery(qOpt)
return runOptimizedLastNResultsQuery(qctxOpt, offset, limit, writeBlock)
}
// localStorage 在 vlstorage 节点中才会被赋值
if localStorage != nil {
return localStorage.RunQuery(qctx, writeBlock)
}
return netstorageSelect.RunQuery(qctx, writeBlock)
}vlselect 和 vlstorage 都会调用 RunQuery。先看 netstorageSelect 分支:
代码片段
// app/vlstorage/netselect/netselect.go
func (s *Storage) RunQuery(qctx *logstorage.QueryContext, writeBlock logstorage.WriteDataBlockFunc) error {
// nqr 代表 NetQueryRunner,用于执行分布式查询
nqr, err := logstorage.NewNetQueryRunner(qctx, s.RunQuery, writeBlock)
if err != nil {
return err
}
search := func(stopCh <-chan struct{}, q *logstorage.Query, writeBlock logstorage.WriteDataBlockFunc) error {
qctxLocal := qctx.WithQuery(q)
return s.runQuery(stopCh, qctxLocal, writeBlock)
}
// nqr.Run 实际执行的还是 Storage.runQuery 方法
// nqr.Run 其中会将 pipe 分成 remote pipe 和 local pipe
// 分别执行分布式查询和本地查询
concurrency := qctx.Query.GetConcurrency()
return nqr.Run(qctx.Context, concurrency, search)
}
// app/vlstorage/netselect/netselect.go
func (s *Storage) runQuery(stopCh <-chan struct{}, qctx *logstorage.QueryContext, writeBlock logstorage.WriteDataBlockFunc) error {
// ...
// runQuery 并发请求所有 vlstorage 节点
for i := range s.sns {
go func(nodeIdx int) {
err := sn.runQuery(qctxLocal, func(db *logstorage.DataBlock) {
writeBlock(uint(nodeIdx), db)
})
}(i)
}
// ...
}
func (sn *storageNode) runQuery(qctx *logstorage.QueryContext, processBlock func(db *logstorage.DataBlock)) error {
path := "/internal/select/query"
responseBody, reqURL, err := sn.getResponseBodyForPathAndArgs(qctx.Context, path, args)
if err != nil {
return err
}
defer responseBody.Close()
// 解析响应,省略
}小结
vlselect 会并发请求本集群的所有 vlstorage 节点,调用各节点的 /internal/select/query 接口。
这种机制允许集群直接加入新的 vlstorage 节点,无需重分布存量数据,但查询必须覆盖全部节点,因此也有相应代价:
- 木桶效应:查询响应时间受最慢节点影响;
- 即使某些节点不存在相关数据,仍会参与查询执行;
- 网络开销不可避免,延迟会增加。
这也解释了为什么 VictoriaLogs 建议优先评估单节点的纵向扩展能力,只有确有需要时再使用集群模式。
vlstorage#
vlstorage 是 VictoriaLogs 的存储组件,负责持久化日志,并向 vlselect 提供查询接口 /internal/select/query,向 vlinsert 提供写入接口 /internal/insert。
查询接口#
/internal/select/query 映射为 processQueryRequest 方法:
代码片段
// app/vlselect/internalselect/internalselect.go
func processQueryRequest(ctx context.Context, w http.ResponseWriter, r *http.Request) error {
// ...
// 调用与 vlselect 查询链相同的 RunQuery
if err := vlstorage.RunQuery(qctx, writeBlock); err != nil {
return err
}
}
// app/vlstorage/main.go
func RunQuery(qctx *logstorage.QueryContext, writeBlock logstorage.WriteDataBlockFunc) error {
// ...
// vlstorage 使用本地存储执行查询
if localStorage != nil {
return localStorage.RunQuery(qctx, writeBlock)
}
return netstorageSelect.RunQuery(qctx, writeBlock)
}localStorage.RunQuery 进入 VictoriaLogs 的本地存储层,具体细节见后文“存储原理”。
写入接口#
/internal/insert 对应的处理方法在 app/vlinsert/internalinsert/internalinsert.go#RequestHandler 中:
代码片段
// app/vlinsert/internalinsert/internalinsert.go
func RequestHandler(w http.ResponseWriter, r *http.Request) {
// ...
// CommonParams 包含了写入时的一些公共参数, 如:
// TenantID logstorage.TenantID
// TimeFields []string
// MsgFields []string
// StreamFields []string
// IgnoreFields []string
// 等等
cp, err := insertutil.GetCommonParams(r)
// 根据压缩类型解析 body 中的数据, 再通过 parseData 转为 InsertRow 类型
err = protoparserutil.ReadUncompressedData(r.Body, encoding, maxRequestSize, func(data []byte) error {
lmp := cp.NewLogMessageProcessor("internalinsert", false)
irp := lmp.(insertutil.InsertRowProcessor)
err := parseData(irp, data)
lmp.MustClose()
return err
})
}
// app/vlinsert/internalinsert/internalinsert.go
func parseData(irp insertutil.InsertRowProcessor, data []byte) error {
// 从对象池中获取 InsertRow 对象
r := logstorage.GetInsertRow()
src := data
i := 0
for len(src) > 0 {
tail, err := r.UnmarshalInplace(src)
src = tail
i++
// 将解析后的 InsertRow 添加到 InsertRowProcessor 中
irp.AddInsertRow(r)
}
}InsertRowProcessor 的实现是 logMessageProcessor:
代码片段
// app/vlinsert/insertutil/common_params.go
func (lmp *logMessageProcessor) AddInsertRow(r *logstorage.InsertRow) {
// 超过 MaxFieldsPerLine(默认 1000) 个字段的日志行会被丢弃
if len(r.Fields) > *MaxFieldsPerLine {
return
}
// 调用 logstorage.LogRows 将 InsertRow 添加到 LogRows 中
lmp.lr.MustAddInsertRow(r)
// 如果需要 flush,则调用 flushLocked 方法
if lmp.lr.NeedFlush() {
lmp.flushLocked()
}
}AddInsertRow 最终也会进入本地存储层,后文继续追踪。
存储原理#
根据官方文档 How does VictoriaLogs work?,VictoriaLogs 的存储设计有以下特点:
- 日志被作为 JSON 条目存储。
- 日志的字段会保存到不同的数据块中。
- 不同日志的相同字段会保存在同一个数据块中。
- 数据块会压缩存储,以减少磁盘空间占用。
- 较小的数据块会在后台合并成较大的数据块。
- 每个数据块都可以独立读取,查询时由多个 worker 并发处理。
为减少查询扫描量,VictoriaLogs 还采用了以下优化:
- 使用 Bloom filter 跳过不包含目标关键字的数据块。
- 针对不同数据类型的字段采用自定义编码和压缩。
- 同一 Stream 的数据在物理上相邻,Stream filter 可以排除无关数据块。
- 为时间范围维护稀疏索引,Time filter 可以进一步缩小扫描范围。
下面从源码中的写入和读取路径继续拆解这些设计。
这部分源码在 lib/logstorage/ 目录下。
写入流程#
先把完整路径串起来,再进入各层源码。以开头的 checkout 日志为例,它会经历以下步骤:
- JSON、Loki、OpenTelemetry、Syslog 等协议处理器把输入转换为时间戳和
[]Field。insertutil在这一层应用_msg_field、_time_field、_stream_fields、忽略字段和额外字段等写入参数; LogRows.MustAdd选出service、instance两个 Stream fields,生成 canonical Stream tags,并计算带 Tenant 信息的streamID;- 单节点模式直接调用本地
Storage.MustAddRows。集群模式先把LogRows转换为原生InsertRow,再由vlinsert按streamHash选择vlstorage; Storage.MustAddRows根据_time把批次拆到对应的 UTC 日 Partition。超出保留期、允许回填范围或未来时间范围的日志行会被跳过,而不是创建任意日期的目录;partition.mustAddRows先检查 Stream 是否已经注册。新 Stream 写入indexdb的三类索引条目,随后所有日志行进入datadb;datadb的分片缓冲把多个小批次合并成logRows。缓冲达到大小阈值或定时器触发后,日志按streamID和时间排序并转换为inmemoryPart;- 同一 Stream 的连续行被切成一个或多个 Block。Block 内识别常量列和普通列,分别写入时间戳、消息、列头、Bloom filter 与 values 数据;
- 内存 Part 继续合并,并按刷盘周期转成磁盘 Part。磁盘上的小 Part 和大 Part 还会在后台分层合并。
这条路径有两个重要结果。第一,每个 Partition 的 indexdb 只在首次遇到 Stream 时增加索引,后续同 Stream 日志可以通过缓存跳过重复注册。第二,排序、列转换和压缩都以批次或 Part 为单位完成,写入请求不需要为每行日志单独维护磁盘文件。
前面我们跟踪到 logMessageProcessor 中的 AddInsertRow 方法,其中调用了 flushLocked 方法:
代码片段
// app/vlinsert/insertutil/common_params.go
func (lmp *logMessageProcessor) flushLocked() {
logRowsStorage.MustAddRows(lmp.lr)
}
// app/vlstorage/main.go
// 这个方法在 vlinsert 中也被调用,只不过就是使用的 netstorageInsert.AddRow 方法
func (*Storage) MustAddRows(lr *logstorage.LogRows) {
if localStorage != nil {
localStorage.MustAddRows(lr)
} else {
// Store lr across the remote storage nodes.
lr.ForEachRow(netstorageInsert.AddRow)
}
}localStorage(*logstorage.Storage)是 VictoriaLogs 的本地存储层实现:
代码片段
// lib/logstorage/storage.go
type Storage struct {
path string // 存储目录的路径
retention time.Duration // 数据保留时间
flockF *os.File // 用于确保 Storage 仅被单个进程打开的文件锁
partitions []*partitionWrapper // 分区列表,按时间排序,例如 partitions[0] 是最早的分区
ptwHot *partitionWrapper // 最新的分区,用于写入新的日志行
// ... 省略其他字段
}partition 按 UTC 日期组织数据,每天对应一个分区。再看 MustAddRows:
代码片段
func (s *Storage) MustAddRows(lr *LogRows) {
// Fast path:
// 尝试将 LogRows 全部写入到 ptwHot 中
ptwHot := s.ptwHot
if ptwHot != nil {
if ptwHot.canAddAllRows(lr) {
ptwHot.pt.mustAddRows(lr)
return
}
}
// Slow path:
// 如果 LogRows 不能被 ptwHot 全部写入,那么就需要把 LogRows 拆分成多个 LogRows,写入到不同的分区中
// PS: 历史的分区可能需要从磁盘加载,会增加写入延迟
now := time.Now().UnixNano()
minAllowedDay := s.getMinAllowedDay(now) // 保留策略:过去的时间
maxAllowedDay := s.getMaxAllowedDay(now) // 保留策略:未来的时间
minAllowedTimestamp := now - s.maxBackfillAge.Nanoseconds() // 最大能接受的历史日志时间
// 遍历 LogRows 中的每个日志行,过滤掉不符合要求的日志行,
// 并根据日志行的时间戳,将其添加到同一分区(日期)的 LogRows 中
m := make(map[int64]*LogRows)
for i, ts := range lr.timestamps {
day := ts / nsecsPerDay
// 根据保留策略 和 最大能接受的历史日志时间 过滤掉不符合要求的日志行
// ...
lrPart := m[day]
lrPart.mustAddInternal(lr.streamIDs[i], ts, lr.rows[i], lr.streamTagsCanonicals[i])
}
for day, lrPart := range m {
ptw := s.getPartitionForWriting(day)
if ptw != nil {
ptw.pt.mustAddRows(lrPart)
}
}
}Storage 按日志时间戳所属的日期分区,具体写入由 partitionWrapper.partition 完成。
代码片段
// lib/logstorage/partition.go
type partition struct {
s *Storage
// path 是分区的完整目录路径。例如 /data/logstorage/partitions/20230801
// /data/logstorage 可配置,partitions/20230801 是 VictoriaLogs 的目录结构
path string
name string // name 是分区的名称,是目录名。如:20230801
idb *indexdb // idb 是索引数据库,用于存储日志行的索引信息
ddb *datadb // ddb 是数据数据库,用于存储日志行的原始数据
// ....
}
// lib/logstorage/partition.go
// mustAddRows 将数据写入 indexdb 和 datadb,并跳过已注册 Stream 的重复索引写入
func (pt *partition) mustAddRows(lr *LogRows) {
// 将 新增 的 Stream 注册到 indexdb 中
if !pt.idb.hasStreamID(streamID) {
streamTagsCanonical := streamTagsCanonicals[rowIdx]
pt.idb.mustRegisterStream(streamID, streamTagsCanonical)
}
// 最后把 LogRows 写入到 datadb 中
pt.ddb.mustAddRows(lr)
}每个 VictoriaLogs 分区主要由索引数据库 indexdb 和数据数据库 datadb 构成。下面依次跟踪两者的写入流程。
indexdb 写入#
indexdb 用于存储 Stream 索引,底层基于 mergeset.Table 实现。
代码片段
// lib/logstorage/indexdb.go
type indexdb struct {
path string // path 是索引数据库的目录路径:如 path/to/partition/indexdb 目录
partitionName string // partitionName 是索引数据库所属的分区名称(天)
tb *mergeset.Table // tb 是索引数据库的存储,用于存储索引信息
s *Storage // s 是 indexdb 所属的 Storage
}- 索引条目准备
mustRegisterStream 会写入三类索引条目:
tenantID:streamID:表示该 Stream 已注册的存在条目,不包含标签值;tenantID:streamID -> streamTagsCanonical:从 Stream ID 到规范化 Stream 标签的映射;tenantID:name:value -> streamID:从标签名和值到 Stream ID 的倒排索引。
代码片段
// lib/logstorage/indexdb.go
func (idb *indexdb) mustRegisterStream(streamID *streamID, streamTagsCanonical string) {
// Register tenantID:streamID entry.
bufLen := len(buf)
buf = marshalCommonPrefix(buf, nsPrefixStreamID, tenantID)
buf = streamID.id.marshal(buf)
items = append(items, buf[bufLen:])
// Register tenantID:streamID -> streamTagsCanonical entry.
bufLen = len(buf)
buf = marshalCommonPrefix(buf, nsPrefixStreamIDToStreamTags, tenantID)
buf = streamID.id.marshal(buf)
buf = append(buf, streamTagsCanonical...)
items = append(items, buf[bufLen:])
// Register tenantID:name:value -> streamIDs entries.
tags := st.tags
for i := range tags {
bufLen = len(buf)
buf = marshalCommonPrefix(buf, nsPrefixTagToStreamIDs, tenantID)
buf = tags[i].indexdbMarshal(buf)
buf = streamID.id.marshal(buf)
items = append(items, buf[bufLen:])
}
// Add items to the storage
idb.tb.AddItems(items)
}以一个具体例子说明。假设 Stream 为 {"tag1":"value1", "tag2":"value2"},tenantID 为 0,streamID 为 0x12345678(hi) 90abcdef(lo),会生成以下索引条目:
- 0x00(indexType) 0x00000000(tenantID) 0x12345678 0x90abcdef(streamID)
- 0x01(indexType) 0x00000000(tenantID) 0x12345678 0x90abcdef(streamID) 0x02(tagCount) 0x04(tag1NameLen) "tag1"(tag1Name) 0x06(tag1ValueLen) "value1" 0x04(tag2NameLen) "tag2" 0x06 "value2"
- 0x02(indexType) 0x00000000(tenantID) 0x04(tagNameLen) "tag1"(tagName) 0x06(tagValueLen) "value1"(tagValue) 0x12345678 0x90abcdef (streamID)
- 0x02(indexType) 0x00000000(tenantID) 0x04(tagNameLen) "tag2"(tagName) 0x06(tagValueLen) "value2"(tagValue) 0x12345678 0x90abcdef (streamID)- 写入内存块
这些索引条目随后加入 indexdb 的 mergeset.Table。写入端使用多个 raw-item shard 降低锁竞争,分片数为 cpus * min(cpus, 16);各分片先将条目追加到 inmemoryBlock,达到阈值后再触发合并。
代码片段
type rawItemsShard struct {
ibs []*inmemoryBlock
}
type inmemoryBlock struct {
commonPrefix []byte // 公共前缀,减少重复存储
data []byte
items []Item
}
// github.com/VictoriaMetrics/VictoriaMetrics/lib/mergeset/encoding.go
func (ib *inmemoryBlock) Add(x []byte) bool {
data := ib.data
// 如果写入当前索引条目后,内存块大小超过了 maxInmemoryBlockSize(64KB),
// 那么就不能写入当前索引条目,返回 false
if len(x)+len(data) > maxInmemoryBlockSize {
return false
}
dataLen := len(data)
data = append(data, x...)
ib.items = append(ib.items, Item{
Start: uint32(dataLen),
End: uint32(len(data)),
})
ib.data = data
return true
}向 inmemoryBlock 追加条目时,原始字节连续写入 data,items 只记录每条数据的起止位置。块达到 64 KiB 上限后,后续条目写入下一个块。例如:
{
data([]byte): [
0x00(indexType) 0x00000000(tenantID) 0x12345678 0x90abcdef(streamID) // 26B
0x01(indexType) 0x00000000(tenantID) 0x12345678 0x90abcdef(streamID) 0x02(tagCount) 0x04(tag1NameLen) "tag1"(tag1Name) 0x06(tag1ValueLen) "value1" 0x04(tag2NameLen) "tag2" 0x06 "value2" // 56B
],
items: [
Item{Start: 0, End: 26},
Item{Start: 26, End: 82},
],
}- 内存块合并
mergeset.Table 会把多个 inmemoryBlock 合并成 inmemoryPart;inmemoryPart 仍会在后台继续合并,并在合适的时机持久化。
单个 shard 累积到 maxBlocksPerShard(256)个块时,会把它们转移到共享的 ibsToFlush 队列;当待处理块达到 maxBlocksPerShard * availableCPUs 时,系统执行一次合并。
代码片段
// github.com/VictoriaMetrics/VictoriaMetrics/lib/mergeset/table.go
func (riss *rawItemsShards) addIbsToFlush(tb *Table, ibsToFlush []*inmemoryBlock) {
if len(ibsToFlush) == 0 {
return
}
var ibsToMerge []*inmemoryBlock
riss.ibsToFlushLock.Lock()
if len(riss.ibsToFlush) == 0 {
riss.updateFlushDeadline()
}
// 将当前分片的内存块加入到 ibsToFlush 中
riss.ibsToFlush = append(riss.ibsToFlush, ibsToFlush...)
// 如果待合并的内存块数量超过了 maxBlocksPerShard(256) * cpus 这一阈值,
// 那么就需要进行合并操作,将这些内存块合并成一个更大的内存块(inmemoryPart)
if len(riss.ibsToFlush) >= maxBlocksPerShard*cgroup.AvailableCPUs() {
ibsToMerge = riss.ibsToFlush
riss.ibsToFlush = nil
}
riss.ibsToFlushLock.Unlock()
// 将内存块合并成一个更大的内存块(inmemoryPart)
tb.flushBlocksToInmemoryParts(ibsToMerge, false)
}在 flushBlocksToInmemoryParts 中将 ibsToMerge 按最多 defaultPartsToMerge(16) 大小分成多个 chunk,再把这些 chunk 合并成一个更大的内存块(inmemoryPart)。注意单个 inmemoryPart 大小不能超过 5% 的可用系统内存。
代码片段
func (tb *Table) flushBlocksToInmemoryParts(ibs []*inmemoryBlock, isFinal bool) {
pws := make([]*partWrapper, 0, (len(ibs)+defaultPartsToMerge-1)/defaultPartsToMerge)
// 按最多 defaultPartsToMerge(16) 大小分成多个 chunk,
// 每个 chunk 合并成一个更大的内存块(inmemoryPart)
for len(ibs) > 0 {
n := defaultPartsToMerge
if n > len(ibs) {
n = len(ibs)
}
go func(ibsChunk []*inmemoryBlock) {
if pw := tb.createInmemoryPart(ibsChunk); pw != nil {
pws = append(pws, pw)
}
}(ibs[:n])
ibs = ibs[n:]
}
// 处理所有的 chunk (partWrapper), 将其合并成一个更大的内存块(inmemoryPart)
// 除非超过大小限制
maxPartSize := getMaxInmemoryPartSize()
for len(pws) > 1 {
// 合并 inmemoryPart
pws = tb.mustMergeInmemoryParts(pws)
pwsRemaining := pws[:0]
for _, pw := range pws {
// 如果大小超过 maxPartSize,那么就直接加入到 inmemoryParts 中
if pw.p.size >= maxPartSize {
tb.addToInmemoryParts(pw, isFinal)
} else {
pwsRemaining = append(pwsRemaining, pw)
}
}
pws = pwsRemaining
}
if len(pws) == 1 {
tb.addToInmemoryParts(pws[0], isFinal)
}
}默认的可用内存 = 系统内存 * 60% (可以通过
memory.allowedPercent配置)
合并不只是上面提到的部分,程序还会在后台对 inmemoryPart/filePart 进行合并。
- 持久化
经过合并后的 inmemoryPart 会被写入到磁盘中,如果满足以下条件,则会被刷入到磁盘中:
- flushInterval 定时触发
- part(inmemoryPart / filePart) 合并时
代码片段
// github.com/VictoriaMetrics/VictoriaMetrics/lib/mergeset/table.go
func (tb *Table) nextMergeIdx() uint64 {
return tb.mergeIdx.Add(1)
}
// github.com/VictoriaMetrics/VictoriaMetrics/lib/mergeset/table.go
func (tb *Table) mergeParts(pws []*partWrapper, stopCh <-chan struct{}, isFinal bool) error {
mergeIdx := tb.nextMergeIdx()
dstPartPath := ""
if dstPartType == partFile {
// 输出是磁盘 filePart,目录名格式为 %016X,如 18804EBAD6A80650
dstPartPath = filepath.Join(tb.path, fmt.Sprintf("%016X", mergeIdx))
}
// 只有一个 inmemoryPart 时,直接将其写入到磁盘中
if isFinal && len(pws) == 1 && pws[0].mp != nil {
mp := pws[0].mp
mp.MustStoreToDisk(dstPartPath)
// 打开新创建的 filePart
pwNew := tb.openCreatedPart(pws, nil, dstPartPath)
// 用新建的 filePart 替换参与合并的源 part
tb.swapSrcWithDstParts(pws, pwNew, dstPartType)
return nil
}
}inmemoryPart 会写入以 %016X 命名的磁盘目录(如 18804EBAD6A80650)。实际写文件由 (*inmemoryPart).MustStoreToDisk 完成。
代码片段
// github.com/VictoriaMetrics/VictoriaMetrics/lib/mergeset/inmemory_part.go
func (mp *inmemoryPart) MustStoreToDisk(path string) {
fs.MustMkdirFailIfExist(path)
metaindexPath := filepath.Join(path, metaindexFilename) // metaindex.bin
indexPath := filepath.Join(path, indexFilename) // index.bin
itemsPath := filepath.Join(path, itemsFilename) // items.bin
lensPath := filepath.Join(path, lensFilename) // lens.bin
var psw filestream.ParallelStreamWriter
psw.Add(metaindexPath, &mp.metaindexData)
psw.Add(indexPath, &mp.indexData)
psw.Add(itemsPath, &mp.itemsData)
psw.Add(lensPath, &mp.lensData)
// 并发执行前面添加的刷写任务
psw.Run()
mp.ph.MustWriteMetadata(path) // 写入当前 part 目录下的 metadata.json
fs.MustSyncPathAndParentDir(path)
}小结:
indexdb 位于 path/to/partitions/<YYYYMMDD>/indexdb,保存 (tenantID, streamID) 存在条目、streamID → streamTagsCanonical 映射,以及标签到 streamID 的倒排索引。
indexdb 的存储按照 part 单位进行管理,每个 part 对应一个文件夹(18804EBAD6A80650),每个文件夹中包含如下的文件:
| 文件名 | 描述 |
|---|---|
| metadata.json | (Part) 索引元数据: {索引数量,块数量, 第一个项的索引, 最后一个项的索引} |
| metaindex.bin | 存储 metaindexRow(用来索引 indexBlock 或者 blockHeaders) |
| index.bin | 存储 blockHeader (commonPrefix、items 数量、索引 items 偏移、索引 lens 偏移等) |
| items.bin | (Block) stream 索引 items 数据(会使用 commonPrefix 手段进行压缩) |
| lens.bin | (Block) lens 索引 items 长度信息(用于支持 commonPrefix 压缩) |
datadb 写入#
datadb 保存日志数据。pt.ddb.mustAddRows 最终调用 (*datadb).mustFlushLogRows 执行实际写入:
代码片段
// lib/logstorage/datadb.go
func (ddb *datadb) mustFlushLogRows(lr *logRows) {
mp.mustInitFromRows(lr)
p := mustOpenInmemoryPart(ddb.pt, mp)
// 将 inmemoryPart 加入到 datadb.inmemoryParts 中
ddb.inmemoryParts = append(ddb.inmemoryParts, pw)
// 触发 inmemoryParts 合并
ddb.startInmemoryPartsMergerLocked()
}其中 mp.mustInitFromRows(lr) 方法会将 logRows 转换为 inmemoryPart:
代码片段
// lib/logstorage/datadb.go
func (mp *inmemoryPart) mustInitFromRows(lr *logRows) {
sort.Sort(lr) // 根据 streamID 排序,再按 timestamp 排序
lr.sortFieldsInRows() // 根据 field 的 Name 进行排序
// 将 inmemoryPart 的各种 buffer 赋值给 blockStreamWriter,
// 后续 bsw.Finalize 实际上也是将 bsw 的数据写入到 inmemoryPart 中
bsw.MustInitForInmemoryPart(mp)
var sidPrev *streamID
uncompressedBlockSizeBytes := uint64(0)
timestamps := lr.timestamps
rows := lr.rows
streamIDs := lr.streamIDs
for i := range timestamps {
streamID := &streamIDs[i]
if sidPrev == nil {
sidPrev = streamID
}
// 注意:将相同流的日志写入同一个 block 中,如果超过 maxUncompressedBlockSize (2MB) 则会写入下一个 block
// 看 blockStreamWriter.MustWriteRows 方法
if uncompressedBlockSizeBytes >= maxUncompressedBlockSize || !streamID.equal(sidPrev) {
bsw.MustWriteRows(sidPrev, trs.timestamps, trs.rows)
trs.reset()
sidPrev = streamID
uncompressedBlockSizeBytes = 0
}
fields := rows[i]
trs.timestamps = append(trs.timestamps, timestamps[i])
trs.rows = append(trs.rows, fields)
uncompressedBlockSizeBytes += uint64(EstimatedJSONRowLen(fields))
}
bsw.MustWriteRows(sidPrev, trs.timestamps, trs.rows)
// 将 bsw 中的数据更新到 inmemoryPart
bsw.Finalize(&mp.ph)
}
// lib/logstorage/block_stream_writer.go
func (bsw *blockStreamWriter) MustWriteRows(sid *streamID, timestamps []int64, rows [][]Field) {
if len(timestamps) == 0 {
return
}
b := getBlock()
// 将原始日志字段写入 block
b.MustInitFromRows(timestamps, rows)
bsw.MustWriteBlock(sid, b)
putBlock(b)
}
type block struct {
timestamps []int64
columns []column // 所有 row 会根据按列进行存储
constColumns []Field // 所有 row 这一列的值相同就可以存储在 constColumns 中,但值不能超过 256 字节
}小结:
在这个环节,日志数据从行格式(logRows)转换为列格式(inmemoryPart),主要经历以下步骤:
- 对
logRows排序:同一stream的数据相邻并按时间排序,同时对Fields依据Field.Name进行排序,便于生成block; - 将排序后的
logRows中相同stream的数据写入一个或多个block(每个block最大maxUncompressedBlockSize(2MB)); block包含常量列(所有行该字段值相同)与普通列(各行值不完全相同);普通列保存所有行的该字段值(缺失则为空);- 将
block写入inmemoryPart,合并后刷入磁盘。
每个 part 包含下面几种文件:
| 文件名 | 存储内容 |
|---|---|
| metadata.json | (Part) part 元信息:记录数、格式版本、时间范围等 |
| metaindex.bin | (Part) 索引 indexBlockHeader (一组 block) |
| index.bin | (Part) 索引 blockHeader |
| column_names.bin | (Part) 中所有的列和列ID |
| column_idxs.bin | (Block) column 的 bloom/values 存在的 shardID |
| columns_header_index.bin | (Block) 索引 columnID 到 columnHeader 的偏移量 |
| columns_header.bin | (Block) columnHeader 数据 |
| timestamps.bin | (Block) 时间戳数据 |
| message_bloom.bin | (Block) 消息的布隆过滤器 |
| message_values.bin | (Block) 实际的日志消息 |
| bloom.bin{shardIdx} | (Block) 普通列的布隆过滤器,用于快速判断是否包含某个值 |
| values.bin{shardIdx} | (Block) 普通列的实际值 |
其中 {shardIdx} 是从 0 开始的分片索引,不是分片数量。例如 bloom.bin3 与 values.bin3 属于同一分片。
日志列名及其 ID 保存在 column_names.bin,column_idxs.bin 再记录各列的 Bloom filter/value 分片索引。
可查询、落盘与合并#
“写入成功”“可以查询”和“已经持久化”不是同一个时刻。理解这三个状态,可以解释写入延迟、刷盘参数和异常退出时的数据边界。
日志首先进入 rowsBuffer。它按可用 CPU 数量分片,减少并发写入时争用;单个 shard 在数据达到阈值时立即 flush,否则由一秒定时器触发。flush 会把 logRows 转成 inmemoryPart 并加入 Partition 的 Part 列表。到这一步后,查询可以读取该内存 Part,但数据还没有进入磁盘目录。
后台 flusher 再根据 -inmemoryDataFlushInterval 把到期的内存 Part 写成 file Part。该参数默认 5 秒,最小值为 1 秒。它描述的是内存数据获得磁盘持久性保证的周期:间隔越短,异常断电时可能留在内存中的窗口越小,但磁盘同步更频繁;间隔越长,写放大和闪存写入压力较低。正常关闭时,mustCloseDatadb 会先清空 rowsBuffer,等待后台任务退出,再强制把剩余内存 Part 写盘。
file Part 的发布过程不是直接覆盖旧文件。系统先在新目录写完数据文件和 metadata.json,同步目录内容,再打开新 Part;随后在内存列表中用新 Part 替换参与合并的源 Part,并原子更新 parts.json。旧 Part 要等引用计数归零后才会关闭和删除,因此正在执行的查询仍可读完它持有的版本。若进程在新目录创建完成、但 parts.json 更新前异常退出,重启时会清理未被清单引用的目录。
合并同样遵循不可变 Part 的思路。内存 Part、小型 file Part 和大型 file Part 分别选择大小接近的对象合并,减少用一个大 Part 反复吸收小 Part 造成的写放大。合并需要同时保留输入和输出数据,所以磁盘空间不足时,写入与查询都可能因合并跟不上而变慢。
读取流程#
写入路径决定了数据如何被剪枝,LogsQL 则决定了剪枝条件和后续计算如何组合。下面仍用开头的日志举例。假设需要查询 checkout 服务最近一小时的错误,并只返回时间、消息和 trace_id,查询可以写成:
_time:1h AND _stream:{service="checkout"} AND level:=error
| fields _time, _msg, trace_id
| sort by (_time) desc
| limit 100这条查询同时包含三种工作:_stream filter 用 indexdb 找 Stream,时间和 level filter 在存储层筛选 Block 与行,fields、sort、limit pipes 对匹配结果做投影和排序。源码没有把它们混在一个遍历循环中,而是先生成 Query,再把过滤阶段和 pipe 阶段连接起来。
从 HTTP 参数到查询计划#
ProcessQueryRequest 首先调用 parseCommonArgs。后者从请求头取得 AccountID、ProjectID,解析 query、start、end、time 和附加过滤条件,再通过 ParseQueryAtTimestamp 得到包含 filter 与 pipes 的 Query。HTTP 的 end 参数按右开区间处理,进入内部时间过滤前会减去 1 纳秒。
/select/logsql/query 的 limit 和 offset 参数也会转换为 pipes。如果查询适合“最后 N 条”优化,处理器会补上按 _time 降序排序,并把 offset/limit 合并进查询。解析器还会折叠相邻的 limit、offset 和 filter,或把 limit 下推到 sort、uniq 等支持提前截断的 pipe,避免后续阶段处理注定不会返回的行。
单节点直接执行这份 Query。集群模式下,NetQueryRunner 会把 pipe 链拆成远端部分和本地部分:能在各 vlstorage 独立完成的计算尽量下推,必须看到全局结果的排序、聚合或收尾逻辑留在 vlselect。拆分完成后,它还会从本地 pipes 反向推导需要的列,把字段投影加入远端查询。这样 vlstorage 不必把后续阶段不会使用的列通过网络传回来。
到达每个 vlstorage 后,Storage.runQuery 从 filter 中提取时间范围、Stream filter、普通字段 filter 和所需列,形成 storageSearchOptions。此后才进入 Partition、Part 和 Block 的本地搜索。
Block 过滤#
代码片段
// lib/logstorage/storage_search.go
func (s *Storage) runQuery(qctx *QueryContext, writeBlock writeBlockResultFunc) error {
// 从查询条件中解析出查询需要的参数:
// type storageSearchOptions struct {
// tenantIDs []TenantID
// streamIDs []streamID
// minTimestamp int64
// maxTimestamp int64
// streamFilter *StreamFilter
// filter filter
// fieldsFilter *prefixfilter.Filter
// hiddenFieldsFilter *prefixfilter.Filter
// timeOffset int64
// }
sso := s.getSearchOptions(qctx.TenantIDs, q, qctx.HiddenFieldsFilters)
s.searchParallel(workersCount, sso, qctx.QueryStats, stopCh, writeBlockToPipes)
}
// lib/logstorage/storage_search.go
func (s *Storage) searchParallel(workersCount int, sso *storageSearchOptions, qs *QueryStats, stopCh <-chan struct{}, writeBlock writeBlockResultFunc) {
// ...
// 启动多个 blockSearch 协程,每个协程从 workCh 中取出一个 blockSearchWorkBatch 进行搜索
workCh := make(chan *blockSearchWorkBatch, workersCount)
for workerID := 0; workerID < workersCount; workerID++ {
go func(workerID uint) {
for bswb := range workCh {
for i := range bsws {
// bs blockSearch 代表对一个 block 进行搜索
bs.search(qsLocal, bsw, bm)
if bs.br.rowsLen > 0 {
writeBlock(workerID, &bs.br)
}
}
}
}(uint(workerID))
}
// 根据时间范围确定需要查询的 Partition
ptws, ptwsDecRef := s.getPartitionsForTimeRange(sso.minTimestamp, sso.maxTimestamp)
defer ptwsDecRef()
// 并发查询每个 Partition
for i, ptw := range ptws {
go func(idx int, pt *partition) {
psfs[idx] = pt.search(sso, qsLocal, workCh, stopCh)
}(i, ptw.pt)
}
}
// lib/logstorage/storage_search.go
func (pt *partition) search(sso *storageSearchOptions, qs *QueryStats, workCh chan<- *blockSearchWorkBatch, stopCh <-chan struct{}) partitionSearchFinalizer {
// 如果查询条件中包含了 _streamFilter 那么从 indexdb 中获取到对应的 streamID
pso := pt.getSearchOptions(sso)
// 在 datadb 中执行查询
return pt.ddb.search(pso, qs, workCh, stopCh)
}
// lib/logstorage/storage_search.go
func (ddb *datadb) search(pso *partitionSearchOptions, qs *QueryStats, workCh chan<- *blockSearchWorkBatch, stopCh <-chan struct{}) partitionSearchFinalizer {
// 按照时间范围确定需要查询的 Part
pws, pwsDecRef := ddb.getPartsForTimeRange(pso.minTimestamp, pso.maxTimestamp)
// Apply search to matching parts
for _, pw := range pws {
pw.p.search(pso, qs, workCh, stopCh)
}
return pwsDecRef
}
// lib/logstorage/storage_search.go
func (p *part) search(pso *partitionSearchOptions, qs *QueryStats, workCh chan<- *blockSearchWorkBatch, stopCh <-chan struct{}) {
bhss := getBlockHeaders()
if len(pso.tenantIDs) > 0 {
p.searchByTenantIDs(pso, qs, bhss, workCh, stopCh)
} else {
// 这个方法是对 part 中的 index
p.searchByStreamIDs(pso, qs, bhss, workCh, stopCh)
}
putBlockHeaders(bhss)
}Storage.searchParallel 会并发调度命中的 partition,并启动一组 block worker。每个 partition 内部,datadb.search 按顺序遍历符合时间范围的 part,把候选 block 批次写入 workCh;真正的 block 匹配由 worker 并发执行。因此,并发发生在 partition 调度和 block 处理两个层面,part 并不是逐个启动独立协程。
其中 part.searchByStreamIDs 的作用是基于 streamIDs 过滤出符合条件的 block(blockHeader):要求 streamID 落在目标集合内且时间范围命中目标区间,并将其加入 workCh。
根据 streamFilter 获取 streamIDs 的代码片段如下。以 _stream: {service="test-app"} 为例,对应的等值查询方法为 getStreamIDsForNonEmptyTagValue:
代码片段
func (is *indexSearch) getStreamIDsForNonEmptyTagValue(tenantID TenantID, tagName, tagValue string) map[u128]struct{} {
ids := make(map[u128]struct{})
ts := &is.ts
kb := &is.kb
// 构建查询前缀:nsPrefixTagToStreamIDs + tenantID + tagName + tagValue
kb.B = marshalCommonPrefix(kb.B[:0], nsPrefixTagToStreamIDs, tenantID)
kb.B = marshalTagValue(kb.B, bytesutil.ToUnsafeBytes(tagName))
kb.B = marshalTagValue(kb.B, bytesutil.ToUnsafeBytes(tagValue))
prefix := kb.B
// 找到第一条以 prefix 开头的记录
ts.Seek(prefix)
for ts.NextItem() {
if !bytes.HasPrefix(item, prefix) {
break
}
// 解析 streamID 并更新结果
tail := item[len(prefix):]
sp.UpdateStreamIDs(ids, tail)
}
return ids
}part 内部还会经过两级过滤:
- 从
metaindex.bin加载indexBlockHeader; - 第一次过滤:通过
indexBlockHeader排除不包含目标streamID和时间范围的索引块; - 加载
blockHeader:从index.bin读取匹配索引块中的所有blockHeader; - 第二次过滤:依据
blockHeader再次排除不包含目标streamID和时间范围的块。
代码片段
// indexBlockHeader 覆盖多个 block
type indexBlockHeader struct {
streamID streamID // 覆盖范围内最小的 streamID
minTimestamp int64 // 覆盖范围内最小的 timestamp
maxTimestamp int64 // 覆盖范围内最大的 timestamp
indexBlockOffset uint64 // 在 index.bin 中的偏移量
indexBlockSize uint64 // 在 index.bin 中的数据大小
}
// blockHeader 保存单个 block 的元数据
type blockHeader struct {
streamID streamID // block 所属的 streamID
uncompressedSizeBytes uint64 // 未压缩数据大小
rowsCount uint64 // 日志行数
timestampsHeader timestampsHeader // 时间戳元数据
columnsHeaderIndexOffset uint64 // columnsHeader 在 columns_header_index.bin 中的偏移量
columnsHeaderIndexSize uint64 // columnsHeader 索引大小
columnsHeaderOffset uint64 // columnsHeader 在 columns_header.bin 中的偏移量
columnsHeaderSize uint64 // columnsHeader 数据大小
}
type timestampsHeader struct {
blockOffset uint64 // block 在 timestamps.bin 中的偏移量
blockSize uint64 // block 在 timestamps.bin 中的大小
minTimestamp int64 // block 内最小时间戳
maxTimestamp int64 // block 内最大时间戳
marshalType encoding.MarshalType // 时间戳编码类型
}经过时间范围和 streamID 的逐级过滤,系统得到需要扫描的 partition、part 和 blockHeader。候选 block 会被打包成 blockSearchWorkBatch,发送到 workCh,再由 blockSearch worker 执行匹配。
VictoriaLogs 在 partition、part、index block 和 block 层级都保留了可用于时间或 Stream 过滤的信息。查询中给出尽可能窄的时间范围和 Stream filter,通常能显著减少扫描量。
Block 匹配#
blockSearch 运行在 Storage.searchParallel 启动的 worker 中,从 workCh 读取 blockSearchWorkBatch 并处理其中的候选 block。
代码片段
// lib/logstorage/storage_search.go
func (bs *blockSearch) search(qs *QueryStats, bsw *blockSearchWork, bm *bitmap) {
bs.reset()
bs.qs = qs
bs.bsw = bsw
// bitmap 用来存储 block 中命中的 log entry 的‘索引’
bm.init(int(bsw.bh.rowsCount))
// 将所有的位标记为 1(初始状态,代表所有行都需要检测)
bm.setBits()
// filter 是一个接口,用来实现 log entries 的过滤
// 除 =(filterExact)外,还支持 AND(filterAnd)等组合 filter
bs.bsw.pso.filter.applyToBlockSearch(bs, bm)
// 如果没有命中任何 log entry,直接返回
if bm.isZero() {
return
}
// 将 block 中的
bs.br.mustInit(bs, bm)
// 获取需要的列,通过 “ | fields level, streamID, timestamp, message” 这种方式来指定
bs.br.initColumns(bsw.pso.fieldsFilter)
}
type filter interface {
String() string // 返回 filter 的字符串表示
updateNeededFields(pf *prefixfilter.Filter)
matchRow(fields []Field) bool
// 即根据 filter 过滤出 bs 中命中的 log entry 的‘索引’,并更新到 bm 中 (标记或者取消标记)
applyToBlockSearch(bs *blockSearch, bm *bitmap)
// 即根据 filter 过滤出 br 中命中的 log entry 的‘索引’,并更新到 bm 中(标记或者取消标记)
applyToBlockResult(br *blockResult, bm *bitmap)
}继续向下会进入各类 filter 的实现。本文只以精确匹配为例,观察 Block 如何检索列式数据。
举例:filterExact 匹配
// filterExact 匹配指定字段的精确值
//
// Example LogsQL: `fieldName:exact("foo bar")` of `fieldName:="foo bar"`
type filterExact struct {
fieldName string // 待匹配字段名
value string // 待匹配字段值
tokens []string
tokensHashes []uint64
}
// lib/logstorage/filter_exact.go
func (fe *filterExact) applyToBlockSearch(bs *blockSearch, bm *bitmap) {
fieldName := fe.fieldName
value := fe.value
// 先判断是否是 const 列 (所有行这一列的值都是相同的)。
// 如果是,可以直接比较 value 是否相等,且只需要比较一次就能确定是否命中
v := bs.getConstColumnValue(fieldName)
if v != "" {
if value != v {
bm.resetBits()
}
return
}
// 如果不是 const 列,检查列是否存在,并获取到 column header:
// 1. `column_names.bin` 提供 column_name 和 column_id 的映射
// 2. `columns_header_index.bin` 提供 column_id 到 column header 的偏移量映射
// 3. `column_idxs.bin` 提供 column_name 到 bloom/values shardId 的映射
ch := bs.getColumnHeader(fieldName)
if ch == nil {
// 列不存在的情况下,如果查询值不是空,那么需要重置 bitmap
if value != "" {
bm.resetBits()
}
return
}
tokens := fe.getTokensHashes()
switch ch.valueType {
case valueTypeString:
matchStringByExactValue(bs, ch, bm, value, tokens)
case valueTypeDict:
matchValuesDictByExactValue(bs, ch, bm, value)
case valueTypeUint8:
matchUint8ByExactValue(bs, ch, bm, value, tokens)
// ... 省略其他数值类型
case valueTypeIPv4:
matchIPv4ByExactValue(bs, ch, bm, value, tokens)
case valueTypeTimestampISO8601:
matchTimestampISO8601ByExactValue(bs, ch, bm, value, tokens)
default:
logger.Panicf("FATAL: %s: unknown valueType=%d", bs.partPath(), ch.valueType)
}
}
// columnHeader 代表列的元数据信息, 它一定只对应一个 Block 中的一列
// 注意如果 name 为空,那么它代表的 message(_msg) 列
type columnHeader struct {
name string
valueType valueType // 列中存储的值的类型
minValue uint64 // 列中存储的最小值, 用于快速判断是否在给定的范围内。适用于 uint*, ipv4, timestamp 和 float64 类型
maxValue uint64 // 列中存储的最大值, 用于快速判断是否在给定的范围内。适用于 uint*, ipv4, timestamp 和 float64 类型
valuesDict valuesDict // 列中存储的唯一值字典, 适用于 valueType = valueTypeDict 类型
valuesOffset uint64 // 列中存储的 values 的偏移量, 用于快速定位到 values.bin 中的对应位置
valuesSize uint64 // 列中存储的 values 的大小, 用于快速定位到 values.bin 中的对应位置
bloomFilterOffset uint64 // 列中存储的 bloom filter 的偏移量, 用于快速定位到 bloom.bin 中的对应位置
bloomFilterSize uint64 // 列中存储的 bloom filter 的大小, 用于快速定位到 bloom.bin 中的对应位置
}取得 column header 后,还要据此判断具体命中了哪些行:
values 匹配
// lib/logstorage/filter_exact.go
func matchStringByExactValue(bs *blockSearch, ch *columnHeader, bm *bitmap, value string, tokens []uint64) {
// 先判断 bloom filter 是否命中, 如果 bloom filter 不命中, 那么直接返回
if !matchBloomFilterAllTokens(bs, ch, tokens) {
bm.resetBits()
return
}
visitValues(bs, ch, bm, func(v string) bool {
return v == value
})
}
// lib/logstorage/filter_phrase.go
func visitValues(bs *blockSearch, ch *columnHeader, bm *bitmap, f func(value string) bool) {
if bm.isZero() {
// Fast path - nothing to visit
return
}
// 根据 列头 来获取 values
// 1. 根据列名确定 values.bin 所在的分片
// 2. 根据 columnHeader 中的 valuesOffset 和 valuesSize 读取 values
values := bs.getValuesForColumn(ch)
bm.forEachSetBit(func(idx int) bool {
return f(values[idx])
})
}完成过滤后,系统再按命中的行索引读取需要返回的字段值,本文不再展开。
Pipe 管线与结果合并#
Block 匹配产出的不是最终 HTTP 行,而是 blockResult。runPipes 从后向前为每个 pipe 创建 pipeProcessor,最终形成一条处理链。搜索 worker 把 blockResult 送入第一个 processor,结果再逐级流向 fields、stats、sort、limit 等后续阶段。无状态 pipe 可以边读边写;需要全局状态的 pipe 会按 worker 分片保存中间结果,并在 flush 阶段合并。
这套接口同时承担资源释放和提前结束。每层 processor 都收到 stopCh 与取消函数。查询出错、客户端断开,或 limit 已收集到足够结果时,可以取消上游扫描;排序、去重和聚合等有状态 pipe 还会使用内存预算,预算不足时停止继续接收 Block,避免查询无限占用内存。
在集群模式中,每个 vlstorage 以流式 DataBlock 返回远端 pipe 的结果,默认使用 ZSTD 压缩。vlselect 一边读取各节点响应,一边执行剩余的本地 pipes,最后由 ProcessQueryRequest 写成 application/stream+json。因此,结果是否有序取决于查询中的 pipe;并发扫描本身不保证日志顺序。普通 /query 请求满足 last-N 优化条件且带有 limit 时,处理器会显式添加时间倒序 pipe,而不是依赖存储层按全局时间返回。
默认情况下,只要一个 vlstorage 请求失败,整个集群查询就失败。这保证返回结果覆盖所有分片。启用 partial response 后,vlselect 可以忽略不可用节点继续处理,但响应只代表当时成功返回的分片。这一选项改变的是完整性约束,不会从其他节点补出故障节点的数据。
小结:
总结一下查询过程:
- 根据时间范围和 Stream 条件确定候选 partition 与 part;
- 并发调度 partition;每个 partition 内顺序遍历候选 part,并借助
indexBlockHeader和blockHeader筛出可能命中的 block; - 多个 block worker 并发处理候选 block;
- 对每个候选 block:
- 先检查匹配的字段是不是常量列,如果是,那么判断一次即可;
- 普通列先读取 column header;列不存在时可直接排除该 block;
- 列存在时先检查 Bloom filter。若未命中,也可排除该 block,无需读取 values;
- 最后只遍历 bitmap 中值为 1 的位置,而非扫描该列的全部 values;这些位置是前置 filter 保留下来的候选行。
存储模型#
本地启动 VictoriaLogs 并写入数据后,存储目录如下:
代码片段
/storage/
├── partitions/
│ ├── 20251210/
│ │ ├── indexdb
│ │ │ ├── 18804EBAD6A6ECA1
│ │ │ │ ├── index.bin
│ │ │ │ ├── items.bin
│ │ │ │ ├── lens.bin
│ │ │ │ ├── metadata.json
│ │ │ │ └── metaindex.bin
│ │ │ ├── 18804EBAD6A6ECB9
│ │ │ ├── ...
│ │ │ ├── 18804EBAD6A6ECB8
│ │ │ └── parts.json
│ │ │
│ │ └── datadb/
│ │ ├── 18804EBAD6A80650
│ │ │ ├── bloom.bin0
│ │ │ ├── bloom.bin1
│ │ │ ├── bloom.bin2
│ │ │ ├── bloom.bin3
│ │ │ ├── bloom.bin4
│ │ │ ├── column_idxs.bin
│ │ │ ├── column_names.bin
│ │ │ ├── columns_header.bin
│ │ │ ├── columns_header_index.bin
│ │ │ ├── index.bin
│ │ │ ├── message_bloom.bin
│ │ │ ├── message_values.bin
│ │ │ ├── metadata.json
│ │ │ ├── metaindex.bin
│ │ │ ├── timestamps.bin
│ │ │ ├── values.bin0
│ │ │ ├── values.bin1
│ │ │ ├── values.bin2
│ │ │ ├── values.bin3
│ │ │ └── values.bin4
│ │ │
│ │ ├── 18804EBAD6A80FF6
│ │ ├── ...
│ │ ├── 18804EBAD6A81004
│ │ └── parts.json
│ │
│ │── 20251211/
│ │── ...
│ └── 20251213/
│
└── flock.lock从目录和内存结构看,VictoriaLogs 存储可以分为四层:
Storage:存储引擎顶层目录,管理所有按日分区;Partition:对应一个 UTC 日期,内部包含indexdb和datadb;Part:后台合并形成的持久化或内存数据单元,保存数据与索引;Block:Part 内的最小查询任务单元。一个 Block 只属于一个 Stream,同一 Stream 可以跨多个 Block。
每个 Partition 都包含 indexdb 和 datadb:
indexdb保存(tenantID, streamID)存在条目、Stream 标签映射和标签倒排索引;datadb保存日志数据。同一 Stream 的日志按时间排序后相邻写入 Block,消息、时间戳和普通字段分别存储,并通过 Bloom filter 排除不可能命中的值。
对应的存储模型如下图所示:
生命周期与运维边界#
Partition 不只是目录层级,也是保留、迁移和快照操作的边界。Storage 启动时读取 partitions/ 下的 UTC 日期目录,分别打开其中的 indexdb 与 datadb;运行期间,新日期第一次出现时自动创建对应 Partition。批量写入跨越多个日期时,MustAddRows 会按天拆分,而不是把整个批次强行写进当前热分区。
写入前会检查三条时间边界:日志日期必须仍在 retentionPeriod 内,时间戳不能早于允许的 maxBackfillAge,也不能超过 futureRetention。被保留策略删除过、被手动 detach,或目录存在但尚未 attach 的 Partition 不会因为迟到日志而自动恢复写入,这可以避免运维中的旧分区被意外重新创建。
后台 retention watcher 以日分区为单位删除过期数据,不会逐行重写 Part。因此,实际删除边界以天为单位,换来的是低成本删除:关闭引用、移除整个 Partition 目录即可。配置磁盘空间或磁盘使用率上限后,系统也可以从最旧的日分区开始释放空间。这个机制优先保证进程可继续工作,但意味着容量阈值触发时,实际保留时间可能短于 retentionPeriod。
Partition 快照会先把内存 Part 刷到磁盘,再为当前 file Parts 创建硬链接。后台合并发布新 Part 时使用原子 parts.json 清单,重启只打开清单列出的目录,并清理未发布完成的临时 Part。Stream 和 filter 缓存则故意不跨重启持久化,因为分区可能在停机期间被删除、从备份恢复或从其他节点复制;重新构建缓存比使用过期映射更安全。
集群运维也沿用这一边界。扩容不搬迁旧分区,缩容可以等待节点数据自然过期,也可以按日创建快照、复制 Partition,再执行 detach/attach。前者简单但周期长,后者完成得快,却需要明确安排快照、复制、校验和查询节点列表变更,VictoriaLogs 不会自动替运维者完成这些步骤。
设计取舍与适用场景#
前面的实现说明,每项优化都对应一个前提,并非对所有查询都同样有效:
| 设计选择 | 获得的收益 | 需要承担的代价 |
|---|---|---|
| 按稳定字段划分 Stream | 同源日志相邻,压缩更好,Stream filter 可以先缩小范围 | Stream fields 选择错误会导致空 Stream 或高基数 |
| 普通字段使用 Block 级 Bloom filter 与列数据 | 不必为每个高基数值维护全局倒排表,写入和存储成本较低 | 没有时间或 Stream 条件时,仍可能检查大量 Block |
| 不可变 Part 加后台合并 | 写入以追加为主,查询可稳定持有旧版本 | 合并产生写放大,并需要额外 CPU、IO 和临时磁盘空间 |
vlinsert 对 vlstorage 分片而不复制 |
集群容量和写入吞吐随节点增加,拓扑简单 | 单节点历史数据没有内建副本,扩缩容不自动重平衡 |
vlselect 查询全部存储节点 |
无需维护全局分片目录,加入节点后即可参与查询 | 延迟受最慢节点影响,默认完整查询依赖所有节点可用 |
VictoriaLogs 更适合追加量大、查询通常带时间范围、可以稳定定义 Stream fields,并且希望控制基础设施成本的日志场景。查询主要按 service、instance、namespace 等维度定位,再过滤消息或 trace_id 时,Stream 索引、时间剪枝和 Block 过滤能够逐层发挥作用。
需要谨慎评估的则是另一类负载:查询经常横跨整个保留周期,几乎没有 Stream 条件;希望任意普通字段都具备全局倒排索引的响应特征;或者要求存储层自动复制、自愈和重平衡。VictoriaLogs 可以执行其中一部分查询,但它不会消除这些工作量,集群高可用也需要通过独立副本、备份和外部路由来补齐。
总结#
数据模型、物理布局和查询执行方式共同决定了 VictoriaLogs 的资源占用。Stream fields 把同源日志归到一起,indexdb 负责从标签定位 Stream;datadb 再按时间和 streamID 排列 Block,把常量列、普通列、消息和时间戳分别编码。写入端通过内存缓冲、不可变 Part 与分层合并换取连续写和更高压缩率。
查询沿着写入时留下的层级反向缩小范围:先确定 Tenant 和时间窗口,再用 Stream filter 查找 streamID,随后过滤 Partition、Part、index block 与 Block。普通字段先经过列存在性和 Bloom filter 检查,最后只读取候选行需要的列。filter 之间用 bitmap 传递候选位置,pipe processor 则完成投影、排序、聚合和 limit;在集群中,这些 pipes 还会拆成远端与本地两部分,兼顾下推计算和全局结果语义。
集群层没有改变本地存储模型。vlinsert 扩展写入并把数据分片到状态节点,vlselect 扩展查询并合并所有节点结果。这个结构容易扩容,也省去了副本协调和自动重平衡,但节点历史数据需要备份或跨集群复制来保护,完整查询默认依赖所有 vlstorage 可用。
源码中的 sync.Pool、slice 复用、分片缓冲和 goroutine worker 降低了分配与锁竞争,但它们是对上述设计的实现优化。实际运行效果还取决于 Stream fields 是否稳定、查询时间范围是否明确、合并是否有足够资源,以及集群是否按业务需要补齐了副本和故障切换。
源码仍在持续演进,阅读其他版本时应先核对本文开头标出的 commit。