VictoriaLogs 设计与实现

VictoriaLogs 部署和使用都比较直接,适合中小团队,也适用于日志量较大、对资源成本敏感的场景。官方资料称,与 Elasticsearch、Grafana Loki 等方案相比,其 RAM 占用最低可降至 1/30,磁盘占用最低可降至 1/15,并可在 Raspberry Pi 上运行。这些是官方给出的上限数据,实际收益取决于日志结构、查询负载和保留周期。

官方建议:如果可以接受在单节点上垂直扩展来满足业务需求,那么就不必使用集群模式。

本文包含大量源码片段。若不关注代码细节,可依次阅读 从数据模型看整体设计架构概览可查询、落盘与合并从 HTTP 参数到查询计划设计取舍与适用场景

源码基线:本文以 VictoriaLogs commit 995d1b3(2025-12-18,晚于 v1.41.1 发布)为准;其 go.mod 固定使用 VictoriaMetrics commit b6bc186。代码片段为突出主流程有所删减,不保证可独立编译。

从数据模型看整体设计#

VictoriaLogs 把查询加速分成两个层次:先在 Stream 维度确定可能相关的数据,再在 Part 和 Block 内利用时间范围、稀疏索引、Bloom filter 与列编码继续缩小扫描范围。写入端按相反方向组织数据,让同一 Stream 的日志在物理上相邻,并通过批量写入和后台合并摊薄小写入的成本。

单节点与集群模式共用同一套本地存储引擎。单节点进程直接调用 logstorage.Storage;集群模式只是在它前面增加无状态的 vlinsertvlselect,分别承担写入分片和查询扇出。后面的源码可以沿两条主线阅读:

  1. 一条日志怎样从协议输入变成 Partition、Part 和 Block;
  2. 一条 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 用于标识产生日志的应用实例,例如 serviceinstance 或一组 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 位哈希。内部 streamIDTenantID 和这个哈希共同组成:TenantID 隔离租户,哈希标识租户内的标签集合。因此,标签相同但租户不同的日志不会属于同一 Stream。

leveltrace_id 仍是普通字段。尤其不能把 trace_id 这类几乎每条日志都变化的字段放进 Stream;否则每个请求都可能创建一个新 Stream,indexdb、缓存和查询需要处理的 Stream 数量会快速增长。反过来,如果没有配置任何 Stream fields,所有日志都会落入空 Stream {}。功能仍然可用,但按应用实例过滤和压缩的效果都会变差。

两类索引各管一层#

文章后面会反复出现 indexdbdatadb。两者都包含“索引”,但作用不同:

官方文档所说的“自动索引所有字段”容易让人联想到所有字段值都进入同一张全局倒排表。VictoriaLogs 不是这样实现的:Stream fields 在 indexdb 中建立标签倒排关系,普通字段的索引信息随 Block 保存在 datadb,粒度更粗。

结构 负责回答的问题 主要内容
indexdb 哪些 Stream 符合 Stream filter? Stream 是否存在、streamID 到标签集合的映射、标签到 streamID 的倒排索引
datadb 这些 Stream 在哪些时间段和 Block 中可能有匹配行? Part 时间范围、indexBlockHeaderblockHeader、列头、Bloom filter 和列值

例如查询 _stream:{service="checkout"} AND level:=error 时,service="checkout" 可以先通过 indexdb 得到候选 streamIDlevel:=error 不会在 indexdb 中查找,而是在候选 Block 上先检查 level 列及其 Bloom filter,可能命中时才读取该列的 values。这个分工避免为每个高基数字段值维护全局倒排表,但也意味着普通字段过滤的效率更依赖时间范围、Stream filter 和 Block 的组织方式。

架构概览#

VictoriaLogs 的基础集群由 vlinsertvlstoragevlselect 组成,本身不在 vlstorage 节点之间复制数据。若要跨故障域保留副本,可以让 vlagent 将同一批日志分别写入两套独立集群。下图展示的是这种跨集群高可用拓扑,而不是单个集群的内部副本关系:

