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

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

SpringBoot RocketMQ 全操作指南

SpringBoot 集成 RocketMQ:从发送到消费 📖 前置阅读:本文假设读者已理解 RocketMQ 的核心概念(NameServer、Broker、Topic、Queue、ConsumerGroup)。如果还不熟悉,建议先阅读 RocketMQ 核心架构与消息模型。 Part 1:概念与前置 1.1 本文目标 上一篇用原版 RocketMQ Java Client 写了 DefaultMQProducer + DefaultMQPushConsumer。真实的 SpringBoot 项目里不需要那么多样板代码——rocketmq-spring-boot-starter 帮你处理了 NameServer 连接、Producer 启动、Consumer 注册。 读完这篇会掌握: MqHelper 封装——为什么要在 RocketMQTemplate 上再包一层,asyncSend + SendCallback 的真实用法 30 分钟延迟取消订单——RocketMQ 内置 delayLevel 的完整实战流程,和 RabbitMQ 方案的对比 @RocketMQMessageListener 消费者——MessageExt 手动反序列化 vs 泛型自动解析,以及三个真实业务消费者 Domain Entity 直传——为什么不做 DTO 转换,以及什么情况下不能这样做 双 MQ 基础设施先行——RabbitMQ 拓扑已就绪但全用 RocketMQ 的设计决策 1.2 前置条件 前置项 具体要求 验证命令 JDK 17+(8+ 也兼容) java -version SpringBoot 3.x(文中用 3.2) mvn dependency:tree | grep spring-boot RocketMQ 5.1.4(NameServer + Broker 都在运行) docker ps | grep rocketmq 前置知识 NameServer/Broker/Topic/Queue 概念 — 确认 RocketMQ 在跑: ...

十一月 8, 2022 · 16 分钟 · 3233 字 · yaomingye

RabbitMQ 生产环境部署与调优

生产环境部署与调优 📖 前置阅读:本文是 RabbitMQ 系列的终篇,假设读者已经掌握前五篇的全部内容(核心概念、交换机类型、SpringBoot 集成、消息可靠性、高级特性)。 一、⚡ 问题切入:单机 RabbitMQ 什么时候扛不住? 前三篇代码都在本地单机 RabbitMQ 上跑的——一个 Docker 容器,内存 1G,磁盘 10G。跑到生产环境会发生什么? 场景 单机 RabbitMQ 的后果 服务器宕机 整个 MQ 服务中断——所有生产者阻塞、消费者闲置 消息积压 50 万条 内存打满 → RabbitMQ 触发内存告警 → 阻塞所有生产者(flow control) 磁盘写满 RabbitMQ 拒绝所有写入→ 消息丢失 每秒 2 万条消息 单机 CPU 100%,消息延迟从 1ms 飙升到 500ms 单机不是不能用,但要知道它的边界。以下场景必须上集群: 消息不能丢(金融交易、订单处理) 服务不能停(7×24 在线业务) 吞吐量超过单机极限(> 5 万 msg/s) 二、RabbitMQ 集群架构 2.1 集群的基本原理 RabbitMQ 集群是多个 Erlang 节点组成的对等网络——每个节点运行一个 RabbitMQ 实例。 集群中共享的东西: Exchange、Queue、Binding 的元数据(定义信息)——在所有节点上自动同步 用户、vhost、权限——自动同步 集群中不共享的东西: ...

十一月 6, 2022 · 6 分钟 · 1085 字 · yaomingye

RabbitMQ 延迟队列与高级特性

延迟队列与高级特性 📖 前置阅读:本文假设读者已掌握 RabbitMQ 的基础操作(Exchange/Queue/Binding/手动ACK)和 SpringBoot 集成。如果还不熟悉,建议先阅读前三篇。 一、⚡ 问题切入:有些事不能现在做 看几个每天都在发生的业务需求: 用户下单后 30 分钟未支付,自动取消订单 用户注册后 7 天未登录,发一封"想你了"邮件 优惠券到期前 3 小时,发短信提醒用户使用 支付成功后立即通知商家,但每封邮件间隔 5 秒(防被邮箱限流) 这些需求的共同点:消息不能发出去就被立刻消费,需要在未来某个时间点才能被消费。 普通队列是"发了就收"——消息一进队就被消费者拿走。延迟队列是"发了先等着,时间到了再收"。 RabbitMQ 没有原生的延迟队列类型,但有两种方式可以实现。 二、延迟队列方案一:TTL + 死信队列 2.1 原理 这是利用已有的机制组合而成的"曲线救国"方案: Producer → 普通 Exchange → 死信队列(当"延迟缓冲区") → TTL 到期 → 死信 Exchange → 实际消费队列 → Consumer ↑ 消息过期变死信自动转入 核心思想:创建一个没有消费者的队列,设置 TTL。消息在这个队列中"等"TTL 时间后过期变成死信,被自动转发到真正的消费队列——消费者只监听消费队列。消费者感知不到延迟的存在。 flowchart TD 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; subgraph DELAY_FLOW ["TTL + DLQ 延迟队列流程"] P([Producer]) -->|"发送\n带 TTL"| DE[delay.exchange] DE -->|"routing: order.delay.30m"| DQ[queue.order.delay.30m\nTTL=30分钟\ndeadLetterExchange=real.exchange\ndeadLetterRoutingKey=order.real] DQ -->|"30分钟后\n消息过期变死信"| DLX[real.exchange] DLX -->|"routingKey=order.real"| RQ[queue.order.real\n实际消费队列] RQ --> C([Consumer\n订单取消服务]) C -.->|"⚠️ 消费者不监听延迟队列\n只监听实际消费队列"| DQ end class P startEnd; class DE,DQ,DLX,RQ process; class C startEnd; 2.2 配置代码 每段延迟时间需要一个独立的队列——每个队列有固定的 TTL: ...

