kafka-examples完整指南:MirrorMaker自定义Handler实现跨数据中心Topic重命名
【免费下载链接】kafka-examplesSnippets and small examples demonstrating kafka features and configs项目地址: https://gitcode.com/gh_mirrors/kaf/kafka-examples
kafka-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",每个子目录都是一个独立可运行的小例子:
| 模块 | 演示内容 |
|---|---|
MirrorMakerHandler | MirrorMaker 自定义 MessageHandler(本文主角) |
SimpleCounter | 经典 Kafka 计数器 Producer |
AvroProducerExample/AvroConsumerExample | Avro 序列化生产与消费 |
CountingProducerInterceptor | Producer 拦截器计数 |
AdminClientExample | AdminClient 创建 Topic |
KafkaStreamsAvg/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 加入 CLASSPATH
export CLASSPATH=$CLASSPATH:/path/to/MirrorMakerHandler/target/TopicSwitchingHandler-1.0-SNAPSHOT.jar第 3 步:启动 MirrorMaker 并指定 Handler
bin/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.0(KStream 实现),重命名应改用内置的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),仅供参考