日志源经 vlagent 复制到两套独立的 VictoriaLogs 集群,查询入口通过负载均衡和故障切换访问任一健康集群
跨集群复制与查询故障切换拓扑
  • vlagent 将日志扇出到两套独立集群的 vlinsert,副本边界是集群,而不是单个 vlstorage 节点。
  • 每套集群内,vlinsert 根据路由策略为每行日志选择一个 vlstorage 节点;vlselect 查询时再向本集群的全部 vlstorage 节点扇出请求并合并结果。
  • 查询入口可通过负载均衡或故障切换机制选择健康集群。若同时查询两套副本,需要由上层处理重复结果。
  • 扩容 vlstorage 时不必迁移存量数据;新写入会分布到更新后的节点集合,查询仍覆盖全部节点。

组件拆分与故障语义#

三个集群组件的边界很明确:vlinsertvlselect 不保存日志数据,可以按写入量和查询量分别扩容;vlstorage 持有本地 Partition 与 Part,是需要稳定磁盘和备份的状态节点。每个 vlstorage 同时具备单节点 VictoriaLogs 的本地写入与查询能力,所以已有单节点实例也可以加入 -storageNode 列表,被集群统一写入或查询。

系统没有负责副本放置和成员变更的控制面。vlinsertvlselect 都从配置的 -storageNode 列表认识后端,数据属于最初接收它的 vlstorage。增加节点只改变后续写入的分布;已有 Part 不迁移。这样省去了重平衡期间的额外 IO,也让集群容量可以近似随节点线性增加,代价是每次完整查询都必须覆盖保存过数据的全部节点。

写入可用性与查询完整性采用了不同策略。某个 vlstorage 写入失败时,只要还有其他可用节点,vlinsert 就会继续尝试,因此新日志仍可进入集群;已经写在故障节点上的历史日志不会自动产生副本。查询默认要求所有 vlstorage 成功返回,只要一个节点不可用,vlselect 就返回 502 Bad Gateway。显式启用 allow_partial_response 后可以返回其余节点的结果,但调用方必须接受结果不完整。

因此,基础集群解决的是容量和吞吐扩展,不等于数据复制。跨故障域的完整高可用需要 vlagent 向独立集群分别写入,并由查询入口只选择一套健康副本。节点级备份仍然必要。

vlstorage 缩容需要额外处理存量数据。一种做法是:

  1. 先从 vlinsertstorageNode 列表中移除待缩容节点,但仍将其保留在 vlselect 的节点列表中;
  2. 利用 VictoriaLogs 的数据保留策略,让该节点上的数据逐步过期;
  3. 数据全部过期后,再从 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)
}

vlselectvlstorage 都会调用 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 日志为例,它会经历以下步骤:

  1. JSON、Loki、OpenTelemetry、Syslog 等协议处理器把输入转换为时间戳和 []Fieldinsertutil 在这一层应用 _msg_field_time_field_stream_fields、忽略字段和额外字段等写入参数;
  2. LogRows.MustAdd 选出 serviceinstance 两个 Stream fields,生成 canonical Stream tags,并计算带 Tenant 信息的 streamID
  3. 单节点模式直接调用本地 Storage.MustAddRows。集群模式先把 LogRows 转换为原生 InsertRow,再由 vlinsertstreamHash 选择 vlstorage
  4. Storage.MustAddRows 根据 _time 把批次拆到对应的 UTC 日 Partition。超出保留期、允许回填范围或未来时间范围的日志行会被跳过,而不是创建任意日期的目录;
  5. partition.mustAddRows 先检查 Stream 是否已经注册。新 Stream 写入 indexdb 的三类索引条目,随后所有日志行进入 datadb
  6. datadb 的分片缓冲把多个小批次合并成 logRows。缓冲达到大小阈值或定时器触发后,日志按 streamID 和时间排序并转换为 inmemoryPart
  7. 同一 Stream 的连续行被切成一个或多个 Block。Block 内识别常量列和普通列,分别写入时间戳、消息、列头、Bloom filter 与 values 数据;
  8. 内存 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
}
  1. 索引条目准备

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"}tenantID0streamID0x12345678(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)
  1. 写入内存块

这些索引条目随后加入 indexdbmergeset.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 追加条目时,原始字节连续写入 dataitems 只记录每条数据的起止位置。块达到 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},
	],
}
  1. 内存块合并

mergeset.Table 会把多个 inmemoryBlock 合并成 inmemoryPartinmemoryPart 仍会在后台继续合并,并在合适的时机持久化。

单个 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 进行合并。

  1. 持久化

