Kafka 生产环境部署与调优

Kafka 生产环境实战 📖 前置阅读:本文是 Kafka 系列的终篇,假设读者已经掌握前五篇的全部内容(核心架构、SpringBoot 集成、Producer 深入、Consumer 深入、Kafka Streams)。 一、⚡ 问题切入:单节点 Docker 的瓶颈 第一篇搭的单节点 KRaft Kafka(一个 Controller + Broker 合体进程)只能用来学习——生产环境中: 单点 后果 唯一的 Broker 挂了 所有消息不可发送/消费——整个 Kafka 瘫痪 无副本 Broker 磁盘损坏 → 消息永久丢失 JVM 内存不足 Full GC 频繁 → 消息延迟抖动 → Producer 超时失败 磁盘写满 Partition 无法写入 → Producer 阻塞或报错 生产最低配:3 台 Broker + 3 台 Controller(或 3 台合体节点,controller + broker 混合模式),每个 Topic 至少 2 副本。 二、KRaft 三节点集群搭建 2.1 架构设计 第一篇用的 KRaft 模式是单节点(Controller 和 Broker 运行在一个进程中)。生产环境拆开: ...

十一月 18, 2022 · 8 分钟 · 1529 字 · yaomingye

Kafka Streams 与高级特性

Kafka Streams 流处理实战 📖 前置阅读:本文假设读者已掌握 Kafka Consumer/Producer 的使用和 Offset/Partition 概念。如果还不熟悉,建议先阅读 Consumer 深入:位移管理与 Rebalance。 一、⚡ 问题切入:用 Consumer + Producer 写流处理有什么毛病? 假设需要统计每个商品的近 5 分钟销量。用 Consumer + Producer 写: // 用 Consumer + Producer 实现滑动窗口计数——代码量爆炸 Map<String, List<Long>> windowCache = new HashMap<>(); // 还需要处理:窗口过期清理、状态持久化、故障恢复、乱序数据... 这个需求正是流处理引擎的用武之地。Kafka Streams 是 Kafka 官方的流处理库——它不是另一个需要部署的服务(不像 Flink/Spark Streaming),而是一个 Java 库,跑在你的应用进程里。 Kafka Streams 本质:Consume → 计算 → Produce,全部走 Kafka,中间状态存在 Kafka 的本地 RocksDB 实例中。 flowchart LR classDef startEnd fill:#701a4c,stroke:#e11d48,stroke-width:2px,color:#fce7f3,font-weight:bold; classDef process fill:#1e1e24,stroke:#6b7280,stroke-width:1.5px,color:#e5e7eb; classDef data fill:#052e16,stroke:#16a34a,stroke-width:1.5px,color:#bbf7d0,font-weight:bold; classDef highlight fill:#450a0a,stroke:#dc2626,stroke-width:1.5px,color:#fecaca,font-weight:bold; INPUT[("input-topic\n原始数据")] -->|"consume"| KS[Kafka Streams\n计算引擎\n+ RocksDB 本地状态] KS -->|"produce"| OUTPUT[("output-topic\n处理结果")] KS -->|"备份"| CHANGELOG[("changelog-topic\n状态变更日志")] class INPUT,OUTPUT data; class CHANGELOG process; class KS highlight; 二、Kafka Streams 核心概念 2.1 KStream vs KTable Kafka Streams 有两个核心抽象,理解它们的区别是正确使用的前提: ...

十一月 17, 2022 · 6 分钟 · 1176 字 · yaomingye

Kafka Consumer 深入:位移管理与 Rebalance

