拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Kafka核心原理与实战:分区、副本、消息可靠性及面试要点

先说点实际的。Kafka 这东西几乎所有做后端或者数据平台的人早晚都会碰到。不管你是搞日志采集、用户行为追踪、消息削峰还是做实时数仓Kafka 几乎成了事实标准。我最早接触它的时候还是因为业务量上来之后ActiveMQ 频繁卡死才被迫换过来的。那时候资料远没有现在多踩了不少坑走了不少弯路。这篇文章我尽量把 Kafka 的核心原理、安装部署、实战操作和那些面试里经常被问到的点一次讲明白。内容会比较长你可以当成一份手边的参考手册来用。1. 为什么大家都在用 Kafka它到底解决了什么问题在深入技术细节之前先想清楚一个根本问题我们为什么需要 Kafka1.1 从一次系统重构说起我以前维护过一个电商后台的老系统。用户下单之后系统要同步干三件事写订单库、扣库存、发短信通知。高峰期一到大促数据库连接池直接被打满短信接口超时订单丢失用户疯狂投诉。这其实就是典型的耦合问题。后来引入 Kafka 之后架构变成了这样下单接口只干一件事就是把订单消息丢进 Kafka然后立刻返回“下单成功”。至于扣库存、发短信这些操作都由下游的消费者服务自己去 Kafka 里拉取消息异步处理。效果立竿见影接口响应时间从 2000 毫秒降到了 50 毫秒左右数据库压力大减。这就是 Kafka 最核心的价值解耦和异步削峰。1.2 Kafka 的定位分布式消息流平台很多人以为 Kafka 只是一个消息队列其实它更像是一个分布式流处理平台具备三个关键能力发布与订阅像传统的消息队列一样可以广播消息给多个消费者。持久化存储消息被写入磁盘可以保留一段时间消费者可以随时回放数据而不是像 RabbitMQ 那样消费完就删除。这一点让它成了天然的“数据管道中枢”。流式处理基于 Kafka Streams API可以对实时数据流进行聚合、连接、窗口计算等复杂处理相当于内置了一个轻量级流计算引擎。1.3 核心组件的基本盘Kafka 的架构里有几个角色必须先有概念后面聊起来才不费劲Producer生产者发送消息的一方负责把数据发布到指定的 Topic。Consumer消费者拉取消息的一方负责订阅并处理 Topic 中的数据。Broker代理服务器Kafka 集群中的一台服务器节点负责存储消息和处理客户端请求。一个集群通常由多个 Broker 组成。Topic主题消息分类的逻辑名称一条消息必须归属于某个 Topic。Partition分区Topic 的物理拆分单位。一个 Topic 可以分成多个分区每个分区是一个有序的日志文件。分区是 Kafka 实现高吞吐和水平扩展的关键。Consumer Group消费者组一组消费者的集合。组内的每个消费者负责消费不同分区的数据实现负载均衡。ZooKeeper / KRaft负责集群的元数据管理、Broker 选举等协调工作。新版本2.8已经在尝试用 KRaft 协议替代 ZooKeeper去掉一个外部依赖。2. Kafka 核心原理与架构解析分区、副本与消费模型这一部分比较烧脑但却是整个 Kafka 的精华所在。理解了原理你才能在实际问题里知道怎么调参才不会被面试官问倒。2.1 分区机制高吞吐的秘密武器为什么 Kafka 能单机支持几十万甚至上百万的每秒吞吐量核心就是分区。一个 Topic 被拆分成多个 Partition每个 Partition 在物理上对应一个文件夹里面是一段段有序的日志文件Segment。生产者的消息写入时会根据 key 的哈希值或轮询策略决定写入哪个 Partition。消费时消费者组里的每个消费者会被分配一部分分区并行拉取。我打个比方一个快递仓库Topic有 10 个收件窗口Partition每个窗口只有一个队列有序日志。如果只有一个窗口快递员只能排队一件一件交件吞吐低有了 10 个窗口可以同时收 10 个快递员的件并行度高吞吐自然就上去了。Topic: orders ├── partition-0 (文件夹: orders-0) ├── partition-1 (文件夹: orders-1) └── partition-2 (文件夹: orders-2)注意一个经典的坑分区数是不能随便缩小的Kafka 官方不允许降低分区数因为消息在分区中的路由算法决定了缩容会导致数据错乱。所以架构初期分区的规划要留有一定的扩展余量。2.2 副本机制保证数据不丢的保险如果只有一个 Broker它挂掉了那整个 Topic 的数据就全没了。所以 Kafka 为每个 Partition 配置了多个副本。副本分为两类Leader 副本负责处理生产者和消费者的读写请求。Follower 副本只负责从 Leader 同步数据不对外提供服务。一旦 Leader 挂掉Follower 中会有新的副本被选举为 Leader。这里有个很重要的参数acks。它决定了生产者发送消息后的确认机制acks 0Producer 发送后不等待任何确认。吞吐最高但消息可能丢失。acks 1Leader 写入成功后即返回确认。默认值效率与可靠性之间的折中。如果 Leader 在 Follower 同步完成前宕机可能丢数据。acks -1allLeader 和所有 ISRIn-Sync Replicas中的 Follower 都写入成功后才返回确认。数据最安全但延迟最高。注意acksall也不是 100% 不丢它只能保证 ISR 中有存活副本时数据内存不丢。极端情况下所有副本都宕机且磁盘损坏物理丢数据是任何软件都避免不了的这时需要考虑跨机房容灾。2.3 消费者组与消息只被消费一次的模型消费者组是 Kafka 实现“一对多”和“负载均衡”的关键。同一个 Topic 的消息可以被多个消费者组各自消费一次。比如订单消息交易组消费它来发短信风控组消费它来检测风险互不干扰。但在同一个消费者组内一条消息只能被组内的一个消费者实例消费。当组内消费者数量发生变化上线、宕机、扩容时Kafka 会触发Rebalance重平衡重新分配分区。Rewalance 是好多新手调试时容易懵的地方。我举个例子Topic 有 4 个分区消费者组里有 3 个消费者实例C1、C2、C3。那么分配可能是 C1 消费 P0、P1C2 消费 P2C3 消费 P3。如果 C3 挂掉Rebalance 后C1 和 C2 会接管 P3 的消息。这里必须提醒一个容易踩坑的点Consumer 实例数最好不超过分区数。如果 4 个分区你起了 5 个 Consumer那么第五个 Consumer 会被闲置永远收不到消息。这是新手排查“怎么我的消费者没消费到数据”时第一件要检查的事。2.4 消息存储与顺序性保证每个 Partition 内部的消息是严格有序的追加写入。但是 Partition 之间没有全局顺序。如果一个业务场景严格要求全局有序比如金融转账的流水那最好把这类消息都发送到同一个 Partition通过 key 哈希比如订单 ID。顺序性在重试场景下尤其要注意。生产者开了重试机制后如果第一批消息写入失败重试而第二批消息已经写入成功最终顺序就会颠倒。Kafka 解决方式是设置max.in.flight.requests.per.connection为 1保证发送中的请求只有一个并配合enable.idempotence开启幂等。3. Kafka 环境安装与集群部署实战原理讲再多不如动手跑一遍。我先从最常见的单机部署讲起再讲集群部署的关键配置。3.1 单机快速体验从下载到跑通Kafka 官方推荐的方式是直接下载二进制包。这里我用目前最普及的 2.13-3.4.0 版本基于 Scala 2.13Kafka 3.4.0举例。# 1. 下载并解压 wget https://archive.apache.org/dist/kafka/3.4.0/kafka_2.13-3.4.0.tgz tar -xzf kafka_2.13-3.4.0.tgz cd kafka_2.13-3.4.0 # 2. 修改配置KRaft 模式3.x 版本已经内置 KRaft可以不装 ZooKeeper # 先生成集群唯一 ID KAFKA_CLUSTER_ID$(bin/kafka-storage.sh random-uuid) # 3. 格式化存储目录 bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties # 4. 启动 Kafka bin/kafka-server-start.sh config/kraft/server.properties这里有个小细节值得说明如果你用的是 2.x 版本还需要额外启动 ZooKeeper。启动命令是bin/zookeeper-server-start.sh config/zookeeper.properties然后才能启动 Kafka server。很多老教程还在讲这种方式但新项目我建议直接上 KRaft少维护一个组件。3.2 集群部署3 节点配置示例生产环境至少三台机器以便容忍单点故障。以 3 个 Broker 为例每台服务器上的config/server.properties关键配置如下# 每个节点的唯一 ID三台分别为 0、1、2 broker.id0 # 监听地址写清楚内网 IP 和端口 listenersPLAINTEXT://192.168.1.10:9092 advertised.listenersPLAINTEXT://192.168.1.10:9092 # 日志存储目录建议数据盘 log.dirs/data/kafka-logs # ZooKeeper 地址如果是 KRaft 模式则换成 controller 配置 zookeeper.connect192.168.1.10:2181,192.168.1.11:2181,192.168.1.12:2181启动第一台后依次启动另外两台。然后用命令查看集群状态bin/kafka-topics.sh --describe --bootstrap-server 192.168.1.10:9092 --topic test-topic如果看到每个 Partition 都有多个副本均匀分布在不同的 Broker 上集群就部署成功了。3.3 Windows Docker 部署路径日常开发中很多场景是在 Windows 笔记本上。如果你不想在 Windows 上折腾 Java 环境变量用 Docker 最省心。version: 3 services: kafka: image: bitnami/kafka:3.4 container_name: kafka ports: - 9092:9092 environment: - KAFKA_CFG_NODE_ID0 - KAFKA_CFG_PROCESS_ROLEScontroller,broker - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS0kafka:9093 - KAFKA_CFG_LISTENERSPLAINTEXT://:9092,CONTROLLER://:9093 - KAFKA_CFG_ADVERTISED_LISTENERSPLAINTEXT://localhost:9092 - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAPCONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT - KAFKA_CFG_CONTROLLER_LISTENER_NAMESCONTROLLER - KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLEtrue volumes: - kafka_data:/bitnami/kafka volumes: kafka_data:保存为docker-compose.yml运行docker-compose up -d。然后可以用 Kafdrop 或 Offset Explorer 之类的可视化工具在浏览器里直接查看消息。3.4 集群部署成败的几个关键点集群部署的坑比单机多列几个我实际踩过的advertised.listeners必须配如果只配了listeners客户端比如别的机器上的消费者会拿到 Broker 的内网 hostname可能解析不了。务必用可路由的 IP 或域名。JVM 堆内存别贪心Kafka 默认的KAFKA_HEAP_OPTS是 1G生产环境建议设为 4G~6G但不要超过物理内存的一半。因为 Kafka 大量使用操作系统的 PageCache 缓存数据堆内存给多了反而浪费页缓存空间。文件句柄数要放开ulimit -n至少设到 100000否则高负载下 Broker 会报 “Too many open files”。磁盘选择Kafka 是顺序写机械硬盘也能用但延迟较高能用 SSD 尽量用 SSD。log.dirs可以配置多个目录逗号分隔Kafka 会自动做磁盘间的均衡。4. 生产消费全流程实操与命令行调优集群起来了topic 建好了接着就是最核心的实操怎么生产消息、怎么消费消息、怎么排错。4.1 命令行快速验证 Topic 中的消息上手验证 Kafka 最直观的就是命令行工具。手动创建 topic# 创建名为 test-events 的 topic3 个分区1 个副本 bin/kafka-topics.sh --bootstrap-server localhost:9092 \ --create --topic test-events \ --partitions 3 --replication-factor 1启动一个生产者手动输入几条消息bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic test-events hello kafka this is a test message done另开一个终端启动消费者查看数据bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic test-events --from-beginning你会在消费者终端看到刚才那三条消息。这个过程中有几个非常容易困惑的问题问题 1生产者启动一次会一直运行吗很多人搜热词的时候会问“kafka生产消费命令启动一次会一直运行吗”答案是看命令。kafka-console-producer.sh启动后一直等待你输入不会自动退出。kafka-console-consumer.sh启动后也是一直挂着监听新消息。这是因为它们本身就是常驻前台进程。想要退出按Ctrl C。但在代码世界里程序会主动关闭生产者或消费者比如批量发完消息后调用producer.close()或者消费者在拉取到足够数据后调用consumer.close()并退出循环。问题 2如何查看 Topic 中已有的数据除了用消费者--from-beginning还可以用kafka-dump-log.sh直接看日志文件内容# 找到该 topic 分区的日志目录和文件 bin/kafka-dump-log.sh --files /tmp/kafka-logs/test-events-0/00000000000000000000.log这个命令会把二进制日志解析成可读的偏移量、时间戳、key、value 等。排查数据是否落盘时非常有用。4.2 生产端关键参数与计算生产端配置直接影响整体吞吐和延迟。我常用的模板如下Properties props new Properties(); props.put(bootstrap.servers, 192.168.1.10:9092,192.168.1.11:9092); props.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(value.serializer, org.apache.kafka.common.serialization.StringSerializer); // 指定 key 用于决定分区保证相同 key 的消息进入同一分区 props.put(acks, all); props.put(retries, 3); props.put(batch.size, 16384); // 16KB props.put(linger.ms, 5); // 最多等待 5ms 凑批 props.put(buffer.memory, 33554432); // 32MB props.put(max.in.flight.requests.per.connection, 5); props.put(enable.idempotence, true); KafkaProducerString, String producer new KafkaProducer(props); for (int i 0; i 1000; i) { producer.send(new ProducerRecord(test-events, key- (i % 3), value- i)); } producer.close();几个参数的选择逻辑acksall配合enable.idempotencetrue是生产环境最稳妥的组合。幂等之后Broker 会为生产者分配一个 PID 并维护序列号重复消息会被去重。linger.ms和batch.size是一对“双刃剑”。想提高吞吐就增大 batch想降低延迟就减小 linger。追求低延迟的场景如实时推荐可以把linger.ms设为 0这样每条消息立刻发送但吞吐会下降一些。buffer.memory是生产者内存中缓存消息的总大小。如果发送速度超过 Broker 接收速度这个缓冲区会被写满然后send()会阻塞直到超时。4.3 消费端核心参数与位移管理消费端最容易出问题的就是位移offset。先梳理清楚概念每个 Partition 有一个 offset 字段表示 Consumer Group 下一条要读取消息的位置。Consumer 成功处理一条消息后需要提交commit位移否则下次重启还会消费到旧数据。Properties props new Properties(); props.put(bootstrap.servers, 192.168.1.10:9092); props.put(group.id, order-service); props.put(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(value.deserializer, org.apache.kafka.common.serialization.StringDeserializer); // earliest: 从头开始消费latest: 只消费新消息none: 没有位移就报错 props.put(auto.offset.reset, earliest); // 是否自动提交位移默认 true props.put(enable.auto.commit, true); props.put(auto.commit.interval.ms, 1000); KafkaConsumerString, String consumer new KafkaConsumer(props); consumer.subscribe(Arrays.asList(test-events)); while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { System.out.printf(offset %d, key %s, value %s%n, record.offset(), record.key(), record.value()); } }在实际项目中我一般建议把enable.auto.commit设为false改为手动提交。因为业务处理失败时不希望提交位移否则消息会“悄悄丢失”。手动提交的方式// 等一批消息处理完成后再提交本次 poll 的最大位移 consumer.commitSync();但这里有个问题如果你的业务逻辑需要处理 10 条消息处理到第 5 条时失败了用commitSync()提交的是整批 offset重启后还是会从该 offset 消费等于把已经成功的 5 条又重新消费了一遍。想精确控制可以用手动提交特定 offset// 处理完所有消息后逐条或按分区提交 for (TopicPartition partition : records.partitions()) { ListConsumerRecordString, String partitionRecords records.records(partition); long lastOffset partitionRecords.get(partitionRecords.size() - 1).offset(); consumer.commitSync(Collections.singletonMap(partition, new OffsetAndMetadata(lastOffset 1))); }把已处理的最大 offset 加 1 提交是官方推荐的“at least once”语义标准写法。这是最老实的做法——接受可能的重复消费但在业务侧做幂等。4.4 消费者再均衡监听器再均衡发生时有可能会出现一部分消费任务被中断位移没提交。此时配合监听器能优雅处理consumer.subscribe(Arrays.asList(test-events), new ConsumerRebalanceListener() { Override public void onPartitionsRevoked(CollectionTopicPartition partitions) { // 分区被回收前提交当前处理进度避免重复消费 consumer.commitSync(); } Override public void onPartitionsAssigned(CollectionTopicPartition partitions) { // 新分区分配完成后做一些初始化工作如定位 offset } });这个 listener 在消费者实例上下线、分区变化时尤其重要。我在一个实时计算项目里就是因为在 Rebalance 期间没提交位移导致每次发版上线都重复处理了一大批数据后来加上这个监听器才解决。5. 消息不丢失与不重复的全链路保障“数据重复”和“数据丢失”是搜索热词里明显的两个痛点也是 Kafka 面试里最高频的两类问题。5.1 三个环节的丢数据风险Kafka 的消息要从 Producer 到 Broker再到 Consumer。任何一环都可能丢数据环节丢数据原因核心对策Producer 到 Broker网络抖动、Broker 写入失败、重试未开启acksallretries0开启幂等Broker 存储副本数不足、Leader 切换时 Follower 落后于 Leaderreplication.factor3min.insync.replicas2unclean.leader.election.enablefalseConsumer 消费自动提交位移但业务代码处理失败手动提交位移处理完再提交消费逻辑幂等针对中间环节我会在集群里做这样的配置# Broker 端配置 default.replication.factor3 min.insync.replicas2 unclean.leader.election.enablefalsemin.insync.replicas2的含义是当 acksall 时Broker 至少要有 2 个副本同步成功才算写入成功。这能防止 Leader 写完后只同步给自己就返回成功的情况。但要注意如果集群只剩 1 个副本存活此时生产会报NotEnoughReplicasException这是可以接受的因为我们要的是不丢数据而不仅仅是可用性。5.2 重复消费的克星幂等机制Kafka 的幂等机制主要防止 Producer 重复发送导致的重复数据通过 PID 和序列号实现。但 Consumer 拉取消息后的重复消费Kafka 本身管不了需要业务端自己想办法。我总结一个简单的幂等方案数据库唯一键去重消费时把业务 ID 写入带唯一索引的表中插入冲突就说明已处理过直接跳过。Redis 幂等表用 SETNX 命令key 是消息的唯一 IDvalue 是处理结果。存在就跳过。状态机校验比如订单状态已经从 READY 变成 PAID就不能再执行一次支付回调。这里有一个经常被误解的点at-least-once vs exactly-once。Kafka 支持 exactly-once但通常指的是 Kafka Streams 内部计算或者 Producer 到 Broker 这一段。对于外部系统如数据库、Redis、第三方 APIKafka 的 exactly-once 并不能保证外部操作只执行一次你必须自己做好幂等。5.3 OOM 与消息积压排查搜索热词里出现了 “kafka oom” 和 “kafka消息延迟高”这两个问题往往关联在一起。OOM内存溢出常见原因有两个生产者buffer.memory过小队列积压后内存暴涨。消费者单次poll拉取的消息条数过多反序列化后撑爆堆内存。典型案例某业务设置的max.partition.fetch.bytes很大默认 1MB如果 Topic 的消息单条很大比如 5MB拉取 10 条就是 50MB很快就把堆内存打爆了。解决方式是合理设置fetch.max.bytes和max.partition.fetch.bytes并确保堆内存足够。消息延迟高排查路径通常是看消费者 Lagbin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group order-service --describe看 LAG 列。LAG 持续增长说明消费速度跟不上生产速度。优化方向增加分区数、增加消费者实例数、优化消费逻辑比如批量写库、排查下游依赖瓶颈。LAG 不增长但端到端有延迟可能是linger.ms设置得过长或者 Broker 端磁盘 I/O 过慢导致消息滞留。我遇到过一次很诡异的情况消费者没报错Lag 也不大但用户消息就是延迟了十几秒。后来发现是消费者在处理每条消息时调用了外部 HTTP 接口这个接口超时时间设置得很长15 秒并发再高也被阻塞住了。后来改成线程池异步推送延迟立刻降下来了。6. 生产环境监控与常见问题速查Kafka 集群跑起来只是开始日常运维和监控才是大头。6.1 监控哪些指标监控不是越多越好抓到关键指标才能快速定位问题。我日常重点看以下几项Broker 端CPUKafka 会用到不少 CPU 做压缩解压、磁盘使用率日志保留时间别设太长、网络吞吐、GC 耗时注意 Full GC 会导致 Broker 长时间停顿。Topic 端消息生产速率、消费速率、分区 Leader 分布是否均衡。消费者端Consumer Lag延迟量是最核心的指标建议配到告警里超过阈值就通知。工具方面免费方案可以选择 Prometheus JMX Exporter Grafana社区里现成的 dashboard 模板很多。如果公司有条件用 Confluent Control Center 或者商业监控平台则更省心。6.2 常见异常速查表现象可能原因解决方法Connection refusedbroker 没启动 / 防火墙挡了端口 /advertised.listeners配置错误检查监听端口、telnet 测试、修正 listenersUnknownTopicOrPartitionException分区数不足、topic 不存在检查 topic 是否存在重建或扩分区OffsetOutOfRangeException消费者读取的位移已过期日志保留已删除根据业务设置auto.offset.resetearliest或调大日志保留时间NotLeaderOrFollowerException分区 Leader 在切换中客户端访问了旧 Leader重试即可客户端会自动刷新元数据RecordTooLargeException消息超过message.max.bytes限制调整 broker 的message.max.bytes和 producer 的max.request.size消息全部堆积在一个 Partitionkey 设置不当导致哈希不均匀改用轮询策略或对 key 加盐消费者频繁 Rebalance会话超时时间过短 / 消费者处理时间过长 / 心跳线程被阻塞调大session.timeout.ms、max.poll.interval.ms或优化消费逻辑6.3 Kafka 面试题高频点梳理因为热词里出现了“kafka面试题”我把这些年常被问到的题目统一整理一下Kafka 和 RabbitMQ / RocketMQ 的区别核心差异Kafka 吞吐量更高、天生为日志和流数据设计、消息可重复消费RabbitMQ 功能更丰富路由策略更灵活社区插件多RocketMQ 在金融场景下的事务消息更成熟。选型看场景没有绝对优劣。Kafka 为什么吞吐量高三点要答到顺序写磁盘、PageCache 机制、零拷贝sendfile。另外分区并行读写、批量发送与批量拉取也是加分项。如何保证 Kafka 消息的顺序性对指定 key 的消息发送到同一个分区分区内消费者线程数为 1重试时要避免同分区消息乱序。Kafka 中的 offset 存在哪里老版本存在 ZooKeeper 中新版本用__consumer_offsets内部主题保存而且这个 topic 默认有 50 个分区。消费者组 Rebalance 的触发条件有哪些消费者加入或退出订阅的 topic 分区数变化消费者处理消息超时导致被判定为死亡。讲一下 Kafka 的副本同步机制ISRISR 是与 Leader 保持同步的副本集合。Follower 通过 Fetcher 线程持续拉取 Leader 的消息并写入日志。如果 Follower 落后过多或长时间未拉取会被踢出 ISR。生产时acksall时只需 ISR 中的副本确认即可。这些都是考原理和细节的典型问题把前面每一章搞懂面试基本能贯通。7. 一个完整的实战案例订单事件流的收集与处理最后用我之前做过的订单事件流案例串起整篇文章的知识点。7.1 场景描述某外卖平台需要实时收集用户下单、接单、配送、完成的全流程事件用于监控订单成功率异常单自动告警统计各区域实时订单量事后回放事件用于业务复盘7.2 架构选型与 Topic 设计我们设计了三个 Topicorder-events核心订单领域事件包含订单 ID、用户 ID、商家 ID、状态、时间戳。delivery-location骑手实时位置流用于计算预计送达时间。alert-events异常告警事件。order-events的分区数设为 12副本数 3因为订单量峰值较高且需要跨 Broker 容灾。Topic 创建命令bin/kafka-topics.sh --bootstrap-server kafka-1:9092,kafka-2:9092,kafka-3:9092 \ --create --topic order-events \ --partitions 12 --replication-factor 3 \ --config retention.ms604800000其中retention.ms604800000表示消息保留 7 天。7.3 生产者实战代码订单服务在每个业务节点产生事件后异步发送Component public class OrderEventProducer { private final KafkaTemplateString, String kafkaTemplate; public void sendOrderCreated(Order order) { // 使用订单 ID 作为 key保证同一个订单的多个事件进入同一分区顺序可控 String eventJson objectMapper.writeValueAsString(order); kafkaTemplate.send(order-events, order.getOrderId(), eventJson); } }这里用订单 ID 做 key 有两个好处一是同一个订单的后续状态变化下单、接单、配送、完成都会进入同一个分区保持顺序二是方便下游按订单粒度聚合。7.4 消费者实战代码含手动位移下游分析服务对实时性要求高但允许少量重复所以使用手动异步提交、配合 Redis 幂等Component public class OrderEventConsumer { Autowired private StringRedisTemplate redisTemplate; KafkaListener(topics order-events, groupId order-analytics) public void onMessage(ConsumerRecordString, String record, Acknowledgment ack) { String key processed: record.topic() : record.partition() : record.offset(); Boolean firstProcess redisTemplate.opsForValue().setIfAbsent(key, 1, Duration.ofDays(1)); if (Boolean.TRUE.equals(firstProcess)) { // 正常的业务处理解析 JSON、更新实时指标、判断是否异常 processEvent(record.value()); ack.acknowledge(); // 手动提交位移 } else { // 已处理过跳过但也要提交位移 ack.acknowledge(); } } }这里的 RedissetIfAbsent就是一个轻量级幂等方案重复消费时直接跳过避免影响下游数据准确性。7.5 运行效果与问题复盘系统上线后高峰期每秒处理约 8000 条订单事件平均端到端延迟低于 300ms。但在压测阶段出现过一次问题现象消费者 Lag 持续上涨瞬时堆积达到 100 万条。排查过程先看消费者日志发现大量“Offset commit failed”再查发现下游的 Redis 写入超时阻塞了消费线程导致poll超时。解决下游批量写入减少 Redis 请求数同时给消费者单独配置线程池把同步写 Redis 改成异步批量写。这个案例最大的收获是Kafka 本身几乎不会成为瓶颈瓶颈往往在下游的消费处理逻辑。消费者代码写得不好再大的分区数也扛不住。8. 写在最后的实操心得如果只能选一条经验分享我会说Kafka 调优不是调一个参数而是一条链路。从 Kafka 为什么快、底层的 PageCache 怎么用的到生产端怎么配、消费端怎么防御重复与积压再到监控怎么观察延迟和 Lag每一步都要心里有数。另外生产环境里消息的可靠性和实时性往往是跷跷板。追求极致吞吐就要接受稍高的延迟追求数据零丢失就要牺牲一点可用性。没有放之四海皆准的配置只有贴合自己业务的取舍。每次升级版本前记得先看一眼官方升级文档Kafka 这两年版本迭代很快KRaft 模式也在逐步成熟。如果你的集群还在用 ZooKeeper可以开始规划平滑迁移了。最后再分享一个小技巧排查消费问题的时候kafka-consumer-groups.sh --describe是你最好的朋友先看 Lag再看日志最后才动代码。顺序反了很容易越查越乱。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门