经过合并后的 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),主要经历以下步骤:

  1. logRows 排序:同一 stream 的数据相邻并按时间排序,同时对 Fields 依据 Field.Name 进行排序,便于生成 block
  2. 将排序后的 logRows 中相同 stream 的数据写入一个或多个 block(每个 block 最大 maxUncompressedBlockSize(2MB));
  3. block 包含常量列(所有行该字段值相同)与普通列(各行值不完全相同);普通列保存所有行的该字段值(缺失则为空);
  4. 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.bin3values.bin3 属于同一分片。

日志列名及其 ID 保存在 column_names.bincolumn_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 与行,fieldssortlimit pipes 对匹配结果做投影和排序。源码没有把它们混在一个遍历循环中,而是先生成 Query,再把过滤阶段和 pipe 阶段连接起来。

从 HTTP 参数到查询计划#

ProcessQueryRequest 首先调用 parseCommonArgs。后者从请求头取得 AccountIDProjectID,解析 querystartendtime 和附加过滤条件,再通过 ParseQueryAtTimestamp 得到包含 filter 与 pipes 的 Query。HTTP 的 end 参数按右开区间处理,进入内部时间过滤前会减去 1 纳秒。

/select/logsql/querylimitoffset 参数也会转换为 pipes。如果查询适合“最后 N 条”优化,处理器会补上按 _time 降序排序,并把 offset/limit 合并进查询。解析器还会折叠相邻的 limitoffset 和 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 过滤出符合条件的 blockblockHeader):要求 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 内部还会经过两级过滤:

  1. metaindex.bin 加载 indexBlockHeader
  2. 第一次过滤:通过 indexBlockHeader 排除不包含目标 streamID 和时间范围的索引块;
  3. 加载 blockHeader:从 index.bin 读取匹配索引块中的所有 blockHeader
  4. 第二次过滤:依据 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 行,而是 blockResultrunPipes 从后向前为每个 pipe 创建 pipeProcessor,最终形成一条处理链。搜索 worker 把 blockResult 送入第一个 processor,结果再逐级流向 fieldsstatssortlimit 等后续阶段。无状态 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 可以忽略不可用节点继续处理,但响应只代表当时成功返回的分片。这一选项改变的是完整性约束,不会从其他节点补出故障节点的数据。

小结:

总结一下查询过程:

  1. 根据时间范围和 Stream 条件确定候选 partition 与 part;
  2. 并发调度 partition;每个 partition 内顺序遍历候选 part,并借助 indexBlockHeaderblockHeader 筛出可能命中的 block;
  3. 多个 block worker 并发处理候选 block;
  4. 对每个候选 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 日期,内部包含 indexdbdatadb
  • Part:后台合并形成的持久化或内存数据单元,保存数据与索引;
  • Block:Part 内的最小查询任务单元。一个 Block 只属于一个 Stream,同一 Stream 可以跨多个 Block。

每个 Partition 都包含 indexdbdatadb

  • indexdb 保存 (tenantID, streamID) 存在条目、Stream 标签映射和标签倒排索引;
  • datadb 保存日志数据。同一 Stream 的日志按时间排序后相邻写入 Block,消息、时间戳和普通字段分别存储,并通过 Bloom filter 排除不可能命中的值。

对应的存储模型如下图所示:

VictoriaLogs 存储按 Storage、按日 Partition、indexdb 和 datadb Part、Block 逐层组织,并在 Block 中分离索引、时间戳、消息与列文件
VictoriaLogs 从 Storage 到文件的存储层级

生命周期与运维边界#

Partition 不只是目录层级,也是保留、迁移和快照操作的边界。Storage 启动时读取 partitions/ 下的 UTC 日期目录,分别打开其中的 indexdbdatadb;运行期间,新日期第一次出现时自动创建对应 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 和临时磁盘空间
vlinsertvlstorage 分片而不复制 集群容量和写入吞吐随节点增加,拓扑简单 单节点历史数据没有内建副本,扩缩容不自动重平衡
vlselect 查询全部存储节点 无需维护全局分片目录,加入节点后即可参与查询 延迟受最慢节点影响,默认完整查询依赖所有节点可用

VictoriaLogs 更适合追加量大、查询通常带时间范围、可以稳定定义 Stream fields,并且希望控制基础设施成本的日志场景。查询主要按 serviceinstance、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。

参考#

访问量 访客数