十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Milvus DataNode 组件深度解析:从消息流到对象存储的数据落盘链路

Milvus DataNode 组件深度解析:从消息流到对象存储的数据落盘链路 Milvus DataNode 组件深度解析从消息流到对象存储的数据落盘链路【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvusDataNode数据节点是 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结构体聚合了成员类型职责syncMgrsyncmgr.SyncManager管理 binlog/deltalog 同步到对象存储importTaskMgr/importSchedulerimportv2.TaskManager/importv2.Scheduler管理批量导入任务taskScheduler/taskManagerindex.TaskScheduler/index.TaskManager管理索引、统计信息、分析等任务externalCollectionManagerexternal.ExternalCollectionManager外部集合刷新管理compactionExecutorcompactor.Executor执行 compaction 压缩任务etcdCli/sessionetcd 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它以 KVobject key → value语义对外提供持久化读写物理载体是 MinIO、S3含兼容 S3 的各家云厂商对象存储等。DataNode 通过StorageFactory创建 ChunkManager 实例核心实现位于 internal/datanode/chunk_mgr_factory.goNewChunkManager接收*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接收DMLData 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 的调度大脑它告诉 DataNodeflush 信息何时对哪些 segment 执行 flush即把内存缓冲的 binlog 刷写到对象存储订阅信息DataNode 应该订阅哪些消息流通道。在实际调用中DataCoord 通过 gRPC 向 DataNode 下发各类任务与指令DataNode 侧的处理入口集中在 internal/datanode/services.goCompactionV2services.go接收 DataCoord 下发的 compaction 计划校验参数后创建对应类型的 Compactor 任务并入队执行PreImport/ImportV2/QueryImport/DropImportservices.go批量导入任务的创建、查询与删除CopySegment/QueryCopySegment/DropCopySegmentservices.go快照恢复/跨集群复制时 segment 文件的拷贝管理QuerySlotservices.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 / StopDataNode 实现了 Milvus 统一的组件生命周期接口types.Component、types.DataNode见 internal/datanode/data_node.go 的断言var _ types.DataNode (*DataNode)(nil)Initdata_node.go初始化 etcd 会话initSession用于服务注册、SyncManager、import 任务管理器与调度器、初始化 C 内核index.InitSegcore、分析器选项analyzer.InitOptions并按文件资源模式初始化 ChunkManagerStartdata_node.go启动 compaction executor、import scheduler、任务调度器将节点状态置为Healthy并预热 goroutine 池、注册 Prometheus 池指标采集函数Stopdata_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 与 indexDataNode 不止是搬运工还承担大量数据面后台任务Compaction压缩合并在CompactionV2中按CompactionType分派任务见 services.goLevel0DeleteCompaction处理 L0 删除日志合并MixCompaction做混合合并当启用 namespace 时会按主键 分区键排序ClusteringCompaction做聚簇压缩SortCompaction做排序压缩BumpSchemaVersionCompaction用于 schema 版本升级。任务通过compactionExecutor.Enqueue提交DataCoord 可查询槽位与任务状态。批量导入ImportPreImport预导入解析并统计源文件与ImportV2实际导入均通过importv2创建任务支持 L0 导入与普通导入由importScheduler按槽位并发调度QueryImport可查询任务进度与生成的 segment 信息。索引 / 统计 / 分析任务通过CreateTask的taskcommon.Index/Stats/Analyze分支创建对应任务由index.TaskScheduler调度调用 C 内核完成索引构建、统计信息收集如 Text 匹配索引与数据分布分析并上报成本耗时、CPU 数。统一任务框架CreateTask/QueryTask/DropTaskservices.go让 DataNode 成为真正意义上的数据面任务执行器DataCoord 只需按统一协议下发任务与查询状态即可。四、DataNode 关键配置详解DataNode 的配置集中在 configs/milvus.yaml 的dataNode段。下面按功能分组整理核心参数4.1 数据同步与流处理dataSync / segmentdataNode: 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 内存水位与强制同步memorymemory: forceSyncEnable: true # 内存占用过高时强制同步 forceSyncSegmentNum: 1 # 强制同步的 segment 数优先选择缓冲最大的 checkInterval: 3000 # 内存检查间隔毫秒 forceSyncWatermark: 0.5 # 单机内存水位线达到后触发同步另外在 configs/milvus.yaml 中还有全局内存水位配置dataNodeMemoryLowWaterLevel: 0.85与dataNodeMemoryHighWaterLevel: 0.95用于流量控制写入降速。4.3 通道检查点channelchannel: workPoolSize: -1 # 所有通道的全局工作池大小0 时取可执行 CPU 数 updateChannelCheckpointMaxParallel: 10 # 通道检查点更新的全局并行度0 时取 10 updateChannelCheckpointInterval: 60 # 更新通道检查点的间隔秒 updateChannelCheckpointRPCTimeout: 20 # UpdateChannelCheckpoint RPC 超时秒 maxChannelCheckpointsPerPRC: 128 # 每次 RPC 携带的最大检查点数 channelCheckpointUpdateTickInSeconds: 10 # 检查点更新器执行频率秒通道检查点用于记录每个通道已消费到的时间戳是故障恢复时数据不丢不重的基础。4.4 导入与压缩任务import / compaction / slotimport: 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 或 vortexdeltalog选择删除日志格式直接决定存储层文件的组织方式与下游读取能力。此外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.goMessage stream消息接收与发布订阅前缀dataNodeSubNamePrefix: dataNodetime tick 间隔 500msconfigs/milvus.yaml、configs/milvus.yamlRoot Coordinator获取最新唯一 IDDataNode持有 RootCoord gRPC 客户端ID 由internal/allocator批量拉取缓存internal/datanode/data_node.go、internal/allocatorData Coordinatorflush 信息与订阅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),仅供参考
返回列表