
kafka-examples完整指南MirrorMaker自定义Handler实现跨数据中心Topic重命名【免费下载链接】kafka-examplesSnippets and small examples demonstrating kafka features and configs项目地址: https://gitcode.com/gh_mirrors/kaf/kafka-exampleskafka-examples是一个汇集 Kafka 特性代码片段的开源示例库其 MirrorMakerHandler 模块展示了如何用不到 30 行 Java 代码编写TopicSwitchingHandler让MirrorMaker在跨数据中心同步时为 Topic 自动加上数据中心前缀如dc1.mm1轻松实现跨数据中心 Topic 重命名避免双集群 Topic 冲突。本文带你完整走通原理 → 编译 → 部署 → 验证全流程。一、kafka-examples 项目是什么 kafka-examples 定位为 Snippets and small examples demonstrating kafka features and configs每个子目录都是一个独立可运行的小例子模块演示内容MirrorMakerHandlerMirrorMaker 自定义 MessageHandler本文主角SimpleCounter经典 Kafka 计数器 ProducerAvroProducerExample/AvroConsumerExampleAvro 序列化生产与消费CountingProducerInterceptorProducer 拦截器计数AdminClientExampleAdminClient 创建 TopicKafkaStreamsAvg/StreamingAvg流式移动平均如需获取完整代码执行git clone https://gitcode.com/gh_mirrors/kaf/kafka-examples二、为什么跨数据中心同步要重命名 Topic MirrorMaker 是 Kafka 自带的集群间镜像工具。当你在dc1 → dc2双向同步时两边往往存在同名 Topic如都叫orders不加区分地镜像回环同步会把消息无限复制下游消费者无法判断消息来自哪个数据中心。TopicSwitchingHandler的思路很直接给源集群的每个 Topic 加前缀orders同步到对端后变成dc1.orders从命名上隔离出这条消息来自 dc1。三、TopicSwitchingHandler 核心机制拆解 ⚙️核心源码见 TopicSwitchingHandler.java实现只有一个接口MirrorMaker.MirrorMakerMessageHandler。1. 构造函数只接收一个前缀参数public TopicSwitchingHandler(String topicPrefix) { this.topicPrefix topicPrefix; }2. 两个 handle() 重载分别对应 Kafka 旧版MessageAndMetadata与新版BaseConsumerRecord两种消费记录类型保证兼容不同版本的 MirrorMaker 调用入口。核心逻辑只有一行new ProducerRecord(topicPrefix . record.topic(), record.partition(), record.key(), record.message());可以看到分区号、key、消息体原样保留只有 Topic 名被替换为前缀.原Topic。这就是跨数据中心 Topic 重命名的全部秘密。四、三步快速上手从编译到运行 完整命令参考 README.md。第 1 步编译生成 jar在MirrorMakerHandler/目录执行mvn package得到target/TopicSwitchingHandler-1.0-SNAPSHOT.jar。第 2 步把 jar 加入 CLASSPATHexport CLASSPATH$CLASSPATH:/path/to/MirrorMakerHandler/target/TopicSwitchingHandler-1.0-SNAPSHOT.jar第 3 步启动 MirrorMaker 并指定 Handlerbin/kafka-mirror-maker.sh \ --consumer.config config/consumer.properties \ --message.handler com.shapira.examples.TopicSwitchingHandler \ --message.handler.args dc1 \ --producer.config config/producer.properties \ --whitelist mm1两个关键参数--message.handler.args dc1即 Handler 的topicPrefix消息将写入dc1.mm1--whitelist mm1只同步mm1这个 Topic。五、验证同步结果 ✅在源端往mm1生产消息bin/kafka-console-producer.sh --topic mm1 --broker-list localhost:9092在对端从重命名后的dc1.mm1消费能收到刚才的消息即说明重命名生效bin/kafka-console-consumer.sh --topic dc1.mm1 --zookeeper localhost:2181 --from-beginning六、实战注意事项 ⚠️版本兼容示例基于kafka_2.11 0.9.0.0见 pom.xml。若使用 Kafka 2.x 的 MirrorMaker 2.0KStream 实现重命名应改用内置的topics.regex.rewrite配置但Handler 拦截改写记录的思想完全一致可照搬本文handle()的写法元数据不迁移该 Handler 只改写消息不自动创建目标 Topic生产环境需配合 Topic 预创建自定义改写规则想改成dc1-mm1这种下划线风格只需修改handle()里的字符串拼接逻辑一个正则就能扩展成任意命名规范。七、举一反三项目中的其他小示例kafka-examples 里还有大量同风格的片段值得细读比如生产者拦截器 CountingProducerInterceptor.java、Avro 会话化消费 AvroClicksSessionizer.java、Kafka Streams 移动平均 StreamingAvg.java每个都足够短适合作为学习 Kafka API 的活文档。总结用TopicSwitchingHandler为 MirrorMaker 定制重命名逻辑只需实现一个接口、改写一行 Topic 拼接再配合--message.handler启动参数即可上线。kafka-examples 用这个极简样例证明了MirrorMaker 的扩展点非常友好跨数据中心的数据命名隔离可以像加个前缀一样简单。【免费下载链接】kafka-examplesSnippets and small examples demonstrating kafka features and configs项目地址: https://gitcode.com/gh_mirrors/kaf/kafka-examples创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考