十一月 5, 2022 · 6 分钟 · 1205 字 · yaomingye

RabbitMQ 消息可靠性保障

消息可靠性保障:三道防线全覆盖 📖 前置阅读:本文假设读者已掌握 SpringBoot RabbitMQ 的基本操作(RabbitTemplate 发送、@RabbitListener 消费、手动 ACK)。如果还不熟悉,建议先阅读 SpringBoot RabbitMQ 全操作指南。 一、⚡ 问题切入:消息去哪儿了? 先看一段日常的订单处理代码: // 下单成功后发消息 @Service public class OrderService { @Transactional public void createOrder(OrderRequest req) { orderMapper.insert(req.toOrder()); // 1. 写 MySQL rabbitTemplate.convertAndSend( // 2. 发消息 "order.exchange", "order.created", req); } } 表面看起来没问题。但消息真的被消费了吗?在以下任何一个环节都可能丢: Producer → [网络] → RabbitMQ → [网络] → Consumer ① 发送丢失 ② Broker 宕机丢失 ③ 消费失败丢失 环节 丢失原因 后果 ① 生产者 → Broker 网络断连、Exchange 不存在、消息路由失败 消息根本没进队列 ② Broker 存储 RabbitMQ 进程崩溃、服务器断电 内存中的消息全部丢失 ③ Consumer 消费 消费者处理到一半挂了、代码异常没 ACK 消息被取走但实际没处理完 这三个环节必须逐一设防——RabbitMQ 提供了完整的机制,但需要生产者、Broker、消费者三端配合。 ...

十一月 4, 2022 · 8 分钟 · 1512 字 · yaomingye

SpringBoot RabbitMQ 全操作指南

SpringBoot 集成 RabbitMQ:从发送到消费 📖 前置阅读:本文假设读者已理解 RabbitMQ 的核心概念(Exchange、Queue、Binding、RoutingKey)和四种交换机类型。如果还不熟悉,建议先阅读前两篇: RabbitMQ 核心概念与 AMQP 协议 交换机类型完全指南 Part 1:概念与前置 1.1 本文目标 前两篇用 RabbitMQ 原生 Java Client 写了所有代码——channel.basicPublish、channel.basicConsume、手动 basicAck。理解底层是正确的,但真正进项目时,Spring AMQP 帮我们做了 90% 的重复工作。 读完这篇会掌握: 用 RabbitTemplate 一行代码发消息(替代 channel.basicPublish 那一大堆) 用 @RabbitListener 注解收消息(替代手动 basicConsume + DeliverCallback) 用 Jackson2JsonMessageConverter 自动序列化/反序列化 Java 对象 用 @Bean + 声明式配置 管理 Exchange/Queue/Binding(替代每次启动时 channel.exchangeDeclare) 三种交换机在 Spring 中的完整示例代码(Direct / Fanout / Topic) 手动 ACK 的配置和坑 1.2 前置条件 前置项 具体要求 验证命令 JDK 17+(8+ 也兼容) java -version Maven 3.6+ mvn -v SpringBoot 3.x(文中用 3.2) mvn dependency:tree | grep spring-boot RabbitMQ 3.12+(management 版) docker ps | grep rabbitmq 前置知识 前两篇的 Exchange/Queue/Binding/RoutingKey 概念 — 确认 RabbitMQ 在跑: ...

十一月 3, 2022 · 13 分钟 · 2677 字 · yaomingye

Redis + Caffeine 双层缓存:降级与容错

Redis + Caffeine 双层缓存 📖 前置阅读:本文是缓存架构的进阶级文章,假设读者已经掌握了 Redis 的基础操作和 Caffeine 本地缓存的 API。如果还不熟悉,建议先阅读: SpringBoot Redis 全操作指南 —— Redis 实战篇 Caffeine 核心与 SpringBoot 集成 —— Caffeine 入门篇 一、⚡ 问题切入:凌晨三点,Redis 挂了 凌晨三点,Redis 内存用满——大量的 TTL 同时到期 + 新一波定时任务写入,导致内存 OOM,Redis 进程被系统 kill。你的服务所有缓存请求全部报错,瞬间全部穿透到 MySQL,数据库连接池耗尽,整个系统不可用。 值班群炸了。你翻日志发现——服务启动时所有 @Cacheable 都配置了 Redis,Redis 一挂连个兜底的都没有。 Redis 是高可用的——有哨兵(Sentinel)、有集群(Cluster),官方说可用性能到 99.99%。但 99.99% 意味着一年有将近 1 小时的不可用时间。这 1 小时如果发生在双十一,后果就不是"维护了一次",而是"事故"。 本地缓存的价值不只是"快",更是 Redis 挂了时的最后一道防线。就算 Redis 是全宇宙最高可用的服务,网络也可能抖——交换机故障、机房间专线断掉、Kubernetes 网络策略变更——这些事情的发生概率比 Redis 自身故障高得多。 本篇要解决的问题:构建 Redis(远程)+ Caffeine(本地)双层缓存架构,把 Redis 的不可用当成"迟早会发生的事"来设计,而不是寄望于它不会发生。 📌 真实场景:数据字典——最简单的双层缓存 在进入复杂架构之前,先看一个真实项目里怎么用 Spring Cache + Caffeine + Redis 做双层缓存。不是所有场景都需要 200 行的 TieredCacheManager——有时候 3 行配置 + 1 个注解就够了。 ...

十月 31, 2022 · 11 分钟 · 2341 字 · yaomingye
Cat Radio