Milvus DataNode 组件深度解析:从消息流到对象存储的数据落盘链路
【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus
DataNode(数据节点)是 Milvus 向量数据库中负责数据持久化的核心工作节点:它订阅分布式消息流中的 insert / delete 消息,将它们以 binlog / deltalog 的形式写入 MinIO、S3 等持久化 blob 存储,并承担数据同步(sync)、compaction(压缩合并)、批量导入(import)等数据面任务。本文以 internal/datanode/README.md 为骨架,结合仓库源码与 configs/milvus.yaml 配置,系统讲解 DataNode 的职责边界、依赖关系、内部结构、核心流程与关键配置,帮助读者理解"数据如何可靠落盘"这条关键链路。
一、DataNode 的定位与核心职责
在 Milvus 的存算分离架构中,各组件各司其职:Proxy 负责接入与校验请求,RootCoord 负责元数据与全局 ID 分配,QueryNode 负责查询与检索,而 DataNode 负责数据写入与持久化。
README 中对 DataNode 的定义非常凝练:
DataNode is the component to write insert and delete messages into persistent blob storage, for example MinIO or S3.
即:DataNode 是将 insert(插入)和 delete(删除)消息写入持久化 blob 存储(如 MinIO 或 S3)的组件。它在 internal/datanode/data_node.go 的包注释中也有对应描述:"Data node persists insert logs into persistent storage like minIO/S3"。
这意味着 DataNode 处于 Milvus 数据链路的下游——它不直接面向客户端,而是作为"数据消费者 + 落盘执行者",把上游源源不断产生的变更消息转化为对象存储上的物理文件,同时维护 segment 的元数据状态,为后续的 compaction、索引构建和查询提供数据基础。
从源码结构看,DataNode 是一个典型的"多任务汇聚"型组件。在 internal/datanode/data_node.go 中,DataNode结构体聚合了:
| 成员 | 类型 | 职责 |
|---|---|---|
syncMgr | syncmgr.SyncManager | 管理 binlog/deltalog 同步到对象存储 |
importTaskMgr/importScheduler | importv2.TaskManager/importv2.Scheduler | 管理批量导入任务 |
taskScheduler/taskManager | index.TaskScheduler/index.TaskManager | 管理索引、统计信息、分析等任务 |
externalCollectionManager | external.ExternalCollectionManager | 外部集合刷新管理 |
compactionExecutor | compactor.Executor | 执行 compaction 压缩任务 |
etcdCli/session | etcd client / session | 服务注册与发现、状态上报 |
对应地,internal/datanode目录下也按功能划分为 compactor、importv2、index、external 等子包,外加taskcost(任务成本统计)、util(工具函数)等辅助包。
二、DataNode 的四大依赖及其数据流角色
README 明确列出了 DataNode 运行所需的四个外部依赖,它们是理解 DataNode 行为的关键。下面逐一结合源码展开。
2.1 KV store:承载持久化 blob 存储
KV store: a kv store that persists messages into blob storage.
DataNode 的"KV store"实际上指的是底层的对象存储抽象层(ChunkManager),它以 KV(object key → value)语义对外提供持久化读写,物理载体是 MinIO、S3(含兼容 S3 的各家云厂商对象存储)等。
DataNode 通过StorageFactory创建 ChunkManager 实例,核心实现位于 internal/datanode/chunk_mgr_factory.go:NewChunkManager接收*indexpb.StorageConfig(包含存储类型、RootPath、地址、AccessKey/SecretKey、是否启用 SSL/SSL CA 证书、Bucket 名、是否使用 IAM、云厂商、IAM Endpoint、是否使用 VirtualHost、请求超时、Region、GCP 凭证 JSON、TLS 最低版本等),组装成objectstorage.Config后调用storage.NewChunkManagerFactory创建持久化存储客户端。
chunkManagerFactory := storage.NewChunkManagerFactory(config.GetStorageType(), objectstorage.RootPath(config.GetRootPath()), objectstorage.Address(config.GetAddress()), objectstorage.AccessKeyID(config.GetAccessKeyID()), objectstorage.SecretAccessKeyID(config.GetSecretAccessKey()), objectstorage.UseSSL(config.GetUseSSL()), objectstorage.SslCACert(config.GetSslCACert()), objectstorage.BucketName(config.GetBucketName()), objectstorage.UseIAM(config.GetUseIAM()), ... objectstorage.CreateBucket(true), )这段代码印证了:DataNode 对底层存储的访问是高度可配置的,既支持内网 MinIO,也支持各类 S3 兼容服务与云厂商对象存储,且默认CreateBucket(true)会自动创建 bucket。StorageConfig 在 compaction、import、snapshot restore 等场景中由上游(DataCoord)随请求下发,DataNode 按需为每个任务创建或复用 ChunkManager。
2.2 Message stream:消息流的订阅与消费
Message stream: receive messages and publish information.
DataNode 通过消息流(Message Stream,如 Pulsar / Kafka / RocksMQ)接收DML(Data Manipulation Language)变更消息——包括插入、删除操作,同时将处理进度(time tick)发布回消息流供上游感知。
README 提到的"receive messages and publish information"在实际代码中体现在两方面:
- 订阅通道:DataNode 以
dataNodeSubNamePrefix(默认dataNode,见 configs/milvus.yaml)为订阅前缀,消费 DataCoord 分配的 DML 通道。 - 时间戳同步:DataNode 定期发送 time tick 消息,将已消费数据的时间戳上报,保证数据可见性与流式读的时延窗口(对应配置
dataNode.timetick.interval,默认 500ms)。
值得一提的是,Milvus 后续版本将通道管理与数据同步下沉到 StreamingNode / flushcommon 层,WatchDmChannels等旧接口已在 DataNode 侧标记为 "not in use"(见 internal/datanode/services.go),数据消费统一由internal/flushcommon下的 pipeline 组件承载,但"订阅消息流、消费 DML 消息"这一本质职责没有改变。
2.3 Root Coordinator:全局唯一 ID 的来源
Root Coordinator: get the latest unique IDs.
写入对象存储的文件(如 binlog、deltalog)以及新增 segment 都需要全局唯一 ID,这些 ID 由 RootCoord 统一分配。DataNode 在数据同步、compaction 等流程中会向 RootCoord 申请 ID 段(ID range),确保分布式环境下多个 DataNode 生成的文件标识不冲突。
从依赖关系看,DataNode持有 RootCoord 的 gRPC 客户端(见 internal/datanode/data_node.go 的注释说明),同时 Milvus 在internal/allocator中提供了global_id_allocator.go等实现,用于向 RootCoord(或 TSO 服务)批量拉取 ID 并本地缓存分配。ID 的单调递增与全局唯一是对象文件路径、segment 标识可靠性的基础。
2.4 Data Coordinator:落盘计划与订阅信息的中枢
Data Coordinator: get the flush information and which message stream to subscribe.
DataCoord(数据协调节点)是 DataNode 的"调度大脑",它告诉 DataNode:
- flush 信息:何时对哪些 segment 执行 flush(即把内存缓冲的 binlog 刷写到对象存储);
- 订阅信息:DataNode 应该订阅哪些消息流通道。
在实际调用中,DataCoord 通过 gRPC 向 DataNode 下发各类任务与指令,DataNode 侧的处理入口集中在 internal/datanode/services.go:
CompactionV2(services.go):接收 DataCoord 下发的 compaction 计划,校验参数后创建对应类型的 Compactor 任务并入队执行;PreImport/ImportV2/QueryImport/DropImport(services.go):批量导入任务的创建、查询与删除;CopySegment/QueryCopySegment/DropCopySegment(services.go):快照恢复/跨集群复制时 segment 文件的拷贝管理;QuerySlot(services.go):向 DataCoord 上报本节点可用的任务槽位数(总槽位减去 index、compaction、import 已占用槽位),供 DataCoord 做任务调度决策。
此外,DataNode 会向 DataCoord 上报 segment 统计信息、同步进度等,形成"DataCoord 下发计划 → DataNode 执行并回传状态"的闭环。下图概括了四类依赖在数据写入链路中的位置:
Proxy(接入客户端请求) │ 生成 DML 消息(insert/delete) ▼ Message Stream(消息流,如 Pulsar/Kafka/RocksMQ) │ DataNode 订阅消费 ▼ DataNode ──► KV store / 对象存储(MinIO、S3…) ← 数据落盘 │ ▲ │ └── 向 RootCoord 申请全局唯一 ID │ └── 从 DataCoord 获取 flush 计划与通道订阅信息,并回传任务状态三、DataNode 的源码结构与核心工作流
3.1 组件生命周期:Init / Start / Stop
DataNode 实现了 Milvus 统一的组件生命周期接口(types.Component、types.DataNode),见 internal/datanode/data_node.go 的断言var _ types.DataNode = (*DataNode)(nil):
- Init(data_node.go):初始化 etcd 会话(
initSession,用于服务注册)、SyncManager、import 任务管理器与调度器、初始化 C++ 内核(index.InitSegcore)、分析器选项(analyzer.InitOptions),并按文件资源模式初始化 ChunkManager; - Start(data_node.go):启动 compaction executor、import scheduler、任务调度器,将节点状态置为
Healthy,并预热 goroutine 池、注册 Prometheus 池指标采集函数; - Stop(data_node.go):先置状态为
Abnormal并等待在途任务退出,依次关闭 SyncManager、会话、import scheduler、外部集合管理器、任务管理器与任务调度器,最后CloseSegcore并清理 compaction 指标。
组件启动入口位于 cmd/components/data_node.go(角色注册见 cmd/roles/roles.go),分布式部署时各节点通过 etcd 会话与心跳实现服务发现。
3.2 数据同步:SyncManager 与 SyncData
数据落盘的核心是internal/flushcommon/syncmgr包。SyncManager接口(sync_manager.go)暴露:
type SyncManager interface { SyncData(ctx context.Context, task Task, callbacks ...func(error) error) (*conc.Future[struct{}], error) SyncDataWithChunkManager(ctx context.Context, task Task, chunkManager storage.ChunkManager, callbacks ...func(error) error) (*conc.Future[struct{}], error) Close() error TaskStatsJSON() string }NewSyncManager会根据 CPU 核数与配置项dataNode.dataSync.maxParallelSyncMgrTasksPerCPUCore(默认 16)计算 worker 池大小(cpuNum * 每核并发数),并通过配置监听器支持运行时动态扩容/缩容(resizeHandler,见 sync_manager.go)。同步任务按 segment 维度加锁派发,保证同一 segment 的并发写操作有序执行,taskStats使用带 15 分钟过期时间的 LRU 缓存任务状态,供TaskStatsJSON查询。
配套的internal/flushcommon子包构成了完整的"消费-缓冲-落盘"管线:pipeline负责消息消费与处理流,writebuffer负责内存缓冲,metacache缓存 segment 元数据,io负责 binlog 的读写封装,broker封装对 DataCoord 的 RPC 调用。
3.3 后台任务:compaction、import 与 index
DataNode 不止是"搬运工",还承担大量数据面后台任务:
- Compaction(压缩合并):在
CompactionV2中按CompactionType分派任务,见 services.go:Level0DeleteCompaction处理 L0 删除日志合并,MixCompaction做混合合并(当启用 namespace 时会按主键 + 分区键排序),ClusteringCompaction做聚簇压缩,SortCompaction做排序压缩,BumpSchemaVersionCompaction用于 schema 版本升级。任务通过compactionExecutor.Enqueue提交,DataCoord 可查询槽位与任务状态。 - 批量导入(Import):
PreImport(预导入,解析并统计源文件)与ImportV2(实际导入)均通过importv2创建任务,支持 L0 导入与普通导入,由importScheduler按槽位并发调度,QueryImport可查询任务进度与生成的 segment 信息。 - 索引 / 统计 / 分析任务:通过
CreateTask的taskcommon.Index/Stats/Analyze分支创建对应任务,由index.TaskScheduler调度,调用 C++ 内核完成索引构建、统计信息收集(如 Text 匹配索引)与数据分布分析,并上报成本(耗时、CPU 数)。
统一任务框架CreateTask/QueryTask/DropTask(services.go)让 DataNode 成为真正意义上的"数据面任务执行器",DataCoord 只需按统一协议下发任务与查询状态即可。
四、DataNode 关键配置详解
DataNode 的配置集中在 configs/milvus.yaml 的dataNode段。下面按功能分组整理核心参数:
4.1 数据同步与流处理(dataSync / segment)
dataNode: dataSync: flowGraph: maxQueueLength: 16 # 流处理图中任务队列最大长度 maxParallelism: 1024 # 流处理图中最大并行执行任务数 maxParallelSyncMgrTasksPerCPUCore: 16 # SyncManager 每 CPU 核的最大并发同步任务数 skipMode: enable: true # 允许跳部分 timetick 消息以降低 CPU 占用 skipNum: 4 # 每跳过 n 条消息消费 1 条 coldTime: 60 # 仅剩 timetick 消息超过该秒数后开启跳过模式 ioConcurrency: 0 # 对象存储 I/O 池并发度,0/负值表示自动(CPU*2) segment: insertBufSize: 16777216 # 单个 binlog 内存缓冲上限(字节),超过即刷到 MinIO/S3 deleteBufBytes: 16777216 # 单通道 delete 日志刷盘缓冲上限(字节) syncPeriod: 600 # 缓冲非空时 segment 的定期同步周期(秒)其中insertBufSize直接决定"缓冲多少数据刷一次盘":设置过小会导致频繁小文件写入,设置过大会增加内存压力,是写入吞吐与内存之间的关键权衡点。
4.2 内存水位与强制同步(memory)
memory: forceSyncEnable: true # 内存占用过高时强制同步 forceSyncSegmentNum: 1 # 强制同步的 segment 数(优先选择缓冲最大的) checkInterval: 3000 # 内存检查间隔(毫秒) forceSyncWatermark: 0.5 # 单机内存水位线,达到后触发同步另外在 configs/milvus.yaml 中还有全局内存水位配置dataNodeMemoryLowWaterLevel: 0.85与dataNodeMemoryHighWaterLevel: 0.95,用于流量控制(写入降速)。
4.3 通道检查点(channel)
channel: workPoolSize: -1 # 所有通道的全局工作池大小,<=0 时取可执行 CPU 数 updateChannelCheckpointMaxParallel: 10 # 通道检查点更新的全局并行度,<=0 时取 10 updateChannelCheckpointInterval: 60 # 更新通道检查点的间隔(秒) updateChannelCheckpointRPCTimeout: 20 # UpdateChannelCheckpoint RPC 超时(秒) maxChannelCheckpointsPerPRC: 128 # 每次 RPC 携带的最大检查点数 channelCheckpointUpdateTickInSeconds: 10 # 检查点更新器执行频率(秒)通道检查点用于记录每个通道已消费到的时间戳,是故障恢复时数据不丢不重的基础。
4.4 导入与压缩任务(import / compaction / slot)
import: concurrencyPerCPUCore: 4 # 每 CPU 核的导入/预导入任务执行并发单元 maxImportFileSizeInGB: 16 # 单个导入文件大小上限(GB) readBufferSizeInMB: 16 # 导入基础读缓冲(MB),实际按分片数动态计算 readDeleteBufferSizeInMB: 16 memoryLimitPercentage: 10 # 导入任务可用的内存上限百分比 writeRetryInitialInterval: 1 # 导入写重试初始退避(秒) writeRetryMaxInterval: 60 # 导入写重试最大退避(秒) copyObjectTimeout: 3600 # 快照恢复时单个对象拷贝超时(秒),含重试 compaction: levelZeroBatchMemoryRatio: 0.5 # L0 批量压缩执行所需的最小空闲内存比例 levelZeroMaxBatchSize: -1 # L0 压缩单批最大 L1/L2 segment 数,<1 表示不限 useMergeSort: true # mix compaction 是否启用 mergeSort 模式 maxSegmentMergeSort: 30 # mergeSort 模式最大合并 segment 数 lobHoleRatioThreshold: 0.3 # TEXT 列压缩空洞率阈值,>= 阈值则重写 LOB 文件 text: inlineThreshold: 65536 # 小于该字节的 TEXT 值内联存储,不写入 LOB 文件 maxLobFileBytes: 67108864 # 单个 TEXT LOB 文件大小上限 flushThresholdBytes: 16777216 # TEXT 列写缓冲刷盘阈值 slot: slotCap: 16 # DataNode 上并发任务(compaction/import 等)上限slotCap与QuerySlot接口配合,是 DataCoord 做任务调度、防止 DataNode 过载的关键约束。
4.5 通信与存储格式(grpc / storage / 网络)
storage: format: parquet # insert 数据存储格式,可选 [parquet, vortex] deltalog: json # delete 日志格式,可选 [json, parquet] ip: # DataNode 监听地址,未指定时取第一个单播地址 port: 21124 # DataNode gRPC 端口 grpc: serverMaxSendSize: 536870912 # 单次 RPC 发送上限(字节) serverMaxRecvSize: 268435456 # 单次 RPC 接收上限(字节) clientMaxSendSize: 268435456 clientMaxRecvSize: 536870912 gracefulStopTimeout: 1800 # 优雅停止超时(秒),超时强制停止storage.format选择 insert 数据的物理文件格式(parquet 或 vortex),deltalog选择删除日志格式,直接决定存储层文件的组织方式与下游读取能力。此外,configs/milvus.yaml 中fileResource.dataNode: sync配置了 DataNode 的文件资源模式(sync/ref/close),与fileresource管理逻辑对应。
五、从源码验证"依赖即边界"
回到 README 的四个依赖,可以在代码中找到一一对应的证据:
| README 依赖 | 源码证据 | 对应文件 |
|---|---|---|
| KV store(持久化 blob 存储) | StorageFactory.NewChunkManager组装 MinIO/S3 客户端并CreateBucket(true) | internal/datanode/chunk_mgr_factory.go |
| Message stream(消息接收与发布) | 订阅前缀dataNodeSubNamePrefix: dataNode;time tick 间隔 500ms | configs/milvus.yaml、configs/milvus.yaml |
| Root Coordinator(获取最新唯一 ID) | DataNode持有 RootCoord gRPC 客户端,ID 由internal/allocator批量拉取缓存 | internal/datanode/data_node.go、internal/allocator |
| Data Coordinator(flush 信息与订阅) | CompactionV2/PreImport/ImportV2/QuerySlot等 RPC 入口 | internal/datanode/services.go |
对应的单元测试覆盖了同步、导入、索引等关键路径,例如 internal/datanode/data_node_test.go、internal/datanode/services_test.go、internal/flushcommon/syncmgr/sync_manager_test.go 中对SyncManager.SyncData的并发与回调验证,可作为深入阅读的起点。
六、小结
DataNode 是 Milvus 数据写入链路的"最后一公里":它以消息流为输入、以对象存储为输出,通过 SyncManager 的并发落盘、compaction 的空间回收、import 的批量写入和 index 任务的执行,把"流式写入"最终沉淀为可供查询与检索的持久化数据。理解其四大依赖(KV store、Message stream、RootCoord、DataCoord)与dataNode配置段,是排查写入延迟、优化刷盘策略、规划对象存储容量时的必备知识。生产环境中建议重点关注insertBufSize与内存水位配置的匹配、slotCap与机器规格的匹配,以及存储格式与下游工具的兼容性。
【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考