尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Kafka主题创建全链路解析:从API调用到分区副本分配与元数据同步

Kafka主题创建全链路解析:从API调用到分区副本分配与元数据同步 这个系列写到第九篇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到CreateTopicsRequest2.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 不考虑机架的分配策略RackUnawareRackUnaware的算法核心很简单把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 机架感知分配策略RackAwareRackAware的算法比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.idrack-a。第二创建主题时要么使用新版Admin接口自动识别broker机架信息要么用命令行时确保分配算法走的是机架感知逻辑。在KRaft模式下机架信息也在元数据里注册逻辑类似。我踩过一次坑接了机架感知的broker配置但线上主题创建用的还是自己写的老脚本结果分配结果完全没体现机架隔离。后来排查发现脚本里有一行代码强制指定了replicasAssignments自定义分配直接绕过了自动机架算法。所以这里引出一个重要的点只要创建时手动指定了副本分布任何自动分配策略都不会生效。3.4 自定义分配手动指定副本分布除了自动分配Kafka还允许完全由用户指定每个分区的副本分配方案。在Admin接口里通过如下方式传入MapInteger, ListInteger 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的集群你建主题时写了副本因子3Controller会直接报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请求先发给任意一个brokerbroker根据管理类请求路由规则把请求转发给当前ControllerController进行合法性校验校验失败直接返回错误Controller计算分区副本分配方案自动或根据自定义分配Controller将元数据写入ZooKeeper或KRaft日志Controller向所有broker广播UpdateMetadataRequest各broker更新本地元数据缓存响应ControllerController把创建成功结果返回给AdminClientAdminClient通过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失联导致“可用”数偏小。我建议生产环境副本因子统一设3broker数量至少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在监控运维方面的一些实战踩坑记录。如果你在根据这篇内容实操时遇到什么奇怪的问题欢迎按文中的链路逐层排查大概率能找到问题所在。
返回列表