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)来确认数据同步状态,并支持可配置的同步超时与进度日志输出。
二、运行前置条件
在运行测试前,需要准备:
- 两个正在运行的 Milvus 实例:上游集群(数据源)与下游集群(数据目标);
- Python 依赖,版本要求
pymilvus>=2.6.0:pip install pymilvus>=2.6.0 pytest numpy - 网络连通性:测试环境必须能访问两个集群,且具备两个集群的鉴权凭据。
需要说明的是,测试默认配置的地址与凭据(见下文)是仓库开发环境的内网地址,实际运行时请务必通过命令行参数覆盖为你自己的集群地址。
三、快速开始与拓扑自动搭建
最简单的运行方式是在 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在测试会话开始时就执行了以下三步:
- 构建集群配置:为源集群与目标集群分别生成包含
cluster_id、connection_param(uri 与 token)以及pchannels列表的配置对象。物理通道(pchannel)命名遵循{cluster_id}-rootcoord-dml_{i}的规范,数量由--pchannel-num控制,默认 16 个; - 建立单向复制拓扑:通过
cross_cluster_topology声明source_cluster_id -> target_cluster_id的单向复制关系,即上游 → 下游; - 初始化 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_DATABASEALTER_DATABASE_PROPERTIES:设置database.max.collections、database.diskQuota.mb等属性并验证下游describe_database结果一致DROP_DATABASE_PROPERTIES:删除指定属性键后,验证下游对应键消失、未删除的键值保持不变
每个用例遵循"上游执行 → 断言上游生效 → 轮询下游直到一致"的三段式模式。
2. 资源组操作 ——test_resource_group.py
覆盖CREATE_RESOURCE_GROUP、DROP_RESOURCE_GROUP、TRANSFER_NODE、TRANSFER_REPLICA。从该文件的用例命名(如test_create_resource_group_not_replicated)可以推断:资源组相关操作预期不会被复制到下游,这体现了 CDC 对不同元数据类型有选择性的复制策略。
3. RBAC 操作 ——test_rbac.py
覆盖:
CREATE_ROLE/DROP_ROLECREATE_USER/DROP_USERGRANT_ROLE/REVOKE_ROLEGRANT_PRIVILEGE/REVOKE_PRIVILEGE
在源码中,RBAC 用例还进一步扩展了密码更新(test_update_password)、权限组(privilege group)的创建/删除、v2 版授权 API 以及组内权限的增删等场景,验证鉴权体系变更也能随 CDC 同步。
4. Collection DDL 操作 ——test_collection.py
覆盖CREATE_COLLECTION、DROP_COLLECTION、RENAME_COLLECTION。其中创建用例(test_create_collection)会先在上游创建带默认 schema 的集合,断言上游has_collection为真,再轮询下游确认同步。
5. 索引操作 ——test_index.py
覆盖CREATE_INDEX、DROP_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_COLLECTION、RELEASE_COLLECTION、FLUSH、COMPACT,验证加载/释放状态与数据落盘、压缩动作的同步。
8. 别名操作 ——test_alias.py
覆盖CREATE_ALIAS、DROP_ALIAS、ALTER_ALIAS(test_create_alias 等)。
9. 分区操作 ——test_partition.py
覆盖:
CREATE_PARTITION/DROP_PARTITION(test_create_partition)LOAD_PARTITION/RELEASE_PARTITION- 分区数据操作(INSERT、DELETE,见
test_partition_insert、test_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 URI | http://10.104.17.154:19530 |
--upstream-token | 上游鉴权 token | root:Milvus |
--downstream-uri | 下游 Milvus URI | http://10.104.17.156:19530 |
--downstream-token | 下游鉴权 token | root:Milvus |
对应地,conftest.py 提供了upstream_client、downstream_client、upstream_uri、downstream_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(操作持续时间,如30m、1h、60s)、--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"的落地方式:
- verify_data_sampling:先从上游拉取全部主键(
filter=""+limit=16384),按sample_ratio(默认 0.2)随机抽样,再对每个样本主键分别在上游、下游执行id == {pk}点查,逐字段比较。浮点字段采用1e-6容差,最终返回匹配数、不匹配数与差异明细; - verify_search_consistency:对同一批查询向量分别在上、下游执行 ANN 搜索,计算每次查询返回主键集合的 Jaccard 重叠率并求平均,用于评估复制后的检索结果一致性;
- verify_query_consistency:以相同 filter 表达式查询两端,对比主键集合的重叠数与各自独有主键;
- verify_iterator_consistency:通过
query_iterator以batch_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 在真实故障下的收敛能力。
九、实践建议
- 务必覆盖默认连接参数:仓库默认的
10.104.17.154/156是开发环境地址,任何实际运行都应显式传入--upstream-uri/--downstream-uri/ token; - 根据数据量与网络带宽调整
--sync-timeout:默认 30 秒适合小样本用例,大数据量 DML 建议提升至 120–180 秒,避免轮询超时导致的误报; - 善用按文件/按方法粒度执行:排查具体同步问题时,优先运行单个测试方法并结合其进度日志(
[WAITING]/[SYNC_OK]/[VERIFY])定位是"上游未生效"还是"下游未同步"; - 用
diff_upstream_downstream.py做最终对账:在用例跑完后,用该脚本做一次全库级深度对比,作为 CDC 收敛性的最终确认; - 混沌验证前先确认拓扑就绪:涉及容器杀死的故障注入需要
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),仅供参考