news 2026/9/10 11:59:18

Milvus CDC 双集群同步测试指南:架构、配置与全量操作验证实践

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Milvus CDC 双集群同步测试指南:架构、配置与全量操作验证实践

Milvus CDC 双集群同步测试指南:架构、配置与全量操作验证实践

【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus

本篇技术指南以 Milvus 开源仓库中 tests/python_client/cdc/README.md 为核心,系统讲解 Milvus CDC(Change Data Capture,变更数据捕获)双集群同步测试套件的设计思路、运行方式与验证原理。该测试套件用于验证上游(upstream,源)集群上的各类操作能否被正确复制到下游(downstream,目标)集群,是 Milvus 集群级数据复制与容灾能力的关键质量保障。读完本文,你将掌握如何搭建 CDC 测试拓扑、按操作类别精准执行用例、自定义连接与同步参数,并理解其"查询驱动"的一致性校验机制。

一、套件定位:它验证什么

该测试套件由位于 tests/python_client/cdc/ 目录下的一组 pytest 用例、会话级 Fixture 与辅助脚本组成,核心使命是:对上游集群执行操作,然后通过查询下游集群,验证变更是否按预期同步到位。它覆盖的同步对象包括:

  • 数据库操作:库的创建、删除、属性变更(ALTER / DROP DATABASE_PROPERTIES);
  • Collection DDL 与管理:建/删/改名集合,以及加载、释放、flush、compact 等生命周期管理;
  • 分区操作:分区创建/删除、加载/释放、分区内数据写入与删除;
  • 数据操作:insert、delete、upsert、bulk insert 等 DML;
  • 索引操作:索引创建与删除;
  • 别名操作:create / drop / alter alias;
  • RBAC 操作:用户、角色、权限的授予与回收;
  • 资源组操作:资源组的创建、删除与节点/副本迁移。

套件在启动时会自动完成 CDC 拓扑搭建(无需手动干预),运行过程中采用基于查询的一致性校验(Query-based verification)来确认数据同步状态,并支持可配置的同步超时与进度日志输出。

二、运行前置条件

在运行测试前,需要准备:

  1. 两个正在运行的 Milvus 实例:上游集群(数据源)与下游集群(数据目标);
  2. Python 依赖,版本要求pymilvus>=2.6.0
    pip install pymilvus>=2.6.0 pytest numpy
  3. 网络连通性:测试环境必须能访问两个集群,且具备两个集群的鉴权凭据。

需要说明的是,测试默认配置的地址与凭据(见下文)是仓库开发环境的内网地址,实际运行时请务必通过命令行参数覆盖为你自己的集群地址。

三、快速开始与拓扑自动搭建

最简单的运行方式是在 CDC 目录下直接执行:

cd /path/to/milvus/tests/python_client/cdc # 以默认配置运行全部用例 pytest testcases/

默认配置下,测试框架在会话启动时自动配置 CDC 拓扑:

  • 上游 URI:http://10.104.17.154:19530
  • 下游 URI:http://10.104.17.156:19530
  • 鉴权:root:Milvus

拓扑是如何自动搭建的

从 conftest.py 可以看到,会话级自动 Fixturecdc_topology_setup在测试会话开始时就执行了以下三步:

  1. 构建集群配置:为源集群与目标集群分别生成包含cluster_idconnection_param(uri 与 token)以及pchannels列表的配置对象。物理通道(pchannel)命名遵循{cluster_id}-rootcoord-dml_{i}的规范,数量由--pchannel-num控制,默认 16 个;
  2. 建立单向复制拓扑:通过cross_cluster_topology声明source_cluster_id -> target_cluster_id的单向复制关系,即上游 → 下游;
  3. 初始化 CDC 连接:调用 pymilvus 的update_replicate_configurationRPC 将配置同时下发到两个集群,随后等待 5 秒让 CDC 完成初始化。

