Redpanda Connect Unified Migrator 深度指南:Kafka 与 Redpanda 集群间 Topic、Schema Registry 与消费组的一体化迁移
【免费下载链接】connectFancy stream processing made operationally mundane项目地址: https://gitcode.com/GitHub_Trending/con/connect
导读
redpanda_migrator是 Redpanda Connect(本仓库GitHub_Trending/con/connect)内置的一套统一数据迁移系统,用于在 Apache Kafka 与 Redpanda 集群之间进行完整、可运维的迁移。它以一对输入/输出组件的形式协同工作,自动完成 Topic 创建与配置同步、Schema Registry 的 Subject/Schema/兼容性迁移,以及消费组偏移量的时间戳关联翻译与提交。读完本文,你将掌握该组件的整体架构、核心工作流程、完整配置参数与实战 YAML 示例、指标监控体系、保证与限制,以及源码级实现细节和测试组织方式,能够独立搭建一条从源集群到目标集群的生产级迁移管道。
本文主体基于仓库文档 internal/impl/redpanda/migrator/README.md,并结合 migrator.go、migrator_topic.go、migrator_schema_registry.go、migrator_groups.go 等源码与测试展开说明。
一、架构总览:三个专职迁移器协同
统一迁移器(Unified Migrator)由一个中央协调器Migrator和三个专职子迁移器组成,三者协同完成集群到集群的完整迁移:
Migrator—— 中央协调器,负责管理输入/输出生命周期、将服务消息转换为 franz-go 记录、协调子迁移器的执行时机、处理来源头(provenance headers)与 Schema ID 翻译。topicMigrator—— Topic 基础设施迁移,负责目标 Topic 名称解析(支持插值)、按镜像分区数创建 Topic、复制受支持的配置键、可选地执行带安全转换的 ACL 复制。schemaRegistryMigrator—— Schema 同步,负责按正则模式列出与过滤 Subject、以 ID 翻译或固定 ID 方式复制 Schema、传播各 Subject 的兼容性设置、执行一次性或周期性同步循环。groupsMigrator—— 消费组偏移量翻译,负责按名称与状态过滤发现消费组、使用时间戳关联翻译偏移量、借助内嵌偏移头精化翻译结果、通过缓存防止偏移回退。
从源码结构看(migrator.go),Migrator以组合方式持有topicMigrator、schemaRegistryMigrator、groupsMigrator三个子迁移器,并各自维护自己的缓存结构:knownTopics(Topic 映射)、knownSubjects/knownSchemas(Schema 映射)、commitedOffsets(已提交偏移量)。底层的KadmClient、SrClient、KgoClient均基于 franz-go 库构建。
关键设计点
- 输入组件不做任何同步工作:
redpanda_migrator输入只是从源集群消费消息并向下游转发,Topic/Schema/Group 的所有同步逻辑都位于配对的输出组件中(migrator.go)。 - 输入输出必须配对:每个管道必须同时配置
redpanda_migrator输入与输出;当单个管道中存在多对迁移器时,通过label字段精确匹配输入与输出的对应关系。 - 同一流共享状态:
Migrator通过GetOrSetGeneric按label + stream作用域存储(migrator.go),确保同一管道内的输入与输出共享同一个Migrator实例。
二、记录构建管道:从 service.Message 到 franz-go Record
输入消息在被写入目标集群之前,会经历一次完整的转换流程。核心转换逻辑位于 messageBatchToFranzRecords,每一步都从消息元数据中提取字段并映射到kgo.Record:
| 源元数据 | 目标字段 | 说明 |
|---|---|---|
kafka_key | kgo.Record.Key | 可选,键会原样保留 |
kafka_value | kgo.Record.Value | 必需;若开启 Schema ID 翻译则先解析并改写 ID |
kafka_topic | kgo.Record.Topic | 经插值解析为目标 Topic;不存在则自动创建 |
kafka_partition | kgo.Record.Partition | 分区被保留,保证写入顺序与源一致 |
kafka_timestamp_ms | kgo.Record.Timestamp | 必需,毫秒时间戳转换为time.Time |
kafka_offset | 偏移头 | 消费组迁移开启时写入offset_header(8 字节大端编码) |
kafka_headers | kgo.Record.Headers | 原头部透传 |
四项关键转换
- Schema ID 翻译(Schema ID Translation):当
translate_ids: true时,源码先通过parseSchemaID解析 Confluent 线格式前缀中的 Schema ID,再调用DestinationSchemaID从knownSchemas缓存/目标注册表中查出目标 ID,最后用updateSchemaID改写值前缀(migrator.go)。ID 在同一批次内做了lastSchemaID缓存,避免重复查询。 - Topic 名称解析(Topic Name Resolution):
topic字段为可插值字符串,从kafka_topic元数据解析出目标 Topic 名。 - 偏移头注入(Offset Header Injection):当消费组迁移启用且
offset_header非空时,将源偏移量以 8 字节大端无符号整数编码写入头部,供目标侧做精确消费组翻译。 - 来源追踪(Provenance Tracking):默认头名
redpanda-migrator-provenance,携带源集群 ID,用于双向迁移时防止消息回流(详见“双向迁移”一节)。
来源头保护逻辑
源码在构造记录时有一层防数据损坏校验(migrator.go):
- 若记录已带来源头但值为空,直接报错(可能的数据损坏);
- 若来源头值等于源集群 ID,报错(说明消息被错误地绕了一圈);
- 若来源头值等于目标集群 ID,则该记录是“回流”消息,注入
kafka.SkipRecord跳过写入; - 若无来源头,则自动追加携带源集群 ID 的来源头。
自定义headers中若与provenance_header或offset_header同名,会被忽略,从而保证迁移关键头永不冲突(migrator.go)。
三、Topic 迁移器:按需创建 + 幂等同步 + 安全 ACL
Topic 同步是“按需(on-demand)”执行的:第一条消息触发初始同步,后续消息遇到新 Topic 时按需创建。源码中SyncOnce仅在knownTopics为空时执行完整同步,之后不再重复(migrator_topic.go)。
创建流程与幂等处理
每次创建走 createTopicLocked:
- 名称解析:用
NameResolver(插值模板)将源 Topic 名转换为目标名;解析时以kafka_topic元数据构造消息(migrator_topic.go)。 - 读取源详情:通过
ListTopics获取分区数、DescribeTopicConfigs获取资源配置。 - 确定副本因子:默认继承源 Topic 的副本数;可配置
topic_replication_factor覆盖;serverless 模式下使用-1(由服务端决定)。 - 过滤配置:仅复制受支持的配置键(见下)。
- 创建或校验:
- 若
TopicAlreadyExists,则检查分区数:源分区数 > 目标分区数时调用CreatePartitions扩容并更新映射;目标分区数 > 源分区数时记录告警并沿用目标分区数(migrator_topic.go)。 - 创建成功则记录指标并写入
knownTopics缓存。
- 若
受支持的配置键子集
supportedTopicConfigs()(migrator_topic.go)决定哪些 Topic 配置会被复制:
- 普通模式:
cleanup.policy、flush.bytes、flush.ms、initial.retention.local.target.ms、retention.bytes、retention.ms、segment.ms、segment.bytes、compression.type、message.timestamp.type、max.message.bytes。 - Serverless 模式:收窄为
cleanup.policy、retention.ms、max.message.bytes、write.caching。
ACL 安全转换
开启sync_topic_acls: true后,迁移器执行与 Kafka MirrorMaker 2 一致的“只读安全转换”(migrator_topic.go):
- 排除
ALLOW WRITE条目(防止目标侧被写入权限污染); - 将
ALLOW ALL降级为ALLOW READ; - 保留资源模式类型(Pattern)与主机(Host)过滤条件;
- 当源集群安全功能未启用(
SecurityDisabled)时,跳过 ACL 同步并记录告警而非失败。
Topic 同步特性小结
- 按需执行:首条消息触发初始同步,后续按需创建。
- 幂等:已存在 Topic 会被校验,分区不足时自动扩容。
- 配置过滤:仅复制受支持键(serverless 感知的子集)。
- 周期补充:
sync_topic_interval(默认5m)控制周期同步,用于覆盖无消息流量的空 Topic 与迁移后新增的 Topic;设为0s则禁用周期同步(仍会在首条消息时创建)。相关逻辑见 SyncLoop。
四、Schema Registry 迁移器:版本选择、ID 翻译与兼容性传播
Schema 同步在输出连接时执行一次初始同步,并由schema_registry.interval(默认5m)控制后台周期循环;设为0s则仅在连接时同步一次(migrator.go)。
主体流程
- 校验目标模式:目标 Schema Registry 必须处于
READWRITE或IMPORT模式(migrator_schema_registry.go);源与目标 URL 必须不同。 - 列出 Subject:调用
Subjects(include_deleted时带ShowDeleted参数),再按 include/exclude 正则过滤,并对列表做随机洗牌以分散负载(migrator_schema_registry.go)。 - 选择版本:
versions: latest仅取最新版本;versions: all(默认)遍历全部版本并按 ID 升序处理。同步采用 DFS 遍历,会递归纳入依赖引用(references)的 Subject/版本,保证 Avro/Protobuf 引用链完整(migrator_schema_registry.go)。 - Subject 重命名:
subject字段支持插值,可用metadata("schema_registry_subject")与metadata("schema_registry_version")构造目标 Subject 名。 - 两种 ID 处理模式(migrator_schema_registry.go):
translate_ids: true:调用CreateSchema走create-or-reuse语义,目标注册表分配新 ID,写入时消息里的 ID 被翻译为新的目标 ID。- 固定 ID:调用
CreateSchemaWithIDAndVersion保留源 ID 与版本;若遇 ID 冲突,源码会回查SchemaByID确认 Schema 内容是否一致,一致则复用现有 Schema(同时给出“可尝试启用 translate_ids”的提示)。
- 兼容性传播:仅当源端某 Subject显式设置了兼容级别时才在目标端设置;若源端用的是全局兼容模式,则不强制写目标端全局模式(migrator_schema_registry.go)。
- Serverless 特殊处理:serverless 模式下剥离 Schema 元数据(
SchemaMetadata)与规则集(SchemaRuleSet);若目标全局模式非 IMPORT,则通过importModeManager为每个 Subject 临时切换到 IMPORT 模式,同步完成后再恢复原模式(migrator_schema_registry.go)。
未知 Schema ID 的处理
写入时若遇到尚未同步到目标端的 Schema ID,不会触发按需重同步:默认(strict: false)直接透传原值;开启strict: true则报错。注意:0 字节前缀的消息(如 Protobuf)无法与 Schema Registry 头区分,strict 模式下可能误报。此参数仅在translate_ids: true时有效——输出端 lint 规则会校验这一点(migrator.go)。
Schema 同步特性小结
- 连接时初始同步一次 + 可选周期循环;
- 未知 Schema 透传或按
strict报错; - ID 翻译(create-or-reuse)与固定 ID 两种模式;
- 仅显式设置的兼容级别会被传播;
- 并发控制由
max_parallel_http_requests(默认 10)决定。
五、消费组迁移器:基于时间戳的偏移量翻译与精确精化
消费组偏移量同步由consumer_groups.interval(默认1m)控制的后台循环执行,且只有在 Topic 同步完成之后才启动(因为翻译依赖目标端 Topic 已存在)。
过滤规则
listGroupsOffsets 依次应用三层过滤:
- 名称过滤:include/exclude 正则;
- 状态过滤:
only_empty: true时仅迁移Empty状态(无活跃成员但保留元数据);默认only_empty: false时迁移除Dead外的所有状态; - Topic 过滤:没有任何 Topic 有已提交偏移量的组被剔除;
- 另外始终跳过迁移器自己的消费组(源码从输入配置读取
consumer_group并记录为SkipSourceGroup)。
偏移量翻译算法(5 步)
- 读取前一条记录:从源集群读取
offset - 1处的记录,取其时间戳(translateOffset)。 - 近似翻译:用
ListOffsetsAfterMilli在目标端查找该时间戳之后的第一个偏移量o1;若返回时间戳恰好等于请求时间戳,则o1 += 1以获得正确的翻译结果。 - 精确精化:仅对
Empty状态组且配置了offset_header时执行 tryFindExactOffset:读取目标端o1处的记录,解码其内嵌的源偏移头。 - 增量调整:计算
delta = 源偏移 - 内嵌偏移,令o1 += delta后重试。 - 收敛:最多尝试 5 次,直到找到精确偏移(
delta == 0)、命中目标端结束偏移eo、或超出边界报错。
不回退保证与并行处理
- 提交前会读取目标端当前已提交偏移量,仅当翻译结果大于当前值时提交(migrator_groups.go),配合
commitedOffsets缓存(记录[源偏移, 目标偏移]对),保证偏移量永不回退。 - 翻译阶段按分区并行、提交阶段按组并行,指标按
group维度打点。
消费组同步特性小结
- 周期执行,由
interval控制; - 基于状态的过滤(默认排除 Dead,可选仅 Empty);
- 基于时间戳的近似翻译(
ListOffsetsAfterMilli); - 内嵌偏移头的精确精化;
- 缓存保证不回退;
- 翻译与提交均按组并行。
六、执行模型:启动序列、消息处理与后台任务
启动序列
- 输入连接:获取源集群元数据并初始化 admin 客户端,同时读取源集群 ID(
onInputConnected,migrator.go)。 - 输出连接:获取目标集群元数据、admin 客户端与集群 ID;启动 Topic 周期同步循环;执行一次 Schema 注册表同步并启动其周期循环;最后启动消费组同步循环(
onOutputConnected,migrator.go)。 - 初始 Schema 同步:一次性同步。
- 后台循环:Schema 循环(
interval > 0时)与消费组循环(interval > 0时)并发运行,与消息处理互不阻塞。
消息处理
- 首条消息触发 Topic 同步:所有被消费的 Topic 按需创建;
- 逐消息操作:按需创建 Topic、开启时翻译 Schema ID;
- 批量写入:转换后的记录以保留分区的方式写入目标端。
并发模型(migrator.go)
- 消息处理:
max_in_flight = 1(单批在途)以保证顺序——注意:这是 README 中的表述,实际配置字段默认值为 10,输出端 lint 规则明确禁止设置key、partitioner、partition、timestamp、timestamp_ms等会破坏消费组迁移或分区保序的字段; - 偏移量翻译:单次同步内按分区并行;
- 偏移量提交:单次同步内按组并行;
- 后台循环:Schema 与消费组同步为独立 goroutine。
错误处理策略
- Topic 创建失败:导致消息批次失败,下个批次重试;
- Schema 同步失败:记录日志,下轮同步重试;
- 消费组同步失败:记录日志,下轮同步重试;
- 偏移量翻译失败:跳过该分区,其余分区继续。
七、配置模式与完整示例
下面所有配置均取自 migrator.go 中注册的官方示例与文档示例,可直接复制使用。
7.1 基础迁移(Basic Migration)
input: redpanda_migrator: seed_brokers: ["source:9092"] topics: ["orders", "payments"] consumer_group: "migration" output: redpanda_migrator: seed_brokers: ["destination:9092"] topic: ${! @kafka_topic } # Preserve namestopic默认值即为${! @kafka_topic }(保持源名称),也可显式写出。
7.2 Topic 名称变换(Topic Name Transformation)
output: redpanda_migrator: topic: prod_${! @kafka_topic } # Add prefix7.3 Schema Registry + ID 翻译
output: redpanda_migrator: schema_registry: url: "http://dest-registry:8081" translate_ids: true # Create-or-reuse mode versions: all # Migrate all versionsschema_registry完整字段如下(源码定义见 migrator_schema_registry.go):
| 字段 | 类型 | 默认值 | 说明 |
|---|---|---|---|
url | string | — | 注册表基础 URL,必需 |
timeout | duration | 5s | HTTP 客户端超时 |
tls | object | — | TLS 配置 |
enabled | bool | true | 是否启用 Schema 迁移 |
interval | duration | 5m | 同步周期,0s仅启动时同步一次 |
include | []string | 空 | Subject 包含正则(空则全包含) |
exclude | []string | 空 | Subject 排除正则,优先级高于 include |
subject | string | — | Subject 名称插值模板 |
versions | enum | all | latest或all |
include_deleted | bool | false | 是否包含软删除的 Schema |
translate_ids | bool | false | 是否翻译 Schema ID |
normalize | bool | false | 创建时是否规范化 Schema |
strict | bool | false | 未知 Schema ID 是否报错(仅 translate_ids 时有效) |
max_parallel_http_requests | int | 10 | 并发 HTTP 请求数上限 |
schema_registry块还支持标准 HTTP 请求认证字段(如basic_auth,见 7.6 Serverless 示例)。
7.4 消费组 + 过滤(Consumer Groups with Filtering)
output: redpanda_migrator: consumer_groups: interval: 1m include: ["app-.*"] # Only app- prefixed groups exclude: ["migration"] # Exclude migrator itself only_empty: true # Only Empty state groupsconsumer_groups完整字段如下(源码定义见 migrator_groups.go):
| 字段 | 类型 | 默认值 | 说明 |
|---|---|---|---|
enabled | bool | true | 是否启用消费组迁移 |
interval | duration | 1m | 同步周期,0s禁用 |
fetch_timeout | duration | 10s | 读取记录用于时间戳翻译的最大等待时间,低吞吐集群可调大 |
include | []string | 空 | 组名包含正则 |
exclude | []string | 空 | 组名排除正则,优先级高于 include |
only_empty | bool | false | true仅迁移 Empty 组;false迁移除 Dead 外所有组 |
7.5 Serverless 模式
output: redpanda_migrator: serverless: true # Restrict configs to serverless subset schema_registry: url: "https://serverless.redpanda.com:8081" translate_ids: trueserverless: true会将 Topic 配置与 Schema 功能收窄到 Redpanda Cloud serverless 支持的子集(见上文“受支持的配置键”与“Serverless 特殊处理”)。
7.6 迁移到 Redpanda Serverless 的完整官方示例
input: redpanda_migrator: seed_brokers: ["source-kafka:9092"] regexp_topics_include: - '.' regexp_topics_exclude: - '^_' consumer_group: "migrator_cg" schema_registry: url: "http://source-registry:8081" output: redpanda_migrator: seed_brokers: ["serverless-cluster.redpanda.com:9092"] tls: enabled: true sasl: - mechanism: SCRAM-SHA-256 username: "migrator" password: "migrator" schema_registry: url: "https://serverless-cluster.redpanda.com:8081" basic_auth: enabled: true username: "migrator" password: "migrator" translate_ids: true consumer_groups: exclude: - "migrator_cg" # Exclude the migration consumer group itself serverless: true # Enable serverless mode for restricted configurations7.7 其他常用字段速查
| 字段 | 类型 | 默认值 | 说明 |
|---|---|---|---|
topic_replication_factor | int | 继承源 | 目标 Topic 副本因子,迁移到不同规模集群时很有用 |
sync_topic_interval | duration | 5m | Topic 周期同步间隔,0s禁用(首条消息仍会创建) |
sync_topic_acls | bool | false | 是否同步 Topic ACL(带安全转换) |
headers | map[string]string | — | 追加到迁移记录的自定义头(插值),与 provenance/offset 头同名会被忽略 |
provenance_header | string | redpanda-migrator-provenance | 来源头名,置空则不添加 |
offset_header | string | redpanda-migrator-offset | 偏移头名,置空则禁用精确偏移翻译 |
max_in_flight | int | 10 | 在途批次上限;建议设为并行复制的分区总数 |
高吞吐调优
官方文档给出的调优建议:
- 输入侧:
partition_buffer_bytes: 2MB(提高单分区缓冲)、max_yield_batch_bytes: 1MB(允许产出更大批次); - 输出侧:
max_in_flight设置为并行复制的分区总数(最多可覆盖集群全部分区),高于被消费分区数不会带来收益。
八、指标监控体系
迁移器暴露完整的 Prometheus 指标,源码中分别在 migrator_topic.go、migrator_schema_registry.go 与 migrator_groups.go 中定义。
Topic 迁移指标
redpanda_migrator_topics_created_total(counter)—— 成功创建的 Topic 数;redpanda_migrator_topic_create_errors_total(counter)—— Topic 创建失败数;redpanda_migrator_topic_create_latency_ns(timer)—— Topic 创建延迟。
Schema Registry 迁移指标
redpanda_migrator_sr_schemas_created_total(counter)—— 成功创建的 Schema 数;redpanda_migrator_sr_schema_create_errors_total(counter)—— Schema 创建失败数;redpanda_migrator_sr_schema_create_latency_ns(timer)—— Schema 创建延迟;redpanda_migrator_sr_compatibility_updates_total(counter)—— 兼容性更新次数;redpanda_migrator_sr_compatibility_update_errors_total(counter)—— 兼容性更新失败数;redpanda_migrator_sr_compatibility_update_latency_ns(timer)—— 兼容性更新延迟。
消费组迁移指标(带group标签)
redpanda_migrator_cg_offsets_translated_total(counter)—— 成功翻译的偏移量数;redpanda_migrator_cg_offset_translation_errors_total(counter)—— 偏移量翻译失败数;redpanda_migrator_cg_offset_translation_latency_ns(timer)—— 偏移量翻译延迟;redpanda_migrator_cg_offsets_committed_total(counter)—— 成功提交的偏移量数;redpanda_migrator_cg_offset_commit_errors_total(counter)—— 偏移量提交失败数;redpanda_migrator_cg_offset_commit_latency_ns(timer)—— 偏移量提交延迟。
消费滞后指标(带topic、partition标签)
redpanda_lag(gauge)—— 迁移器输入在每个分区上的当前消费滞后(高水位与当前消费位置的差值),用于观察迁移进度是否跟得上生产速度。
九、保证与限制
保证(Guarantees)
- Topic 分区数:目标 Topic 以匹配的分区数创建;
- 不回退(No offset rewind):消费组偏移量绝不向回移动;
- ACL 安全:排除 WRITE 操作、ALL 降级为 READ;
- 幂等:重复同步安全(Topic、Schema、消费组均是)。
限制(Limitations)
- 偏移量翻译为尽力而为:若前一条记录的时间戳无法读取,或目标端在该时间戳后没有偏移量,则跳过该分区;
- 分区数一致要求:消费组迁移要求源与目标 Topic 分区数一致(不一致的 Topic 会被
filterTopics跳过,migrator_groups.go); - Schema Registry 模式:目标必须处于
READWRITE或IMPORT模式; - 精确偏移依赖:精确翻译依赖目标端记录中的偏移头(迁移器会自动添加);
- 时间戳单调性:近似翻译依赖时间戳单调递增的假设(源码
translateOffset注释明确说明); - 输出端禁用的字段:
key、partitioner、partition、timestamp、timestamp_ms均被 lint 规则禁止,设置会破坏消费组迁移或分区保序。
十、高级特性
双向迁移(Bidirectional Migration)
来源头防止环形迁移:
output: redpanda_migrator: provenance_header: "redpanda-migrator-provenance" # Default带有所属目标集群 ID 来源头的记录会被跳过(见“来源头保护逻辑”),从而实现 A↔B 双向复制而不产生死循环。
ACL 复制(ACL Replication)
output: redpanda_migrator: sync_topic_acls: true- 排除
ALLOW WRITE条目; - 将
ALLOW ALL降级为ALLOW READ; - 保留资源模式类型与主机过滤。
Schema 规范化(Schema Normalization)
output: redpanda_migrator: schema_registry: normalize: true创建 Schema 时规范化格式,保证语义等价但格式不同的 Schema 可被识别为相同(源码通过schemaEquals/schemaStringEquals实现 JSON/Avro 的 JSON 语义比较与 Protobuf 的空白归一比较,见 migrator_schema_registry.go)。
精确偏移翻译(Exact Offset Translation)
内嵌偏移头使精确消费组定位成为可能:
- 消费组迁移启用时自动添加到目标记录;
- 由
tryFindExactOffset用来精化时间戳翻译结果; - 可处理非单调时间戳与亚毫秒精度场景。
自定义头注入(Custom Headers)
output: redpanda_migrator: headers: x-migration-processed-at: "${! timestamp_unix_milli() }" x-migration-latency-ms: "${! timestamp_unix_milli() - meta(\"kafka_timestamp_ms\") }"十一、测试组织与覆盖
迁移器拥有覆盖单元、集成与浸泡(soak)三类测试的完整测试体系(详见 TESTING.md 与migrator/目录):
internal/impl/redpanda/migrator/ ├── migrator_test.go # 单元测试:输出 lint 规则校验(key/partitioner/partition/timestamp 等字段) ├── conv_test.go # 单元测试:Topic 名称映射(相同与变换名称) ├── migrator_schema_registry_test.go # 单元测试:版本解析(latest/all/非法输入)、Schema 相等比较 ├── migrator_groups_test.go # 单元测试:从组偏移量提取 Topic ├── integration_test.go # 集成测试:端到端迁移 ├── migrator_topic_integration_test.go # 集成测试:Topic 迁移 ├── migrator_schema_registry_integration_test.go # 集成测试:Schema Registry 迁移 ├── migrator_groups_integration_test.go # 集成测试:消费组迁移 ├── integration_soak_test.go # 长时间运行稳定性测试 └── integration_helpers_test.go # 测试基础设施:Docker 化 Redpanda 集群单元测试要点
- 配置与校验:输出端 lint 规则校验;
- 数据转换:Topic 名映射;
- Schema Registry:版本解析、Schema 相等比较(类型、Schema 字符串、引用);
- 消费组:从组偏移量提取 Topic。
集成测试要点(均使用真实 Redpanda 集群)
- 端到端迁移(integration_test.go):单分区迁移 + Schema Registry;畸形 Schema ID 头处理;多分区 + 消费组;Kafka 输入与 franz 消费组兼容;真实 Confluent → Redpanda Serverless 迁移(手动);来源头双向迁移;非单调时间戳的精确偏移翻译;
- Topic 迁移(migrator_topic_integration_test.go):Topic 配置同步、ACL 安全转换复制、幂等同步、分区增长处理;
- Schema Registry 迁移(migrator_schema_registry_integration_test.go):include/exclude 过滤、插值名称解析、版本选择(latest vs all)、ID 翻译模式、相同 Schema 的 ID 复用、规范化、幂等、兼容性传播;
- 消费组迁移(migrator_groups_integration_test.go):带过滤的组偏移量列举、记录时间戳读取、多节点时间戳读取(手动)、完整偏移同步(翻译 + 提交)。
Soak 测试
integration_soak_test.go 提供长时间运行稳定性验证:可持续负载下连续迁移,支持配置时长、消息速率与 Topic 数量,支持内存与 CPU 剖析。
测试基础设施
- 所有集成测试通过 Docker 启动真实 Redpanda 集群(含 Schema Registry);
- 测试验证真实 Kafka 协议交互;
- 消费组测试验证偏移提交行为;
- 最终一致性场景使用
assert.Eventually处理。
十二、实现细节与结论
缓存策略
- Topics:
knownTopics映射避免重复创建尝试; - Schemas:
knownSubjects与knownSchemas映射避免冗余 Schema 操作,并支持 ID 冲突检测(checkSchemaIDConflict); - 消费组:
commitedOffsets映射防止偏移回退。
从源码可以得出的工程要点
- 名称转换器(conv.go):
nameConverter只在源/目标名称不同时存储映射,同名时直通,优化内存占用; - 低层读取(migrator_groups.go):
readRecordAtOffset直接构造kmsg.FetchRequest,定位分区 leader 后按精确偏移读取单条记录,是时间戳翻译与精确精化的底层基础; - 吞吐表现:
bench/目录(bench/README.md)提供了在两集群间迁移 30GB 数据的基准测试(task一键运行),展示约 1GB/s 的迁移吞吐与流式压测模式(100MB/s 持续数据流)。
总结
Redpanda Unified Migrator 将「Topic 基础设施 + Schema 注册表 + 消费组状态」这三件集群迁移中最容易出错的环节,封装为一对输入/输出组件即可使用的统一方案。它通过按需 Topic 创建、ID 翻译、时间戳关联偏移翻译与来源头防护等机制,在保证分区保序、偏移不回退与 ACL 安全的前提下,实现了可重复、幂等、可观测的集群迁移;配合完整的单元/集成/浸泡测试与基准测试,是 Kafka ↔ Redpanda 迁移场景中一套生产可用的参考实现。
【免费下载链接】connectFancy stream processing made operationally mundane项目地址: https://gitcode.com/GitHub_Trending/con/connect
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考