主题与分区


主题作为消息的归类,可以细分为一个或多个分区,分区可以看做是对消息的二次归类。分区的划分不仅为Kafka提供了可伸缩性、水平扩展的功能,还通过多副本机制为kafka提供数据冗余以提供数据的可靠性。

分区可以有一个至多个副本,每个副本对应一个日志文件,每个日志文件对应一个至多个日志分段(LogSegment),每个日志分段还可以细分为索引文件、日志存储文件、和快照文件等。

kafka-topics.sh脚本中的 zookeeper、partitions、replication-factor和topic这4个参数分别代表ZooKeeper连接地址、分区数、副本因子和主题名称。另一个 create 参数表示的是创建主题的指令类型,在 kafka-topics.sh 脚本中对应的还有list、describe、alter和delete这4个同级别的指令类型,每个类型所需要的参数也不尽相同。

生产者的分区分配是指为每条消息指定其所要发往的分区,消费者中的分区分配是指为消费者指定其可以消费消息的分区,分区副本分配是指为集群制定创建主题时的分区副本分配方案,即在哪个broker中创建哪些分区的副本。

为什么分区只能增加而不能减少?

按照Kafka现有的代码逻辑,此功能完全可以实现,不过也会使代码的复杂度急剧增大。实现此功能需要考虑的因素很多,比如删除的分区中的消息该如何处理?如果随着分区一起消失则消息的可靠性得不到保障;如果需要保留则又需要考虑如何保留。直接存储到现有分区的尾部,消息的时间戳就不会递增,如此对于Spark、Flink这类需要消息时间戳(事件时间)的组件将会受到影响;如果分散插入现有的分区,那么在消息量很大的时候,内部的数据复制会占用很大的资源,而且在复制期间,此主题的可用性又如何得到保障?与此同时,顺序性问题、事务性问题,以及分区和副本的状态机切换问题都是不得不面对的。反观这个功能的收益点却是很低的,如果真的需要实现此类功能,则完全可以重新创建一个分区数较小的主题,然后将现有主题中的消息按照既定的逻辑复制过去即可。

kafka-configs.sh脚本使用entity-type参数来指定操作配置的类型,并且使用entity-name参数来指定操作配置的名称。entity-type只可以配置4个值:topics、brokers、clients和users,entity-type与entity-name的对应关系:

entity-type的释义entity-name的释义
主题类型的配置,取值为topics 指定主题的名称
broker类型的配置,取值为brokers 指定brokerid值,即broker.id参数配置的值
客户端类型的配置,取值为clients 指定clientid值, 即kafkaProducer或kafkaConsumer的client.id参数配置的值
用户类型的配置,取值为users 指定用户名

再使用alter指令变更配置时,需要配合add-config和delete-config一起使用。

更改配置

./bin/kafka-configs.sh --zookeeper localhost:2181/kafka --alter --entity-type topics --entity-name topic-demo --add-config cleanup.policy=compact,max.message.bytes=10000

显示配置

./bin/kafka-configs.sh --zookeeper localhost:2181/kafka --describe --entity-type topics --entity-name topic-demo

删除配置

./bin/kafka-configs.sh --zookeeper localhost:2181/kafka --alter --entity-type topics --entity-name topic-demo --delete-config cleanup.policy,max.message.bytes

image.png

在特殊情况下,Kafka集群中的Broker节点有时会出现宕机或者崩溃的问题,当分区的leader节点发生故障时,分区的其中一个follower节点会成为新的leader节点,这样会导致集群的负载不均衡,从而影响整体的健壮性和稳定性。此外,当原来的leader节点恢复之后重新加入集群,它只能成为一个新的follower节点而不再对外提供服务。

为了有效地治理负载失衡的情况,Kafka引入了优先副本的概念。所谓的优先副本也就是AR集合列表中的第一个副本。Kafka确保所有主题的优先副本在Kafka集群中均匀分布,这样就保证了所有分区的leader均衡分布。如果leader过于集中,就会造成集群负载不均衡。事实上,优先副本让之前失去leader副本的节点重新将节点上的优先副本扶持上位成为leader。而之前通过选举的leader又成为了follower节点。

优先副本的选举可以通过命令实现:

./bin/kafka-preferred-replica-election.sh --zookeeper localhost:2181/kafka

当集群中新增broker节点时,只有新创建的主题分区才有可能被分配到这个节点上,而之前的主题分区并不会自动分配到新加入的节点中,因为在它们被创建时还没有这个新节点,这样新节点的负载与原先节点的负载之间严重不均衡。另外一种情形,当集群中的一个节点突然宕机下线时,如果节点上的分区是单副本的,那么这些分区就变得不可用了;如果节点上的分区是多副本的,那么位于这个节点上的leader副本的角色会转交到集群的其它follower副本中。如果不把失效的分区副本迁移到其它可用的broker 节点,那么会影响整体的可用性和可靠性。

为了解决上述问题,需要将分区副本再次进行合理的分配,即分区重分配。Kafka提供了kafka-reassign-partitions.sh脚本来进行分区重分配的工作,以便在集群扩容、broker节点失效的场景下对分区进行迁移。

分区重分配步骤:

创建主题清单json文件:reassign.json

{
    "topics":[
        {
            "topic":"topic-test"
        }
    ],
    "version":1
}

根据这个json文件和指定要分配的broker节点列表来生成一份候选的重分配方案:

./bin/kafka-reassign-partitions.sh --zookeeper localhost:2181/kafka --generate --topics-to-move-json-file reassign.json --broker-list 0,2

project.json:

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

执行重分配方案,其中project.json保存的是前面生成的重分配方案策略。

./bin/kafka-reassign-partitions.sh --zookeeper localhost:2181/kafka --execute --reassignment-json-file project.json

验证重分配是否成功

./bin/kafka-reassign-partitions.sh --zookeeper localhost:2181/kafka --verify --reassignment-json-file project.json

image.png

从上图可以看出,leader副本主要在broker2节点,该节点的负载较大,我们重新选取leader副本

./bin/kafka-preferred-replica-election.sh --zookeeper localhost:2181/kafka

image.png

增加副本因子:

原来的副本分布:

image.png

创建配置文件project.json:

{
    "versions":1,
    "partitions":[
        {
            "topic":"topic-test",
            "partition":1,
            "replicas":[
                1,
                0,
                2
            ],
            "log.dirs":[
                "any",
                "any",
                "any"
            ]
        },
        
        {
            "topic":"topic-test",
            "partition":0,
            "replicas":[
                2,
                1,
                0
            ],
            "log.dirs":[
                "any",
                "any",
                "any"
            ]

        },
        {
            "topic":"topic-test",
            "partition":2,
            "replicas":[
                1,
                0,
                2
            ],
            "log.dirs":[
                "any",
                "any",
                "any"
            ]
        },
        {
            "topic":"topic-test",
            "partition":3,
            "replicas":[
                2,
                0,
                1
            ],
            "log.dirs":[
                "any",
                "any",
                "any"
            ]
        }
    ]
}

执行命令:

./bin/kafka-reassign-partitions.sh --zookeeper localhost:2181/kafka --execute --reassignment-json-file bin/project.json

再次查看副本:

image.png