使用Kafka MirrorMaker 1复制记录时是否可以保留原始分区?
Kafka MirrorMaker 1 向Azure EventHub复制时保留原始分区的解决方法
MirrorMaker 1默认使用内置的DefaultPartitioner重新计算目标分区,不会继承源端记录的原始分区号,配合EventHub的Kafka兼容层默认分区路由规则,就会出现原始分区丢失的问题,可通过以下方法解决:
方法一:使用自定义分区器强制绑定源分区号
这是最稳定的解决方案,前提是你已经保证目标EventHub主题和源端主题分区数完全一致:
- 编写自定义Kafka生产者分区器,核心逻辑直接返回源记录的
partition()值,无需按key/时间等重新计算分区,示例核心逻辑如下:
public class KeepSourcePartitioner implements Partitioner { @Override public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) { // 直接返回源记录携带的原始分区号即可 return sourceRecord.partition(); } // configure、close等默认方法按标准实现即可 }
- 将编写好的分区器打包为jar文件,放到MirrorMaker 1运行环境的类路径下
- 启动MirrorMaker 1时,给生产者配置添加参数:
partitioner.class=com.yourpackage.KeepSourcePartitioner
方法二:调整现有Handler配置实现分区保留
如果不想编写自定义代码,可直接调整MirrorMaker 1的启动配置:
- 给MirrorMaker消费者添加配置:
include.message.headers=true,确保源端消息的分区元数据被完整携带到生产者侧 - 给MirrorMaker生产者添加两个配置:
partitioner.ignore.keys=true,关闭默认分区器按消息key重新计算分区的逻辑azure.eventhub.producer.disable.partition.id.validation=false,关闭EventHub默认的自动分区路由功能,允许生产者指定分区号
验证方式
配置完成后可先做小流量验证:给源端每个分区发送带唯一标识的测试消息,分别消费EventHub目标主题的每个分区,确认消息的分区归属和源端完全一致即可上线运行。
内容的提问来源于stack exchange,提问作者jaysc
相关产品推荐
相关产品推荐

