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

RocketMQ 生产环境部署与调优

RocketMQ 生产部署 📖 前置阅读:本文是 RocketMQ 系列的终篇,假设读者已经掌握前五篇的全部内容(核心架构、SpringBoot 集成、高级消息类型、可靠性、消费者模式)。 一、⚡ 问题切入:单机 Docker 的瓶颈 第一篇搭的单机 RocketMQ(一个 NameServer + 一个 Broker)只能用来学习——生产环境中: 单点 后果 一台 NameServer 挂了 Producer/Consumer 无法获取路由——整个 MQ 瘫痪 一台 Broker 挂了 所有消息不可用——消息无法发送/消费 JVM 内存不足 Full GC 频繁 → 消息延迟抖动 → 超时重试雪崩 磁盘写满 CommitLog 无法写入 → 生产者阻塞 生产最低配:2 台 NameServer + 至少 2 台 Broker(主从)。 二、双主双从高可用集群搭建 2.1 架构设计 NameServer 集群:2 台(互不通信,各自独立) Broker 集群:2 组主从(Master-A + Slave-A、Master-B + Slave-B) NameServer-1 NameServer-2 (192.168.1.10:9876) (192.168.1.11:9876) ↑ ↑ ┌─────────┴───────┬───────────────┘ │ │ Broker-A (Master) Broker-A (Slave) 192.168.1.20:10911 192.168.1.21:10911 brokerId=0 brokerId=1 Broker-B (Master) Broker-B (Slave) 192.168.1.22:10911 192.168.1.23:10911 brokerId=0 brokerId=1 路由发现:Producer/Consumer 配置所有的 NameServer 地址——只要有一台 NameServer 活着,路由就能工作。 ...

十一月 12, 2022 · 5 分钟 · 909 字 · yaomingye

RocketMQ 消费者模式与过滤器

RocketMQ 消费者模式 📖 前置阅读:本文假设读者已掌握 SpringBoot RocketMQ 的基本消费操作。如果还不熟悉,建议先阅读 SpringBoot RocketMQ 全操作指南。 一、⚡ 问题切入:一条消息,谁来消费? 前面五篇的消费者代码都默认了一件事——一条消息只被一个消费者实例处理。但实际业务中: 订单消息——只能被一个实例消费(一个订单不能被两个服务处理两次) 配置刷新消息——所有实例都要收到(所有缓存节点刷新缓存) 部分 Tag 的消息——只关心"订单创建",对"订单支付"不感兴趣 消费不过来——10 个 Queue,2 个实例,怎么分? 这四大问题的答案都在这一篇里。 二、集群消费 vs 广播消费 2.1 集群消费(CLUSTERING)——默认模式 同一个 ConsumerGroup 内的所有实例共享消费一个 Topic 的消息——每条消息只被组内一个实例消费。 @Component @RocketMQMessageListener( topic = "order-topic", consumerGroup = "order-consumer-group", messageModel = MessageModel.CLUSTERING // 集群模式(默认值,可以不写) ) public class OrderClusteringListener implements RocketMQListener<OrderMessage> { @Override public void onMessage(OrderMessage msg) { System.out.printf("[实例A] 处理订单: orderId=%d%n", msg.getOrderId()); } } 部署 3 个实例,Topic 有 8 个 Queue: ...

十一月 11, 2022 · 4 分钟 · 730 字 · yaomingye

RocketMQ 消息可靠性与容错

RocketMQ 消息可靠性 📖 前置阅读:本文假设读者已掌握 SpringBoot RocketMQ 的发送和消费操作。如果还不熟悉,建议先阅读 SpringBoot RocketMQ 全操作指南。 一、⚡ 消息可能丢在哪? 和 RabbitMQ 一样,RocketMQ 的消息丢失也分三个环节: Producer → [网络] → Broker → [网络] → Consumer ① 发送丢失 ② Broker 宕机丢失 ③ 消费失败丢失 在 RocketMQ 中,生产者端没有 RabbitMQ 的 Publisher Confirm——替代方案是同步发送 + 返回值判断。Broker 端的持久化取决于刷盘策略和主从同步。消费端靠消费状态返回 + 重试 + 死信。 二、① 生产者端:同步发送判断返回值 RocketMQ 没有 RabbitMQ 的 ConfirmCallback 异步通知机制,但 syncSend 本身就是同步等待 Broker 确认——返回 SEND_OK 说明消息已写入 CommitLog: @Service public class ReliableProducer { @Autowired private RocketMQTemplate rocketMQTemplate; public void sendReliably(OrderMessage msg) { // syncSend 是同步阻塞的——返回 SEND_OK 才说明 Broker 已接收 SendResult result = rocketMQTemplate.syncSend( "order-topic:created", msg); if (SendStatus.SEND_OK.equals(result.getSendStatus())) { log.info("消息已确认到达 Broker: msgId={}", result.getMsgId()); } else { // 只有 SEND_OK 才认为成功——FLUSH_DISK_TIMEOUT 等状态说明刷盘超时 log.error("消息发送未确认: status={}", result.getSendStatus()); // 补偿:写入 DB 重试表 saveToRetryTable(msg); } } } SendStatus 的四种返回值: ...

十一月 10, 2022 · 5 分钟 · 853 字 · yaomingye

RocketMQ 顺序消息、延迟消息与事务消息

顺序消息、延迟消息与事务消息 📖 前置阅读:本文假设读者已掌握 SpringBoot RocketMQ 的基本操作(RocketMQTemplate、@RocketMQMessageListener)。如果还不熟悉,建议先阅读 SpringBoot RocketMQ 全操作指南。 一、⚡ 问题切入:三种 RabbitMQ 做不到或做不好的事 RabbitMQ 六篇系列学完时留了几个坑——有些场景 RabbitMQ 不是不能用,而是做起来别扭: 需求 RabbitMQ 方案 痛点 订单创建→支付→发货严格按序 单队列 + 单消费者,关并发 吞吐量压到一条线;一旦重试入队顺序全乱 30 分钟后自动取消 Delayed Message 插件或 TTL+DLX 插件生产不可靠;TTL+DLX 有消息时序问题 下单 + 扣库存 + 发消息三件事原子执行 自己实现本地消息表 + 定时补偿 代码量大,维护麻烦 RocketMQ 对这三种场景都有原生支持——不是插件,不是 workaround,是设计时就考虑进去了。 二、顺序消息 —— 深度篇 2.1 上一篇回顾 + 补充 上一篇讲了基本用法:syncSendOrderly 用 orderId 哈希选 Queue,同一个 orderId 进同一个 Queue → 该 Queue 内 FIFO。消费端 consumeMode = ConsumeMode.ORDERLY。 但这只讲了正常流程。重试会破坏顺序——这是最容易踩的坑。 2.2 顺序消费的重试机制:挂起而非重入队 并发消费中,失败的消息通过 RECONSUME_LATER 进入重试 Topic,然后延迟重新投递。但顺序消费不能这么干——如果第 2 条消息失败后进了重试队列,第 3 条消息先被消费,顺序就乱了。 ...

十一月 9, 2022 · 6 分钟 · 1151 字 · yaomingye
Cat Radio