文章 587
评论 5
浏览 228699
RocketMQ Tag/SQL92 过滤实战:客户端拉取无效消息浪费带宽?Broker 端精准过滤,网络开销降 70%!

RocketMQ Tag/SQL92 过滤实战:客户端拉取无效消息浪费带宽?Broker 端精准过滤,网络开销降 70%!

公司做物流系统,订单状态变更通过 RocketMQ 广播。刚开始只有一个消费者,所有消息照单全收。后来业务拆分,多了十几个微服务各自订阅感兴趣的消息。问题来了——每个服务都从 Broker 拉全量消息,然后自己过滤。一天几千万条消息,每个服务实际需要的不到 10%,90% 拉过来就扔了。内网带宽跑满,消息堆积,消费延迟越来越大。 这个问题特别容易被忽视。因为消息队列用起来太简单了,Producer 发、Consumer 收,中间不用管。但当你有了十几个 Consumer 各自只需要不同子集的消息时,全量拉取就是在用带宽换便利。 今天聊聊 RocketMQ 的 Tag 和 SQL92 两种 Broker 端过滤方式,把过滤逻辑从消费端前移到 Broker,让不需要的消息根本不出 Broker 的门。 客户端过滤的问题在哪 默认的 Push Consumer 模式,客户端从 Broker 拉消息的流程是这样的: Broker ──拉取──→ Consumer(全量消息) │ ├─ 需要的消息(10%)→ 处理 └─ 不需要的消息(90%)→ 丢弃 看起来没什么问题?你丢你的,又不占硬盘....

RocketMQ 事务消息半消息清理:Half Message 堆积导致 Broker 磁盘告警?自动补偿机制!

RocketMQ 事务消息半消息清理:Half Message 堆积导致 Broker 磁盘告警?自动补偿机制!

做过分布式事务开发的朋友肯定都遇到过这个问题:使用 RocketMQ 的事务消息时,由于网络抖动、服务宕机、消费者超时等原因,部分 Half Message(半消息)无法被正确处理,导致在 Broker 上不断堆积。这不仅占用磁盘空间,严重时还会触发磁盘告警,影响整个消息队列的稳定性。 我之前就遇到过这样一个案例:某天凌晨,监控告警显示 RocketMQ Broker 的磁盘使用率突然飙升至 85%,马上就要触达 90% 的告警阈值。排查后发现,是某个服务的数据库在凌晨进行大批量数据迁移时,部分事务消息的本地事务执行失败,但由于网络重试机制的问题,相关的 Half Message 没有被正确回滚,大批量堆积在了 Broker 上。 今天我们就来聊聊 RocketMQ 事务消息 Half Message 的自动清理方案,让您的系统远离磁盘告警的困扰。 Half Message 的产生与堆积原因 1. 事务消息的执行流程 首先,让我们回顾一下 RocketMQ 事务消息的完整流程: ┌─────────────────────────────────────────────────────....

SpringBoot + RocketMQ 异步批量发送优化:生产端吞吐提升 5 倍,RT 降低 80%!

SpringBoot + RocketMQ 异步批量发送优化:生产端吞吐提升 5 倍,RT 降低 80%!

在高并发场景下,消息发送的性能直接影响系统的整体吞吐量。传统单条消息发送模式存在以下问题: 每次发送都需要网络往返,RT(响应时间)高 broker 压力增大, TPS 上不去 服务器资源利用率低 RocketMQ 的批量发送和异步处理机制可以有效解决这些问题。本文将详细介绍如何在 SpringBoot 中实现 RocketMQ 异步批量发送优化,让生产端吞吐提升 5 倍,RT 降低 80%。 为什么需要批量发送? 先看一下单条发送和批量发送的对比: 单条发送模式: ┌─────────────────────────────────────────────────────────────┐ │ Msg1 ──→ Broker ──→ ACK │ │ ↓ │ │ Msg2 ──→ Broker ──→ ACK │ │ ↓ │ │ Msg3 ──→ Broker ──→ ACK │ │ │ │ 3次网络往返,3次broker处理,RT = N × RTT │ └─────────────────────────────────────────────────────────────┘....

SpringBoot + RocketMQ 异步批量发送优化:生产端吞吐提升 5 倍,RT 降低 80%!

SpringBoot + RocketMQ 异步批量发送优化:生产端吞吐提升 5 倍,RT 降低 80%!