Kafka Consumer 深入 📖 前置阅读:本文假设读者已掌握 SpringBoot Kafka 的基本消费操作(@KafkaListener)。如果还不熟悉,建议先阅读 SpringBoot Kafka 全操作指南。 一、⚡ 问题切入:消费者重启后,怎么知道上次读到哪了? RabbitMQ 的答案是"消息消费后就删了,不需要记位置"。RocketMQ 的答案是"Broker 帮你记 offset"。Kafka 的答案是——消费者自己记,记在一个叫 __consumer_offsets 的内部 Topic 里: 消费者在 Partition-2 上消费到 offset=1500 ↓ 提交 offset __consumer_offsets Topic: Key: (order-consumer-group, order-topic, 2) Value: offset=1500 消费者重启 ↓ ↓ 读取 offset 从 offset=1501 继续消费 这个设计是 Kafka 和 RabbitMQ/RocketMQ 最核心的消费端差异——Kafka 的消费者对自己的消费进度负全责。如果消费者忘记提交 offset,重启后就会从上次提交的位置重新消费,产生重复消息。 二、Offset 提交机制 2.1 自动提交 vs 手动提交 Kafka 提供了两种 Offset 提交方式: 提交方式 配置 行为 风险 自动提交 enable-auto-commit: true 每隔 auto.commit.interval.ms(默认 5s)自动提交 poll 返回的最大 offset 消息可能没处理完就提交了——进程挂了会丢消息 手动提交 enable-auto-commit: false + ack-mode: manual 消费者处理完消息后显式调用 ack.acknowledge() 消息可能处理完了但没提交——重启后重复消费 flowchart TD classDef startEnd fill:#701a4c,stroke:#e11d48,stroke-width:2px,color:#fce7f3,font-weight:bold; classDef condition fill:#2a1147,stroke:#a855f7,stroke-width:1.5px,color:#ede9fe,font-weight:bold; classDef process fill:#1e1e24,stroke:#6b7280,stroke-width:1.5px,color:#e5e7eb; classDef highlight fill:#450a0a,stroke:#dc2626,stroke-width:1.5px,color:#fecaca,font-weight:bold; POLL([poll 拉取消息]) --> PROCESS[处理消息] PROCESS --> MODE{提交模式} MODE -- "自动提交" --> AUTO["每隔 auto.commit.interval.ms\n自动提交最后一次 poll 的 offset"] MODE -- "手动提交" --> MANUAL["业务处理成功后\n显式调用 ack.acknowledge()"] AUTO --> RISK1["风险:消息还没处理完\n但 offset 已提交\n→ 进程挂了丢消息"] MANUAL --> RISK2["风险:消息已处理完\n但 offset 没提交\n→ 重启后重复消费"] class POLL startEnd; class MODE condition; class AUTO,MANUAL highlight; class RISK1,RISK2 process; 手动提交比自动提交更安全——至少你知道什么时候提交了。消息重复消费可以用幂等解决,但消息丢失无法恢复。 ...

十一月 16, 2022 · 6 分钟 · 1173 字 · yaomingye

Kafka Producer 深入:分区、ACK 与幂等

Kafka Producer 深入 📖 前置阅读:本文假设读者已掌握 SpringBoot Kafka 的基本发送操作(KafkaTemplate.send)。如果还不熟悉,建议先阅读 SpringBoot Kafka 全操作指南。 一、⚡ 问题切入:消息到底发到了哪个 Partition? 上一篇用 kafkaTemplate.send("order-topic", key, msg) 发送消息时,生产者背后发生了三件事: Partitioner 决定消息进入哪个 Partition 消息攒批——在内存 Buffer 中等待 batch.size 或 linger.ms 条件触发 根据 acks 配置决定什么时候认为发送成功 当你看到日志里 partition=1, offset=0 时,背后是这三个步骤的协作。每个步骤都有配置项可以调整——它们直接影响消息顺序、可靠性、吞吐量。 二、分区策略 —— 消息路由的第一个环节 2.1 默认分区策略 Kafka Producer 的 Partitioner 接口决定了每条消息进入哪个 Partition。默认实现是 DefaultPartitioner: // DefaultPartitioner 的逻辑(简化版) public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) { int numPartitions = cluster.partitionsForTopic(topic).size(); if (keyBytes == null) { // Key 为 null → 使用 Sticky 分区(粘性分区) // 不是轮询!是把一批消息都发到同一个 Partition,等 batch 满了才换下一个 return stickyPartition(topic, numPartitions); } else { // Key 不为 null → 用 murmur2 哈希 % Partition 数量 return Utils.murmur2(keyBytes) % numPartitions; } } 规则一:Key 为 null → Sticky Partition。Kafka 2.4 之前是轮询(Round Robin——每条消息换一个 Partition),2.4+ 改为 Sticky——把一批消息"粘"在同一个 Partition 上,等这个 Batch 满了或时间到了才切到下一个 Partition。这减少了网络请求次数——把同一批的消息打成一个请求发给同一个 Broker。 ...

