这个系列写到第九篇,Kafka的核心骨架已经基本过了一遍:整体架构、broker部署、生产者和消费者的客户端用法、主题与分区的基本概念,前面都聊过。按正常的学习路线,接下来应该进入一个看似平平无奇、但实际坑最深的地方——主题创建。很多人在生产上用Kafka的第一步就是建主题,但建主题这个动作背后,藏着一整条完整的链路:客户端请求协议、服务端Controller状态机、元数据存储与广播、分区副本分配算法。你如果没把这层纸捅破,后面遇到“分区不均匀”“副本一直UnderReplicated”“创建主题超时”这类问题,会完全无从下手。
这篇内容会聚焦三块:主题创建代码怎么写的、分区副本自动分配的三种策略怎么选、一个Topic从API调用到全集群可见的底层流程是什么。适合已经能把Kafka跑起来、但想进一步理解集群内部机制的读者。看完之后,你不光能写出健壮的主题创建代码,还能在面试和排障时把底层逻辑讲清楚。
1. 主题创建的前置认知:为什么Kafka把建Topic这件事做得这么重
1.1 主题在Kafka中的地位
如果拿快递系统做类比,Kafka集群就是一个大型分拨中心,主题就是分拨中心里一条条独立的传送带线路。每条线路有自己的分拣通道(分区),每个通道有备用货物槽位(副本),而货物本身(消息)只有放到对应通道里,下游消费者才能按顺序提走。
这个类比想说明一件事:创建主题不是在“组织里加一条记录”这么简单,它意味着集群要为它分配物理存储、启动副本同步、广播元数据,并且让所有broker和客户端都知道“这条通道现在开通了”。所以Kafka把主题创建做成了集群级的元数据变更操作,而不是本地写一条配置。
这个认知很重要。我见过不少刚上手的人,以为创建一个Topic就和往MySQL里insert一条记录一样,执行完命令就万事大吉。实际上一个Topic的创建,涉及Controller的选举协调、元数据落盘、全集群broker的元数据刷新、分区leader的选举等多个环节。只看到一条命令执行成功,看不到背后的链路,后面排查问题就会很被动。
1.2 创建主题本质上是在做集群元数据变更
主题创建的命令行工具kafka-topics.sh --create,底层调用的其实是Kafka的AdminClient API;而AdminClient发出的CREATE_TOPICS请求,最后会交给集群中Controller角色的broker处理。也就是说,主题创建不是某个broker独自完成的,而是由Controller统一协调整个集群完成的。
这里引出一个关键点:Controller在Kafka集群中扮演的相当于“元数据总指挥”的角色。所有主题的增删改、分区的扩缩容、leader的选举,都要经过Controller。创建主题就是给Controller下发一条“我要新增这些分区和副本”的指令,Controller做完校验和分配后,把元数据写入到ZooKeeper(老架构)或KRaft元数据日志(新架构),然后通知所有broker更新自己的元数据缓存。
所以,如果你在代码里创建主题成功,只能说明Controller已经接受了请求,不代表所有broker都已经感知到新主题。这个时间差在极端情况下会带来“刚建完topic马上生产就报LEADER_NOT_AVAILABLE”的现象。后面第五部分专门讲这个问题。
1.3 命令行工具与代码管理的取舍
先给个结论:临时调试用命令行,生产环境建议用代码管理。
命令行kafka-topics.sh的优点是快,一条命令建完,特别适合本地开发环境测试。但它的缺点也很明显:操作不可审计,参数容易输错,而且多人共用一套集群的时候很难追溯“这个主题是谁建的、为什么副本因子只有2”。一旦线上主题被误删或建错,找原因都找不到。
代码管理的好处在于,你可以把主题创建封装成一个标准化的服务:创建前检查是否已存在、配置统一的副本因子和分区数、记录操作日志、加上审批流程。我自己在团队里就维护了一个小工具类,所有主题创建都走这个入口,半年下来几乎没再出过因手动误操作导致的主题配置问题。
当然,这里的代码管理不是说要自己另造一个平台,而是建议写一个封装了KafkaAdminClient的统一工具,甚至只是一个带参数校验的脚本都可以。核心目标是让“创建主题”这个操作可重复、可控、可审计。
2. 创建主题的代码简析:从KafkaAdminClient到CreateTopicsRequest
2.1 KafkaAdminClient为什么是我推荐的管理入口
Kafka的Java客户端里有一个独立的组件叫KafkaAdminClient,专门负责集群管理操作。早期版本还有AdminUtils和AdminClient两种入口,但后来的版本已经收敛到org.apache.kafka.clients.admin.Admin这个接口里。日常管理主题、查看消费者组、查询分区状态、调整配置,都可以用它搞定。
推荐用Admin接口而不是自己拼ZooKeeper命令去写节点,原因有三点。第一,Admin接口的协议是走Kafka原生TCP协议,不是依赖ZooKeeper的临时节点操作,对ZooKeeper集群的压力更小;第二,它天然兼容新老版本,KRaft模式下操作方式也一样;第三,它返回的Future和Result对象可以方便地处理异步场景,避免阻塞主线程。
有一点必须提醒:KafkaAdminClient不是轻量组件,它内部会创建一套NetworkClient连接池,所以用法上要尽量复用,不要每次操作都new一个。正确的姿势是用一次,用完关闭,或者在应用里做成单例。创建主题这种低频操作,直接在工具方法里用try-with-resources方式关闭就行。
2.2 最小可运行示例:代码先跑通再说
下面这段代码是我在实际项目里精简后的创建主题核心逻辑。先说明,这个例子故意去掉了复杂的校验和异常分类,只保留主干,方便理解。
import org.apache.kafka.clients.admin.Admin; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.admin.CreateTopicsResult; import java.util.Collections; import java.util.Map; import java.util.Properties; import java.util.concurrent.TimeUnit; public class CreateTopicExample { public static void main(String[] args) throws Exception { Properties props = new Properties(); // 指向Kafka集群地址,多个broker用逗号分隔 props.put("bootstrap.servers", "192.168.1.10:9092,192.168.1.11:9092"); // 请求超时时间,建议显式设置 props.put("request.timeout.ms", 30000); // 推荐使用try-with-resources方式,用完后自动释放连接 try (Admin admin = Admin.create(props)) { // 构造主题描述:名称、分区数、副本因子 NewTopic newTopic = new NewTopic("order-events", 12, (short) 3); // 可选:给主题设置单独的配置覆盖,不设置则沿用broker默认值 newTopic.configs(Collections.singletonMap("retention.ms", "604800000")); // 发起创建请求 CreateTopicsResult result = admin.createTopics(Collections.singleton(newTopic)); // 阻塞等待创建结果,注意timeout一定要给,避免永久等待 result.all().get(30, TimeUnit.SECONDS); System.out.println("topic created successfully"); } catch (Exception e) { System.err.println("create topic failed: " + e.getMessage()); throw e; } } }这里有几个点值得展开。
第一,NewTopic构造器里的replicationFactor是short类型,传int会编译报错,新手经常踩。第二,createTopics()接受的是集合,所以一次可以创建多个主题,批量创建时会合并成一个请求,效率更高。第三,result.all().get(30, TimeUnit.SECONDS)这一步非常关键,它把异步请求变成了同步等待,而且显式限制了等待时间。如果不加超时,一旦集群Controller出现异常,主线程可能一直阻塞在那里。
我通常还会在每个主题创建前先做一次存在性检查,用admin.listTopics()看看名字是否已经被占用。虽然服务端也有校验,但客户端提前拦一下,可以避免把“主题已存在”这种业务错误混入真正的异常流程里。
2.3 请求在客户端内部的流转过程
很多人写代码只看到API层面,不知道一个createTopics()调用背后做了什么。其实AdminClient和普通生产者一样,底层也是通过NetworkClient发送二进制协议请求。
大致的流转是这样的:admin.createTopics()会构建一个CreateTopicsRequest对象,里面包含主题名称、分区数、副本因子、配置项等字段。然后这个请求被交给KafkaClient的发送队列,经过Kafka协议编码,通过Socket连接发给集群中任意一个broker。broker收到后,从请求头里解析出API Key,发现是CREATE_TOPICS请求,就会把它路由给对应的处理器,最终转发到Controller。
这里有个细节:AdminClient的请求不是必须发给Controller,发给任意broker都可以,因为broker之间的内部协议会自动把管理类请求转发给Controller处理。这算Kafka设计上的一个便利性设计,但也意味着在broker特别多的大集群里,管理请求的链路会比数据请求略长。
2.4 服务端处理入口:AdminManager的职责
服务端收到CreateTopicsRequest之后,实际处理逻辑主要在AdminManager。从Kafka源码看,KafkaApis.handleCreateTopicsRequest()会做第一层校验,比如请求协议版本是否支持、主题名称是否为空;然后会把请求交给AdminManager.handleCreateTopicsRequest()做真正的处理。
AdminManager里做的事情可以拆成四步。第一步,遍历请求里的每个主题,检查命名是否合法、是否已存在。第二步,检查副本因子是否超过broker数量、分区数是否合法。第三步,调用元数据管理器创建主题,这步会触发分区副本的分配计算。第四步,把创建结果封装成CreateTopicsResponse返回给客户端。
其中第三步是整个链路的核心。在ZooKeeper架构下,它会在ZooKeeper里创建/brokers/topics/{topic}节点,并写入分区副本分配信息。在KRaft架构下,它会向__cluster_metadata主题写入一条记录,由Controller Quorum复制。听完这层,你就明白为什么说创建主题是集群级操作了。
2.5 配置参数的优先级与校验
创建主题时,除了分区数和副本因子,还可以通过NewTopic.configs()指定topic级别的配置,比如retention.ms、segment.bytes、cleanup.policy等。这些配置如果没有指定,就会使用集群的默认值;指定了就会单独记录在这个topic的元数据里,覆盖broker端的默认配置。
优先级从高到低是:topic级配置 > broker端动态配置 > broker端静态配置文件(server.properties)。这个优先级关系平时容易忽略,我见过有人改了broker的default.replication.factor,发现新主题仍然是2副本,后来一查才发现是脚本里显式指定了副本因子是2。
另外要注意,不是所有配置都能topic级覆盖。比如log.dir这种属于broker物理目录的配置就不能在topic级别设置。好在Kafka本身有配置校验,传了不支持的配置会被直接拒绝,不会出现半生效的情况。
3. 分区副本分配策略详解:自动分配背后的三套算法
3.1 分配入口与版本差异
主题创建时,分区和副本具体落在哪几个broker上,可以由客户端指定,也可以让服务端自动算。自动算的逻辑,在旧版本里叫AdminUtils.assignReplicasToBrokers,新版本虽然把代码挪到了Controller的元数据管理模块里,但核心算法思路基本没变。
两套算法有必要区分开。第一套是不考虑机架信息的分配,叫RackUnaware;第二套是考虑机架信息的分配,叫RackAware。适合什么场景,取决于你的broker是不是部署在多机架(比如多个物理机房、多个交换机分区)环境下。
3.2 不考虑机架的分配策略(RackUnaware)
RackUnaware的算法核心很简单:把broker列表当做一个环形队列,从随机的一个起始位置开始,依次给每个分区分配副本,分配时再引入一个“副本偏移量”来尽量避免连续分区使用完全相同的副本组合。
口头描述比较抽象,举一个具体例子。假设集群有5个broker,编号0到4,要创建6个分区、副本因子3。第一次计算时可能是这样的:
- 分区0的副本:broker[0], broker[1], broker[2]
- 分区1的副本:broker[1], broker[2], broker[3]
- 分区2的副本:broker[2], broker[3], broker[4]
- 分区3的副本:broker[3], broker[4], broker[0]
- 分区4的副本:broker[4], broker[0], broker[1]
- 分区5的副本:broker[0], broker[1], broker[2]
看到规律了吗?每个分区的第一个副本作为leader,后续副本依次往后挪一位。这样从整体看,每个broker上的leader数量会比较均衡。但这只是理想情况,因为算法里有个随机起点,所以不是每次分配都这么规整。引入随机起点的目的是避免多个Topic同时创建时形成固定模式,造成局部热点。
RackUnaware的优点是计算快、实现简单,适合单机架或broker分布比较均匀的集群。缺点也很明显:它完全不感知broker的物理位置。如果两个broker其实坐在同一个机柜下,同一个交换机的下游,一旦这个机柜断电,某个分区的多个副本可能同时下线,数据就彻底不可用了。所以生产环境有条件的话,尽量还是用RackAware。
3.3 机架感知分配策略(RackAware)
RackAware的算法比RackUnaware复杂一些,核心约束是:同一分区的副本尽量分散到不同机架。这样任意一个机架故障,至少还能保证其他机架上保留完整副本。
算法流程可以这样理解。先把所有broker按机架分组,比如机架A有broker0、broker1,机架B有broker2、broker3,机架C有broker4、broker5。然后给每个分区分配副本时,先轮流从每个机架里选一个broker,保证第一轮每个机架都能分到副本;如果副本因子大于机架数,再从机架列表开始第二轮分配,此时才允许某个机架出现多个副本。
用前面的例子,6个分区、副本因子3、机架3个,分配结果会倾向于:
- 分区0的副本:broker0(机架A)、broker2(机架B)、broker4(机架C)
- 分区1的副本:broker1(机架A)、broker3(机架B)、broker5(机架C)
这样每个分区都横跨三个机架,单机架故障不会丢数据。
启用RackAware需要两个条件。第一,broker的server.properties里必须配置rack.id,比如rack.id=rack-a。第二,创建主题时要么使用新版Admin接口(自动识别broker机架信息),要么用命令行时确保分配算法走的是机架感知逻辑。在KRaft模式下机架信息也在元数据里注册,逻辑类似。
我踩过一次坑:接了机架感知的broker配置,但线上主题创建用的还是自己写的老脚本,结果分配结果完全没体现机架隔离。后来排查发现脚本里有一行代码强制指定了replicasAssignments,自定义分配直接绕过了自动机架算法。所以这里引出一个重要的点:只要创建时手动指定了副本分布,任何自动分配策略都不会生效。
3.4 自定义分配:手动指定副本分布
除了自动分配,Kafka还允许完全由用户指定每个分区的副本分配方案。在Admin接口里通过如下方式传入:
Map<Integer, List<Integer>> replicasAssignments = new HashMap<>(); // 分区0副本放在broker 1、2、3 replicasAssignments.put(0, Arrays.asList(1, 2, 3)); // 分区1副本放在broker 4、5、0 replicasAssignments.put(1, Arrays.asList(4, 5, 0)); NewTopic newTopic = new NewTopic("manual-topic", replicasAssignments);这种写法的好处是完全可控,特别适合做跨机房容灾。比如公司有两个机房,要求奇数分区的主副本在A机房、偶数分区的主副本在B机房,自动算法做不到这么细,就必须手动指定。
但手动指定的代价是要自己承担分配合理性。一旦broker宕机,手动指定的副本分布如果过于集中,就很容易出现数据不可用。所以我的建议是:除非有明确的容灾约束,否则优先使用自动分配;如果非要用手动指定,一定要用脚本校验同一分区副本没有落在同一个broker上,并且各broker的leader分布尽量均衡。
3.5 顺带补充:消费端的Assignor策略别搞混
很多人把“分区副本分配”和“消费者分区分配”混为一谈。前者是主题创建时,broker上副本的物理分布;后者是消费者组启动时,各个消费者实例瓜分哪些分区。
消费者端的分配策略主要有三种。RangeAssignor按主题顺序连续分配,对大数量主题不友好;RoundRobinAssignor把所有分区拉通之后轮流分配,整体更均衡;StickyAssignor在保持上次分配尽量不变的前提下重新均衡,减少rebalance时的不必要分区变动。Kafka从2.3版本开始逐渐倾向于StickyAssignor思路,新版本默认的partition.assignment.strategy参数是一个列表,包含了RangeAssignor和CooperativeStickyAssignor。
这个知识点放在这篇里是因为它在面试里经常和主题创建一起被考察,很多人一紧张就混。简单记法:创建主题看broker,消费组看consumer。
4. 底层流程分析:一个Topic从API到全集群可见的完整链路
4.1 客户端视角的异步与Future
代码层面,admin.createTopics()方法并没有真的把请求发出去,它只是构建了一个CreateTopicsResult对象。真正的网络请求是在调get()或者whenComplete()的时候才被触发完成的。
这里有个内部机制值得了解:CreateTopicsResult内部持有多个KafkaFuture,分别对应每个主题的创建结果。你可以逐个处理,也可以直接用result.all()统一处理所有主题。实际开发里我推荐用all(),因为创建主题通常是一个批量运维操作,失败了直接看整体结果,处理逻辑更简单。
关于超时设置,再啰嗦一遍。future.get(30, TimeUnit.SECONDS)里的30秒不是随便写的。在主题数量多、Controller繁忙的场景下,创建请求可能需要几秒甚至十几秒。给太短容易误报失败,给太长会拖住线程。我一般建议在集群正常状态下,创建单主题控制在5秒内;批量创建20个以内主题,30秒足够。如果经常超时,应该去查Controller所在broker的CPU和元数据写入延迟,而不是盲目调大超时。
4.2 Controller节点的处理链路
请求到了Controller所在的broker,服务端处理器会经历这些步骤。先由KafkaApis识别请求类型,然后交给AdminManager,再调用Controller的元数据更新模块。Controller会做几个校验,比如:
- 主题名是否合法(不能为空、不能包含非法字符、不能以
__开头除非是内部主题) - 是否已存在同名主题
- 分区数和副本因子的值是否在合理范围
- 集群内可用的broker数量是否满足副本因子要求
这些校验都过了之后,Controller才开始真正分配副本。分配结果就是第三部分讲的那几种算法。分配完成后,Controller要更新元数据存储,并在内存中更新自己的元数据缓存。
校验环节最容易出问题的就是“副本因子大于可用broker数”。比如只有2个broker的集群,你建主题时写了副本因子3,Controller会直接报Replication factor: 3 larger than available brokers: 2。这个错误提示很明确,但它的判定依据是“可用”broker,不是“所有”broker。如果一个broker处在/brokers/ids里但当前不可用,也可能导致同样的报错。遇到这种case,优先检查集群是不是有broker失联。
4.3 元数据在ZooKeeper/KRaft中的落地
ZooKeeper架构下,Controller创建主题的最后一步是在ZooKeeper的/brokers/topics/{topic}路径下写一个节点,节点内容包含每个分区的副本分配列表。写完这个节点后,ZooKeeper会触发/brokers/topics路径的watcher事件,其他broker就能感知到有变化。
KRaft架构下流程有些不同。Kafka 3.x版本开始力推KRaft,元数据不再存ZooKeeper,而是存在内部的__cluster_metadata主题里。Controller的Quorum机制负责元数据复制和顺序保证。这也是Kafka号称摆脱ZooKeeper依赖后的核心变化。对于应用开发者来说,底层存储变化不影响Admin API的调用方式,但排查问题时看到的现象略有区别:ZooKeeper模式下可以用zk客户端直接看节点,KRaft模式下要用kafka-metadata-shell.sh去查看元数据日志。
4.4 元数据传播到所有broker
元数据写入完毕,还不是终点。Controller会向所有broker发送UpdateMetadataRequest,通知大家新主题的分区信息、leader和ISR等情况。每个broker收到后更新自己的MetadataCache,这样生产者和消费者的元数据请求才能拿到新主题对应的分区信息。
客户端侧也有一个元数据自动更新机制。生产者默认每隔metadata.max.age.ms(默认5分钟)会重新拉取一次元数据,但当它发现某个Topic不存在时会触发立即刷新。所以一般情况下,主题创建成功后几秒钟,新分区就能正常读写。
不过这里有一个并发时间窗口:如果你在创建主题后立刻发消息,生产者的元数据可能还没刷新到最新状态,就会报LEADER_NOT_AVAILABLE或者UNKNOWN_TOPIC_OR_PARTITION。这不是代码写错了,而是元数据传播有延迟。解决方式就是生产者客户端要配置合理的重试机制,等元数据追上;或者代码层面在创建完主题后主动调一次admin.describeTopics(topic),等返回成功再开始生产。
4.5 全链路时序总结
用文字把整条链路按执行顺序列一遍,排查问题时对着这个清单看,效率会高很多。
- AdminClient发起
CreateTopicsRequest,请求先发给任意一个broker - broker根据管理类请求路由规则,把请求转发给当前Controller
- Controller进行合法性校验,校验失败直接返回错误
- Controller计算分区副本分配方案(自动或根据自定义分配)
- Controller将元数据写入ZooKeeper或KRaft日志
- Controller向所有broker广播
UpdateMetadataRequest - 各broker更新本地元数据缓存,响应Controller
- Controller把创建成功结果返回给AdminClient
- AdminClient通过Future让调用方感知创建结果
- 生产者客户端通过元数据请求获取新分区信息,开始正常生产
这十步里任何一步出问题,都会造成主题创建失败或创建后不可用。平时排障就是对照这个链路,先看客户端报错,再看Controller日志,再看元数据存储,最后看broker的元数据缓存状态。
5. 常见问题与排查技巧实录
5.1 副本因子大于可用broker数量
这是新手最容易踩的坑。报错信息长这样:
ERROR Error while creating topic: 'test-topic' - Replication factor: 3 larger than available brokers: 2原因很简单,副本因子超过了当前可用的broker数,Controller连分配到哪个broker上都做不了。解决办法有三个方向:减小副本因子到小于等于broker数;增加broker节点;如果集群本来就有足够节点,检查是否有broker失联导致“可用”数偏小。
我建议生产环境副本因子统一设3,broker数量至少3个起步。有些开发环境为了省资源只起1个broker,副本因子就只能写1。但只有单副本意味着没有冗余,磁盘故障直接丢数据,所以这个简化只适合测试环境。
5.2 创建成功后立刻生产,报LEADER_NOT_AVAILABLE
这个现象在本地测试时很常见:kafka-topics.sh --create执行成功,紧接着用kafka-console-producer.sh发消息,却报LEADER_NOT_AVAILABLE或UNKNOWN_TOPIC_OR_PARTITION。原因是元数据还没有完全同步到客户端。
元数据的传播是异步的,创建成功只代表Controller处理完了,不代表所有broker和客户端都已经知道新主题的分区leader在哪。生产者的元数据刷新有固定周期,当它发现请求某个不存在的分区时,会触发元数据强制刷新,但刷新本身也需要时间。
遇到这种情况最省事的办法就是重试。生产者的retries参数默认是Integer.MAX_VALUE,所以这种临时的元数据错误一般会自动恢复。如果你是自己写的发送脚本,可以在创建完主题后主动睡眠几秒,或者调用一次AdminClient.describeTopics(),等元数据真正同步后再发数据。
5.3 分区分配不均匀
集群里broker数量很多,但某些broker上的分区数明显比其他broker多,这种“热点”现象需要及时处理。分区分配不均匀往往有几个原因:一是早期集群扩容后,新broker没有自动分担旧分区;二是手动创建主题时指定了不合理的副本分布;三是自动分配算法随机起点加上多次创建后累积偏差。
排查时先用kafka-topics.sh --describe看每个broker的分区数,再用kafka-reassign-partitions.sh做一次全局迁移,把分区从多的broker挪到少的broker。注意迁移会影响磁盘IO和网络流量,尽量在业务低峰期操作,并且迁移完成后要观察ISR是否恢复正常。
5.4 机架感知未生效
如果确认集群里所有broker都配置了正确的rack.id,创建出来的主题分区副本却还是扎堆在同一机架,一般是两个原因:创建方式绕过了自动分配,比如手动指定了replicasAssignments;或者创建代码走的是老版本客户端协议,broker端没能获取到机架信息。
排查方式是先看主题的分区副本分布,确认是否真的没做机架隔离。再查broker端日志里有没有机架相关的告警,最后检查创建代码的版本和调用路径。还有一个隐蔽细节:ZooKeeper架构下,broker启动时注册到ZooKeeper的信息里包含机架信息,如果broker顺序启动有先后,某些broker的rack信息还没注册完就执行了主题创建,Controller很可能拿到不完整的机架列表。所以集群刚启动完的那段时间,不要急着批量建主题。
5.5 创建主题超时
创建主题超时是比较棘手的一类问题,因为报错只是“timed out”,没有具体原因。常见诱因包括Controller侧CPU飙高、ZooKeeper或KRaft元数据写入延迟大、网络分区导致Controller无法和多数broker通信。
排查思路先看Controller日志有没有异常,再看ZooKeeper(或KRaft)集群的健康状态,最后看broker之间的网络延迟。如果Controller所在的broker频繁Full GC,也会导致请求积压。线上Kafka集群的Controller内存要适当调大,同时监控Full GC频率。还有一种情况是元数据积压,broker处理UpdateMetadataRequest太慢,需要检查broker的IO负载。
我在实际中遇到最多的是Kafka环境内存溢出(OOM)导致Controller假死,看起来像是创建超时,实际是broker进程已经半死不活。所以给broker设置合理的内存参数、预留足够堆外内存,比任何参数调优都优先。
6. 实战建议与经验收尾
6.1 分区数规划的几个判断维度
创建主题时必须决定分区数。这个数定少了,后面想扩容有操作成本;定多了,又会白白占用broker的文件句柄和内存。这里分享我的经验法则,供大家参考。
第一看吞吐量。单个分区在Kafka里顺序写,SSD环境下单分区吞吐几百MB/s不是问题,但考虑到消费者处理能力,一般按单分区每秒处理多少条消息来反推。比如业务高峰期每秒需要消费10万条消息,单分区每秒能消费2万条,那就至少需要5个分区,再留一些余量定到8个。
第二看消费者并发度。一个分区只能被同一消费组里的一个消费者实例消费,所以分区数至少不能低于消费者实例数量,否则会有消费者被闲置。反过来,如果消费者实例数大于分区数,多余的实例也闲着,所以两者尽量匹配。
第三看尽可能避免topic不可用。分区数越多,Controller的元数据管理压力越大,单broker上的分区数也不宜过多。经验值是单个broker处理几百个分区没问题,上千个就要谨慎,最好通过监控观察文件句柄和内存占用。
6.2 代码管理中必须养成的三个习惯
把创建主题做成一个标准操作之后,有三个习惯建议坚持下来。
第一个习惯是创建前检查。不管是脚本还是API,创建前都先查一下主题是否存在,避免重复创建导致审批流和监控告警混乱。第二个习惯是显式设置超时。AdminClient的异步操作一定要给get()设置超时时间,否则故障时线程会一直挂着。我之前就在一次线上事故中发现部分线程卡在createTopics()上,最后定位就是没设超时。第三个习惯是创建后验证。用describeTopics()确认分区数和副本因子符合预期,再对外暴露使用。
这三个习惯都很简单,但真正常年坚持下来,线上主题相关的杂事会少很多。
6.3 我个人的一些体会
从这个系列开始到现在,主题创建是我觉得最适合用来串起Kafka核心概念的一个入口。它把客户端编程、Controller协调、元数据存储、副本分配、元数据同步这些零散知识点,全部串成了一条完整的链路。搞懂这条链路后,再看其他Kafka问题,比如某个分区没有leader、消费组rebalance频繁、集群元数据不一致,都会比之前清晰很多。
我最早学Kafka的时候,也是从命令行敲kafka-topics.sh开始的,当时觉得不就是建个主题嘛,没什么技术含量。直到后来在生产环境排查一个“主题创建后分区数据写入失败”的问题,才意识到这条链路的复杂程度。现在回头看,如果当时能尽早把主题创建背后的源码和流程读一遍,后面会少走很多弯路。
系列后续我计划聊聊生产者和消费者底层的工作机制,以及Kafka在监控运维方面的一些实战踩坑记录。如果你在根据这篇内容实操时遇到什么奇怪的问题,欢迎按文中的链路逐层排查,大概率能找到问题所在。