在高并发的业务场景中,消息队列的性能直接影响整个系统的吞吐量和响应时间。特别是在订单处理、日志收集、数据同步等场景下,如何高效地发送消息成为系统性能的关键因素。 RocketMQ 作为一款高性能的消息中间件,在默认配置下已经表现出色,但在极端情况下,单条消息的同步发送仍然会成为性能瓶颈。今天我就跟大家分享一套基于 SpringBoot 的 RocketMQ 异步批量发送优化方案,通过批量发送和异步处理,实现生产端吞吐提升 5 倍,响应时间降低 80% 的显著效果。 为什么需要 RocketMQ 发送优化? 先来说说我们面临的挑战。在高并发场景下,使用默认的 RocketMQ 发送方式会遇到以下问题: 同步发送延迟高:每次发送都需要等待 broker 响应,在网络延迟较大时会严重影响系统性能 频繁网络请求:单条消息发送会产生大量网络请求,增加网络开销 资源消耗大:每个消息都需要独立的线程处理,线程资源消耗大 吞吐量受限:单线程同步发送的吞吐量有限,难以满足高并发需求 以一个电商系统为例,在秒杀活动中,订单创建的峰值可能达到每秒数万个,此时消息发送的性能直接决定了系统能否扛住流量冲击。....

RocketMQ 实战指南:从入门到原理到生产实战、八股面试

RocketMQ 实战指南:从入门到原理到生产实战、八股面试

引言:为什么你需要掌握 RocketMQ? 还记得去年双十一,我们公司核心交易系统因为消息队列性能瓶颈导致订单处理延迟,差点酿成重大事故。事后复盘发现,问题的根源在于团队对消息队列的理解停留在"会用"层面,缺乏深入原理和调优经验。 消息队列作为分布式系统的核心组件,承载着异步解耦、流量削峰、数据分发等关键职责。RocketMQ 作为阿里巴巴开源的分布式消息中间件,凭借其高吞吐量、高可用性、丰富的消息特性,已成为国内互联网公司的首选方案。 本文将从入门到原理,从实战到面试,带你全面掌握 RocketMQ。 一、RocketMQ 入门:10分钟快速上手 1.1 什么是 RocketMQ? RocketMQ 是阿里巴巴于2012年开源的第三代分布式消息中间件,2016年成为 Apache 顶级项目。它借鉴了 Kafka 的高吞吐设计,同时解决了 Kafka 在事务消息、延迟消息、消息轨迹等方面的不足。 核心特点: 高吞吐量:单机写入性能可达10万+ TPS 高可用性:支持多 Master 多 Slave 架构,自动故障切换 丰富的消息类型:普通消息、顺序消息、事务消息、延迟消息 消息轨迹....

SpringBoot + RocketMQ + 事务状态机:订单超时未支付自动取消,消息 100% 可靠触发

SpringBoot + RocketMQ + 事务状态机:订单超时未支付自动取消,消息 100% 可靠触发

为什么订单超时取消这么重要? 在电商系统中,用户下单后通常有30分钟的支付时间。如果用户未在规定时间内支付,系统需要自动取消订单并释放占用的商品库存。这看似简单的功能,实际上涉及多个技术难点: 时间精确控制:必须在指定时间准确触发取消操作 消息可靠性:确保取消指令能被可靠传递和执行 状态一致性:保证订单在整个生命周期中的状态一致性 高并发处理:在大促期间可能有大量订单需要处理 传统的定时轮询方案存在明显缺点:资源消耗大、实时性差、难以处理突发流量。我们需要一个更高效可靠的解决方案。 技术选型:为什么选择RocketMQ + 事务状态机? RocketMQ:可靠的延时消息 RocketMQ提供了强大的延时消息功能,支持预设的延时等级(从秒级到小时级),非常适合处理订单超时场景。其高可用性、高吞吐量的特性,确保了消息的可靠传递。 事务状态机:状态转换的守护者 通过明确定义的状态和转换规则,事务状态机确保订单在任何情况下都保持一致状态,防止非法状态转换。 核心实现:三步走策略 第一步:订单创建时发送延时消息 当用户下单成功后,我们立即发送一条延时消息,指定在30分钟后执行订单检查: //....

服务端开发博客:后端架构、高并发、性能优化与微服务实战教程