Apache Kafka 分区重分配的实现原理解析

一   认识Kafka


Kafka是一个开源流处理平台,它由 Apache 软件基金会开发的,开发它的目的是为了提供一个统一的、高吞吐、低延迟的实时数据处理平台。它的持久化层与同类平台不同,本质上是一个“按照分布式事务日志架构的大规模发布/订阅消息队列”,这使它非常具有价值。

二   Kafka的使用

想要完成分区副本的重分配,需要在 Kafka 的根路径下,执行如下命令

执行

./bin/kafka‐reassign‐partitions.sh ‐‐zookeeper localhost:2181/kafka 
‐‐reassignment‐json‐file reassign‐topic.json ‐‐execute

分区副本的分布情况由eassign‐topic.json 文件指定,如

{
"version": 1,
"partitions": [
{
"topic": "test",
"partition": 2,
"replicas": [
2,
1
],
"log_dirs": [
"any",
"any"
]
}
}

从上我们可以看出opic=test,partition=2 的分区的两副本分别移动到 brokerId=2 和 brokerId=1 的节点的任意磁盘路径上。

三    ZooKeeper 和 Kafka Controller

3.1  ZooKeeper

Kafka 的元数据存储在 ZooKeeper 中。Apache ZooKeeper是可靠的分布式协调服务框架。它凭借着数据模型类似于文件系统的树形结构,实现保存一些元数据协调信息。同时 ZooKeeper具有 Watch 通知功能。一旦 znode 节点被创建、删除,子节点数量发生变化,或是 znode 所存的数据本身变更, ZooKeeper会及时通知客户端,触发对应的处理操作。

3.2  Kafka Controller

Kafka Controller作为 Apache Kafka 的核心组件,它能够在 Apache ZooKeeper 的帮助下管理和协调整个 Kafka 集群。集群中任意一台 Broker 都能充当控制器的角色。事实上,在运行过程中,只能有一个 Broker 成为控制器,来发送各种操作指令。

四   分区重分配流程

Kafka需要在client、broker 和 controller 的协同运行下完成分区重分配。

流程图如下:

1.png

流程图分析


1、kafka-reassign-partitions 客户端 

先由客户端发起分区重分配任务,它的入口主类为 ReassignPartitionsCommand.scala 中,接着调用 executeAssignment 方法。客户端的 executeAssignment 方法主要完成了如下操作:

·  解析 json 文件 ,进行json 文件校验

·  读取 json 文件内容,判断是否继续执行副本重分配

·  校验分区副本数和副本数据路径数是否一致,校验 partition/replica 是否为空/重复

·  检查待重分配的分区在集群中是否存在,检查确认所有目标 broker 均在线,检查是否已存在分区副本重分配任务

·  分配任务记录,发送 alterReplicaLogDirs 请求

2、controller 维护分区的元数据信息

在 controller 启动时会创建 partitionReassignmentHandler,kafkaController 主线程回调 onControllerFailover 时,当/admin/reassign_partitions 发生变化时,会触发分区副本重分配操作,在 maybeTriggerPartitionReassignment 中通过调用 onPartitionReassignment 真正执行分区副本重分配。

onPartitionReassignment 的执行过程如下:

·  在 zk 中将 AR 更新为 RAR+OAR

·  向所有副本(RAR+OAR)中发送 LeaderAndIsr 请求

·  将 RAR-OAR 的副本状态置为 NewReplica,直到所有 RAR 中的副本完成与 leader 的同步

·  将所有 RAR 的副本置为 OnlineReplica 状态,将 RAR 作为 AR

·  判断 leader 不在 RAR 中,检查 leader 状态,如果 leader 健康则更新 LeaderEpoch,否则重新选择 leader

·  将 OAR-RAR 的副本置为 Offline 状态

·  将 OAR-RAR 的副本置为 NonExistentReplica 状态,并将 zk 中的 AR 置为 RAR(/brokers/topics/${topicName}数据格式:{"version":1,"partitions":{"0":[${brokerId}]}})

·  更新 zk 中/admin/reassign_partitions 的值,同步所有 broker,更新元数据信息

3、broker 端数据跨路径迁移

底层数据跨路径迁移需要 broker 端完成的,broker 接收到客户端发来的请求后,调用 alterReplicaLogDirs 方法

步骤如下:

·  确保目的路径/待移动分区在线

·  标记需要进行迁移的分区副本路径

·  对于需要移动的分区副本,创建 future Log

·  停止当前 Log 的清理工作,等待 future Log 同步

·  创建 ReplicaAlterLogDirsThread,逐个数据构造 Fetch 请求

·  通过 ReplicaManager.fetchMessages 从分区副本 leader 获取数据,完成数据同步

原创文章,作者:网友投稿,如若转载,请注明出处:https://www.cloudads.cn/archives/4195.html

发表评论

登录后才能评论