值得注意的是,该 Fixture 使用独立的短生命周期客户端up_tmp/dn_tmp)来执行控制面 RPC,避免与后续 DML 共用 gRPC 通道而产生通道关闭竞态;同时配置下发采用并发扇出apply_replicate_configuration,见 conftest.py),这是因为服务端在waitUntilPrimaryChangeOrConfigurationSame中会阻塞非主集群,顺序调用可能因首个客户端恰好是副本而触发 RPC 超时死锁。这套拓扑搭建对测试用例是透明的——用例直接使用已配置好的upstream_client/downstream_clientFixture 即可。

四、测试分类与用例组织

测试按操作类别组织为多个独立的测试文件,每个文件对应一个TestCDCSync*测试类,全部继承自 base.py 中的TestCDCSyncBase,并以@pytest.mark.tags(CaseLabel.CDC)标记便于过滤。

1. 数据库操作 ——test_database.py

对应类TestCDCSyncDatabase,覆盖:

  • CREATE_DATABASE(test_create_database)
  • DROP_DATABASE
  • ALTER_DATABASE_PROPERTIES:设置database.max.collectionsdatabase.diskQuota.mb等属性并验证下游describe_database结果一致
  • DROP_DATABASE_PROPERTIES:删除指定属性键后,验证下游对应键消失、未删除的键值保持不变

每个用例遵循"上游执行 → 断言上游生效 → 轮询下游直到一致"的三段式模式。

2. 资源组操作 ——test_resource_group.py

覆盖CREATE_RESOURCE_GROUPDROP_RESOURCE_GROUPTRANSFER_NODETRANSFER_REPLICA。从该文件的用例命名(如test_create_resource_group_not_replicated)可以推断:资源组相关操作预期不会被复制到下游,这体现了 CDC 对不同元数据类型有选择性的复制策略。

3. RBAC 操作 ——test_rbac.py

覆盖:

  • CREATE_ROLE/DROP_ROLE
  • CREATE_USER/DROP_USER
  • GRANT_ROLE/REVOKE_ROLE
  • GRANT_PRIVILEGE/REVOKE_PRIVILEGE

在源码中,RBAC 用例还进一步扩展了密码更新(test_update_password)、权限组(privilege group)的创建/删除、v2 版授权 API 以及组内权限的增删等场景,验证鉴权体系变更也能随 CDC 同步。

4. Collection DDL 操作 ——test_collection.py

覆盖CREATE_COLLECTIONDROP_COLLECTIONRENAME_COLLECTION。其中创建用例(test_create_collection)会先在上游创建带默认 schema 的集合,断言上游has_collection为真,再轮询下游确认同步。

5. 索引操作 ——test_index.py

覆盖CREATE_INDEXDROP_INDEX,源码中还包含综合向量索引(FLOAT / FLOAT16 / BINARY / SPARSE)、综合标量索引以及 BFLOAT16 / INT8 新向量类型的索引创建同步验证。

6. 数据操作 ——test_dml.py

覆盖:

  • INSERT:插入后 flush,并在下游用count(*)聚合查询等待记录数达标
  • DELETE:先查询上游真实 ID,再按id in [...]过滤删除,验证下游剩余计数与删除记录不可见
  • UPSERT:同时"更新既有主键 + 插入新主键",在下游以 Strong 一致性级别验证总数与更新/新增记录的过滤命中数
  • BULK_INSERT:批量导入场景

DML 用例还覆盖了含多种数据类型的综合 schema(FLOAT_VECTOR / FLOAT16_VECTOR / BINARY_VECTOR / SPARSE_FLOAT_VECTOR 等最多 4 个向量字段,加标量、数组、JSON 字段),以及 auto_id 情况下上下游主键一致性(见 test_insert_auto_id_consistency)。

7. Collection 管理操作 ——test_collection.py中的TestCDCSyncCollectionManagement

覆盖LOAD_COLLECTIONRELEASE_COLLECTIONFLUSHCOMPACT,验证加载/释放状态与数据落盘、压缩动作的同步。

8. 别名操作 ——test_alias.py

覆盖CREATE_ALIASDROP_ALIASALTER_ALIAS(test_create_alias 等)。

9. 分区操作 ——test_partition.py

