拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Redpanda Connect Unified Migrator 深度指南:Kafka 与 Redpanda 集群间 Topic、Schema Registry 与消费组的一体化迁移

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.goMigrator以组合方式持有topicMigrator、schemaRegistryMigrator、groupsMigrator三个子迁移器并各自维护自己的缓存结构knownTopicsTopic 映射、knownSubjects/knownSchemasSchema 映射、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_keykgo.Record.Key可选键会原样保留kafka_valuekgo.Record.Value必需若开启 Schema ID 翻译则先解析并改写 IDkafka_topickgo.Record.Topic经插值解析为目标 Topic不存在则自动创建kafka_partitionkgo.Record.Partition分区被保留保证写入顺序与源一致kafka_timestamp_mskgo.Record.Timestamp必需毫秒时间戳转换为time.Timekafka_offset偏移头消费组迁移开启时写入offset_header8 字节大端编码kafka_headerskgo.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 Resolutiontopic字段为可插值字符串从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 迁移器按需创建 幂等同步 安全 ACLTopic 同步是“按需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调用Subjectsinclude_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.gotranslate_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 客户端同时读取源集群 IDonInputConnectedmigrator.go。输出连接获取目标集群元数据、admin 客户端与集群 ID启动 Topic 周期同步循环执行一次 Schema 注册表同步并启动其周期循环最后启动消费组同步循环onOutputConnectedmigrator.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 Migrationinput: 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 Transformationoutput: 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字段类型默认值说明urlstring—注册表基础 URL必需timeoutduration5sHTTP 客户端超时tlsobject—TLS 配置enabledbooltrue是否启用 Schema 迁移intervalduration5m同步周期0s仅启动时同步一次include[]string空Subject 包含正则空则全包含exclude[]string空Subject 排除正则优先级高于 includesubjectstring—Subject 名称插值模板versionsenumalllatest或allinclude_deletedboolfalse是否包含软删除的 Schematranslate_idsboolfalse是否翻译 Schema IDnormalizeboolfalse创建时是否规范化 Schemastrictboolfalse未知 Schema ID 是否报错仅 translate_ids 时有效max_parallel_http_requestsint10并发 HTTP 请求数上限schema_registry块还支持标准 HTTP 请求认证字段如basic_auth见 7.6 Serverless 示例。7.4 消费组 过滤Consumer Groups with Filteringoutput: 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字段类型默认值说明enabledbooltrue是否启用消费组迁移intervalduration1m同步周期0s禁用fetch_timeoutduration10s读取记录用于时间戳翻译的最大等待时间低吞吐集群可调大include[]string空组名包含正则exclude[]string空组名排除正则优先级高于 includeonly_emptyboolfalsetrue仅迁移 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_factorint继承源目标 Topic 副本因子迁移到不同规模集群时很有用sync_topic_intervalduration5mTopic 周期同步间隔0s禁用首条消息仍会创建sync_topic_aclsboolfalse是否同步 Topic ACL带安全转换headersmap[string]string—追加到迁移记录的自定义头插值与 provenance/offset 头同名会被忽略provenance_headerstringredpanda-migrator-provenance来源头名置空则不添加offset_headerstringredpanda-migrator-offset偏移头名置空则禁用精确偏移翻译max_in_flightint10在途批次上限建议设为并行复制的分区总数高吞吐调优官方文档给出的调优建议输入侧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_totalcounter—— 成功创建的 Topic 数redpanda_migrator_topic_create_errors_totalcounter—— Topic 创建失败数redpanda_migrator_topic_create_latency_nstimer—— Topic 创建延迟。Schema Registry 迁移指标redpanda_migrator_sr_schemas_created_totalcounter—— 成功创建的 Schema 数redpanda_migrator_sr_schema_create_errors_totalcounter—— Schema 创建失败数redpanda_migrator_sr_schema_create_latency_nstimer—— Schema 创建延迟redpanda_migrator_sr_compatibility_updates_totalcounter—— 兼容性更新次数redpanda_migrator_sr_compatibility_update_errors_totalcounter—— 兼容性更新失败数redpanda_migrator_sr_compatibility_update_latency_nstimer—— 兼容性更新延迟。消费组迁移指标带group标签redpanda_migrator_cg_offsets_translated_totalcounter—— 成功翻译的偏移量数redpanda_migrator_cg_offset_translation_errors_totalcounter—— 偏移量翻译失败数redpanda_migrator_cg_offset_translation_latency_nstimer—— 偏移量翻译延迟redpanda_migrator_cg_offsets_committed_totalcounter—— 成功提交的偏移量数redpanda_migrator_cg_offset_commit_errors_totalcounter—— 偏移量提交失败数redpanda_migrator_cg_offset_commit_latency_nstimer—— 偏移量提交延迟。消费滞后指标带topic、partition标签redpanda_laggauge—— 迁移器输入在每个分区上的当前消费滞后高水位与当前消费位置的差值用于观察迁移进度是否跟得上生产速度。九、保证与限制保证GuaranteesTopic 分区数目标 Topic 以匹配的分区数创建不回退No offset rewind消费组偏移量绝不向回移动ACL 安全排除 WRITE 操作、ALL 降级为 READ幂等重复同步安全Topic、Schema、消费组均是。限制Limitations偏移量翻译为尽力而为若前一条记录的时间戳无法读取或目标端在该时间戳后没有偏移量则跳过该分区分区数一致要求消费组迁移要求源与目标 Topic 分区数一致不一致的 Topic 会被filterTopics跳过migrator_groups.goSchema 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 Replicationoutput: redpanda_migrator: sync_topic_acls: true排除ALLOW WRITE条目将ALLOW ALL降级为ALLOW READ保留资源模式类型与主机过滤。Schema 规范化Schema Normalizationoutput: redpanda_migrator: schema_registry: normalize: true创建 Schema 时规范化格式保证语义等价但格式不同的 Schema 可被识别为相同源码通过schemaEquals/schemaStringEquals实现 JSON/Avro 的 JSON 语义比较与 Protobuf 的空白归一比较见 migrator_schema_registry.go。精确偏移翻译Exact Offset Translation内嵌偏移头使精确消费组定位成为可能消费组迁移启用时自动添加到目标记录由tryFindExactOffset用来精化时间戳翻译结果可处理非单调时间戳与亚毫秒精度场景。自定义头注入Custom Headersoutput: 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.goTopic 配置同步、ACL 安全转换复制、幂等同步、分区增长处理Schema Registry 迁移migrator_schema_registry_integration_test.goinclude/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处理。十二、实现细节与结论缓存策略TopicsknownTopics映射避免重复创建尝试SchemasknownSubjects与knownSchemas映射避免冗余 Schema 操作并支持 ID 冲突检测checkSchemaIDConflict消费组commitedOffsets映射防止偏移回退。从源码可以得出的工程要点名称转换器conv.gonameConverter只在源/目标名称不同时存储映射同名时直通优化内存占用低层读取migrator_groups.goreadRecordAtOffset直接构造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),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门