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 CDCChange Data Capture变更数据捕获双集群同步测试套件的设计思路、运行方式与验证原理。该测试套件用于验证上游upstream源集群上的各类操作能否被正确复制到下游downstream目标集群是 Milvus 集群级数据复制与容灾能力的关键质量保障。读完本文你将掌握如何搭建 CDC 测试拓扑、按操作类别精准执行用例、自定义连接与同步参数并理解其查询驱动的一致性校验机制。一、套件定位它验证什么该测试套件由位于 tests/python_client/cdc/ 目录下的一组 pytest 用例、会话级 Fixture 与辅助脚本组成核心使命是对上游集群执行操作然后通过查询下游集群验证变更是否按预期同步到位。它覆盖的同步对象包括数据库操作库的创建、删除、属性变更ALTER / DROP DATABASE_PROPERTIESCollection DDL 与管理建/删/改名集合以及加载、释放、flush、compact 等生命周期管理分区操作分区创建/删除、加载/释放、分区内数据写入与删除数据操作insert、delete、upsert、bulk insert 等 DML索引操作索引创建与删除别名操作create / drop / alter aliasRBAC 操作用户、角色、权限的授予与回收资源组操作资源组的创建、删除与节点/副本迁移。套件在启动时会自动完成 CDC 拓扑搭建无需手动干预运行过程中采用基于查询的一致性校验Query-based verification来确认数据同步状态并支持可配置的同步超时与进度日志输出。二、运行前置条件在运行测试前需要准备两个正在运行的 Milvus 实例上游集群数据源与下游集群数据目标Python 依赖版本要求pymilvus2.6.0pip install pymilvus2.6.0 pytest numpy网络连通性测试环境必须能访问两个集群且具备两个集群的鉴权凭据。需要说明的是测试默认配置的地址与凭据见下文是仓库开发环境的内网地址实际运行时请务必通过命令行参数覆盖为你自己的集群地址。三、快速开始与拓扑自动搭建最简单的运行方式是在 CDC 目录下直接执行cd /path/to/milvus/tests/python_client/cdc # 以默认配置运行全部用例 pytest testcases/默认配置下测试框架在会话启动时自动配置 CDC 拓扑上游 URIhttp://10.104.17.154:19530下游 URIhttp://10.104.17.156:19530鉴权root:Milvus拓扑是如何自动搭建的从 conftest.py 可以看到会话级自动 Fixturecdc_topology_setup在测试会话开始时就执行了以下三步构建集群配置为源集群与目标集群分别生成包含cluster_id、connection_paramuri 与 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_DATABASEtest_create_databaseDROP_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 用例还覆盖了含多种数据类型的综合 schemaFLOAT_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_ALIAStest_create_alias 等。9. 分区操作 ——test_partition.py覆盖CREATE_PARTITION/DROP_PARTITIONtest_create_partitionLOAD_PARTITION/RELEASE_PARTITION分区数据操作INSERT、DELETE见test_partition_insert、test_partition_delete10. 源码中扩展的更多场景除 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_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同步等待超时秒30sync_timeout会注入到每个用例的等待轮询逻辑中见下文wait_for_sync用于控制上游操作后、下游未达预期状态时的最长等待时间。其他扩展参数稳定性/混沌场景面向稳定性与故障注入场景conftest 还注册了以下参数--request-duration操作持续时间如30m、1h、60s、--is-check是否对 checker 统计做断言、--milvus-nsMilvus 部署的 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先从上游拉取全部主键filterlimit16384按sample_ratio默认 0.2随机抽样再对每个样本主键分别在上游、下游执行id {pk}点查逐字段比较。浮点字段采用1e-6容差最终返回匹配数、不匹配数与差异明细verify_search_consistency对同一批查询向量分别在上、下游执行 ANN 搜索计算每次查询返回主键集合的 Jaccard 重叠率并求平均用于评估复制后的检索结果一致性verify_query_consistency以相同 filter 表达式查询两端对比主键集合的重叠数与各自独有主键verify_iterator_consistency通过query_iterator以batch_size100全量遍历两端主键比较集合是否完全相等。辅助脚本手动对账与拓扑维护scripts/setup_cdc_topology.py可独立运行的拓扑配置脚本除搭建单向复制外还支持多目标集群逗号分隔的 URI/ID 列表与集群下线场景——对要移除的集群下发cross_cluster_topology: []的空拓扑使其从复制关系中剥离。所有客户端在同一线程池内并发执行update_replicate_configurationscripts/diff_upstream_downstream.py深度对账脚本。逐库、逐集合采集两端的num_entities、schema 字段数、索引名集合、分区名集合、副本数与count(*)用deepdiff比较在排除num_entities这一瞬态字段后可反复轮询间隔 60 秒最多 10 次直到无差异用于长时间同步后的全量一致性确认。稳定性与故障注入支撑stablity/目录下的稳定性用例通过chaos.checker中封装的各操作 CheckerCollectionCreateChecker、InsertChecker、SearchChecker、FullTextSearchChecker、Import2PCChecker 等对操作结果与最终一致性进行统计断言conftest.py 中的kubectl_helperFixture 借助 Chaos Mesh 的PodChaosaction 为container-kill对指定 instance 的容器执行不删除 Pod 对象的容器级故障注入并轮询确认所有容器 restartCount 递增、Pod UID 保持不变随后kubectl wait等待 Pod Ready。配合switchover_helperconftest.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),仅供参考