覆盖:

  • CREATE_PARTITION/DROP_PARTITION(test_create_partition)
  • LOAD_PARTITION/RELEASE_PARTITION
  • 分区数据操作(INSERT、DELETE,见test_partition_inserttest_partition_delete

10. 源码中扩展的更多场景

除 README 列举的 9 类外,目录中还包含面向进阶场景的用例文件,可配合 CI/混沌演练使用:

  • test_setup_cdc.py:拓扑搭建本身的自检
  • test_switchover.py:拓扑主备切换(switchover)与故障切换(failover)下的同步行为
  • test_force_promote.py/test_force_promote_cleanup.py:强制提升场景
  • test_fts_and_text.py:全文检索(BM25 FTS)与文本匹配的同步
  • test_multi_database.py:多数据库场景
  • test_schema_features.py:动态字段、可空字段、默认值、分区键、聚集键等 schema 特性
  • test_import_2pc.py:两阶段导入(Import 2PC)同步
  • test_search_verification.py:搜索结果一致性验证
  • test_collection_properties.py:集合属性同步

五、配置参数详解

所有命令行参数均在 conftest.py 的pytest_addoption中注册,通过 pytest 的--key value形式传入。

连接参数

参数说明默认值
--upstream-uri上游 Milvus URIhttp://10.104.17.154:19530
--upstream-token上游鉴权 tokenroot:Milvus
--downstream-uri下游 Milvus URIhttp://10.104.17.156:19530
--downstream-token下游鉴权 tokenroot:Milvus

对应地,conftest.py 提供了upstream_clientdownstream_clientupstream_uridownstream_token会话级 Fixture,其中两个 client Fixture 会在会话结束时自动close()

CDC 拓扑参数

参数说明默认值
--source-cluster-id源集群标识符cdc-test-source-0930
--target-cluster-id目标集群标识符cdc-test-target-0930
--pchannel-num物理通道(pchannel)数量16

--pchannel-num直接决定每个集群配置中生成的rootcoord-dml_{i}通道数量,即 CDC 复制使用的物理通道宽度;--source-cluster-id--target-cluster-id同时作为通道命名的前缀,因此自定义 cluster-id 时通道名会随之变化,需确保与集群实际配置一致。

测试参数

参数说明默认值
--sync-timeout同步等待超时(秒)30

sync_timeout会注入到每个用例的等待轮询逻辑中(见下文wait_for_sync),用于控制"上游操作后、下游未达预期状态"时的最长等待时间。

其他扩展参数(稳定性/混沌场景)

面向稳定性与故障注入场景,conftest 还注册了以下参数:--request-duration(操作持续时间,如30m1h60s)、--is-check(是否对 checker 统计做断言)、--milvus-ns(Milvus 部署的 Kubernetes 命名空间,默认chaos-testing),以及 Import 2PC 相关的--import-2pc-workload--import-2pc-minio-host--import-2pc-minio-bucket--import-2pc-downstream-minio-host--import-2pc-downstream-minio-bucket--import-2pc-rows(默认 20 行/次)。

六、使用示例

运行指定类别的用例

# 数据库操作用例 pytest testcases/test_database.py # RBAC 操作用例 pytest testcases/test_rbac.py # 数据操作用例 pytest testcases/test_dml.py

自定义连接配置

pytest testcases/ \ --upstream-uri http://localhost:19530 \ --upstream-token root:Milvus \ --downstream-uri http://localhost:19531 \ --downstream-token root:Milvus

自定义同步超时

网络较慢或数据量较大时,适当放大超时:

pytest testcases/test_dml.py --sync-timeout 180

自定义 CDC 拓扑

pytest testcases/ \ --source-cluster-id my-source \ --target-cluster-id my-target \ --pchannel-num 32

全参数自定义

pytest testcases/test_database.py \ --upstream-uri http://10.100.1.10:19530 \ --upstream-token root:Milvus \ --downstream-uri http://10.100.1.20:19530 \ --downstream-token root:Milvus \ --source-cluster-id prod-source \ --target-cluster-id prod-target \ --pchannel-num 32 \ --sync-timeout 180

运行单个测试方法

pytest testcases/test_database.py::TestCDCSyncDatabase::test_create_database \ --upstream-uri http://localhost:19530 \ --downstream-uri http://localhost:19531

七、项目结构

tests/python_client/cdc/ ├── conftest.py # pytest 插件入口:命令行参数、会话级 Fixture、CDC 拓扑自动搭建、切换/混沌辅助 ├── scripts/ │ ├── setup_cdc_topology.py # 独立运行的拓扑搭建脚本(支持多目标与集群下线场景) │ └── diff_upstream_downstream.py # 上下游集群全量元数据/数据对比脚本 ├── stablity/ │ ├── test_single_request_operation.py # 单请求操作稳定性用例 │ └── test_concurrent_operation.py # 并发操作稳定性用例 └── testcases/ ├── base.py # 测试基类与工具函数(命名、等待同步、schema 工厂、数据生成、验证助手) ├── test_database.py # 数据库操作用例 ├── test_rbac.py # RBAC 操作用例 ├── test_collection.py # Collection DDL 与集合管理用例 ├── test_index.py # 索引操作用例 ├── test_dml.py # 数据操作用例 ├── test_collection_management.py # 集合管理用例(源码中对应 TestCDCSyncCollectionManagement 位于 test_collection.py) ├── test_alias.py # 别名操作用例 ├── test_partition.py # 分区操作用例 ├── test_resource_group.py # 资源组操作用例 ├── test_switchover.py # 拓扑切换/故障切换用例 ├── test_force_promote.py # 强制提升用例 ├── test_fts_and_text.py # 全文检索与文本匹配用例 ├── test_import_2pc.py # Import 2PC 用例 ├── test_multi_database.py # 多数据库用例 ├── test_schema_features.py # schema 特性用例 ├── test_search_verification.py # 搜索一致性用例 └── test_setup_cdc.py # 拓扑搭建自检用例

八、源码级解析:同步等待与一致性验证机制

wait_for_sync:带进度日志的轮询等待

所有用例的"等待同步"都复用了基类中的静态方法 wait_for_sync。它接受一个返回布尔值的检查函数check_func、超时时间与操作名,以2 秒为间隔轮询

  • 每轮执行check_func(),返回 True 即记录[SUCCESS] {operation} synced successfully in {elapsed:.2f}s
  • 每 10 秒或首次检查时输出进度百分比[WAITING] ... (xx.x% of timeout)
  • 检查函数内部抛出的异常会被捕获并记录为 warning 后继续重试(避免查询尚未就绪的下游导致误判);
  • 超时未达成则记录[FAILED]并返回 False,由用例中的assert决定失败。

查询驱动的四类验证助手

基类提供了四种可复用的数据一致性验证方法,体现了"Query-based verification"的落地方式:

  1. verify_data_sampling:先从上游拉取全部主键(filter=""+limit=16384),按sample_ratio(默认 0.2)随机抽样,再对每个样本主键分别在上游、下游执行id == {pk}点查,逐字段比较。浮点字段采用1e-6容差,最终返回匹配数、不匹配数与差异明细;
  2. verify_search_consistency:对同一批查询向量分别在上、下游执行 ANN 搜索,计算每次查询返回主键集合的 Jaccard 重叠率并求平均,用于评估复制后的检索结果一致性
  3. verify_query_consistency:以相同 filter 表达式查询两端,对比主键集合的重叠数与各自独有主键;
  4. verify_iterator_consistency:通过query_iteratorbatch_size=100全量遍历两端主键,比较集合是否完全相等。

辅助脚本:手动对账与拓扑维护

  • scripts/setup_cdc_topology.py:可独立运行的拓扑配置脚本,除搭建单向复制外,还支持多目标集群(逗号分隔的 URI/ID 列表)与集群下线场景——对要移除的集群下发cross_cluster_topology: []的空拓扑,使其从复制关系中剥离。所有客户端在同一线程池内并发执行update_replicate_configuration
  • scripts/diff_upstream_downstream.py:深度对账脚本。逐库、逐集合采集两端的num_entities、schema 字段数、索引名集合、分区名集合、副本数与count(*),用deepdiff比较;在排除num_entities这一瞬态字段后可反复轮询(间隔 60 秒,最多 10 次)直到无差异,用于长时间同步后的全量一致性确认。

稳定性与故障注入支撑

stablity/目录下的稳定性用例通过chaos.checker中封装的各操作 Checker(CollectionCreateChecker、InsertChecker、SearchChecker、FullTextSearchChecker、Import2PCChecker 等)对操作结果与最终一致性进行统计断言;conftest.py 中的kubectl_helperFixture 借助 Chaos Mesh 的PodChaos(action 为container-kill)对指定 instance 的容器执行不删除 Pod 对象的容器级故障注入,并轮询确认所有容器 restartCount 递增、Pod UID 保持不变,随后kubectl wait等待 Pod Ready。配合switchover_helper(conftest.py)可完成主备方向互换、并在故障窗口内重试直至拓扑恢复,从而验证 CDC 在真实故障下的收敛能力。

九、实践建议

  1. 务必覆盖默认连接参数:仓库默认的10.104.17.154/156是开发环境地址,任何实际运行都应显式传入--upstream-uri/--downstream-uri/ token;
  2. 根据数据量与网络带宽调整--sync-timeout:默认 30 秒适合小样本用例,大数据量 DML 建议提升至 120–180 秒,避免轮询超时导致的误报;
  3. 善用按文件/按方法粒度执行:排查具体同步问题时,优先运行单个测试方法并结合其进度日志([WAITING]/[SYNC_OK]/[VERIFY])定位是"上游未生效"还是"下游未同步";
  4. diff_upstream_downstream.py做最终对账:在用例跑完后,用该脚本做一次全库级深度对比,作为 CDC 收敛性的最终确认;
  5. 混沌验证前先确认拓扑就绪:涉及容器杀死的故障注入需要kubectl与 Chaos Mesh 环境(--milvus-ns指定命名空间),且应保证cdc_topology_setup已成功完成,否则故障注入结果无法归因于复制链路。

通过本文,你可以完整掌握该 CDC 测试套件的拓扑搭建原理、参数体系、用例组织与验证机制,并能够直接复用它来验证自己部署的双集群 Milvus 数据复制链路。

【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/10 11:56:39

对数几率回归从二分类到多分类:西瓜与鸢尾花数据集的Python实战

简介:基于对数几率回归模型实现西瓜与鸢尾花分类识别的期末大作业资料包,面向计算机、数据科学、人工智能等专业学生及教师,覆盖课程设计、期末大作业与毕设参考场景。压缩包共30个文件、约544KB,主要包含Python源码(.py)、Jupyte…

作者头像 李华
网站建设 2026/9/10 11:54:30

InsightFace ArcFace-Paddle 基础训练预测功能测试(TIPC)完整指南

InsightFace ArcFace-Paddle 基础训练预测功能测试(TIPC)完整指南 【免费下载链接】insightface State-of-the-art 2D and 3D Face Analysis Project 项目地址: https://gitcode.com/GitHub_Trending/in/insightface 本文是 InsightFace 仓库中 A…

作者头像 李华
网站建设 2026/9/10 11:53:37

燃气轮机动态建模与Simulink仿真实践

1. 项目概述:回热燃气轮机动态建模的核心价值燃气轮机作为能源动力领域的核心装备,其动态特性研究一直是工程师关注的焦点。这个基于Matlab/Simulink 2021构建的回热燃气轮机动态模型,本质上是一个能够模拟启动、停机和变工况过程的部件级仿真…

作者头像 李华
网站建设 2026/9/10 11:52:55

RP2040低功耗实战:手撕寄存器实现<30μA深度休眠

1. 为什么你写的低功耗代码总“省不下电”?——从RP2040的寄存器真相说起我第一次在Pico上跑低功耗demo时,用官方SDK调了个sleep()函数,万用表一测:电流从8mA掉到7.2mA。心里咯噔一下——这哪是休眠,这是打盹儿。后来拆…

作者头像 李华