十一月 15, 2022 · 6 分钟 · 1274 字 · yaomingye

SpringBoot Kafka 全操作指南

SpringBoot Kafka 实战 📖 前置阅读:本文假设读者已理解 Kafka 的核心概念(Broker、Topic、Partition、ConsumerGroup、Offset)。如果还不熟悉,建议先阅读 Kafka 核心架构与日志存储模型。 🎯 第一步:目标说明 上一篇用原版 Kafka Java Client 写了 KafkaProducer + KafkaConsumer。和 RabbitMQ、RocketMQ 一样——真实的 SpringBoot 项目里不需要那些样板代码。spring-kafka 帮我们处理了连接管理、Producer 生命周期、Consumer 线程池、Offset 提交。 读完这篇会掌握: KafkaTemplate 三种发送方式(同步/异步/回调) @KafkaListener 注解消费——单条和批量 JSON 序列化全链路配置——Producer 端 JsonSerializer + Consumer 端 JsonDeserializer Producer 配置:acks、retries、batch.size、linger.ms、compression.type Consumer 配置:group.id、auto.offset.reset、enable.auto.commit、max.poll.records 📋 第二步:前置条件 前置项 具体要求 验证命令 JDK 17+(8+ 也兼容) java -version SpringBoot 3.x(文中用 3.2) mvn dependency:tree | grep spring-boot Kafka 3.7.0 KRaft 模式(单节点即可) docker ps | grep kafka 前置知识 Broker/Topic/Partition/ConsumerGroup/Offset 概念 — 确认 Kafka 在跑: ...

十一月 14, 2022 · 8 分钟 · 1619 字 · yaomingye

Kafka 核心架构与日志存储模型

Kafka:分布式提交日志,不是消息队列 📖 前置阅读:本文假设读者已理解消息队列的基本价值(异步、解耦、削峰填谷),最好读过 RabbitMQ 或 RocketMQ 的任意一篇基础文章。有了 MQ 概念再学 Kafka 事半功倍。 一、⚡ 问题切入:RabbitMQ 和 RocketMQ 有什么共同的"毛病"? 先回顾 RabbitMQ 和 RocketMQ 的消费模型: RabbitMQ: Consumer 收到消息 → 手动 basicAck → Broker 删除消息 RocketMQ: Consumer 收到消息 → 返回 CONSUME_SUCCESS → offset 推进 共同点:消息被消费者确认(ACK)后,Broker 就把它删了。消息在 Broker 上的生命周期是"暂存"——它存在只是为了等待消费者拿走。 这个模型有一个隐含的限制:一条消息只能被消费一次。想重放消息?RabbitMQ 做不到(消息已经删了),RocketMQ 可以重置 offset 但受 CommitLog 保留时间限制。 这时候再看 Kafka 的设计:消息消费后不删除。消息存在磁盘上,按时间或大小策略统一过期,消费者想从哪个位置读就从哪个位置读。 Kafka: Consumer 自己管 offset,随时可以回到过去的某个位置重读 消息不是被消费掉的——是按时间自然过期的 这就是 Kafka 和 RabbitMQ/RocketMQ 本质上的不同——Kafka 不是一个消息队列,它是一个分布式提交日志(Distributed Commit Log)。 二、🧬 Kafka 是什么:分布式提交日志 2.1 核心定义 Kafka 的官方定位:分布式、分区化、多副本的提交日志服务。 ...

十一月 13, 2022 · 5 分钟 · 1020 字 · yaomingye
Cat Radio