一 认识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、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.ucloud.cn/yun/127584.html
摘要:而在服务器中应该充分利用多线程来处理执行逻辑。能保证所在的失效,该消息仍然可以从新选举的中获取,不会造成消息丢失。这意味着无需等待来自的确认而继续发送下一批消息。 showImg(https://segmentfault.com/img/remote/1460000018373147?w=702&h=369); 1.概述 Apache Kafka最早是由LinkedIn开源出来的分布式...
阅读 1095·2022-09-27 09:47
阅读 953·2022-09-27 09:28
阅读 1461·2022-09-27 09:16
阅读 772·2022-09-27 08:21
阅读 966·2022-09-27 08:08
阅读 1104·2022-09-18 12:33
阅读 774·2022-09-16 08:01
阅读 755·2022-09-15 12:27