十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

高并发削峰填谷实战:MQ选型、分层架构与全链路可靠性设计

高并发削峰填谷实战:MQ选型、分层架构与全链路可靠性设计 1. 这不是“加个MQ就完事”的故事高并发下削峰填谷的真实战场很多人一听到“高并发消息队列”脑子里立刻浮现出一张经典架构图前端流量像洪水一样冲过来后面摆个Kafka或RocketMQ当“水坝”流量被拦住、缓存、慢慢放行——然后就以为“削峰填谷”完成了。我干这行十年亲手参与过6次电商大促、3次政务系统秒杀、2次金融级实时风控系统的压测调优踩过的坑比读过的文档还多。实话讲削峰填谷从来不是把消息丢进队列就万事大吉而是一场在吞吐、延迟、一致性、容错四条钢丝上同时走的平衡术。你用的不是MQ是压力计、缓冲器、时间调节阀更是整个系统稳定性的最后一道保险栓。核心关键词——高并发、消息队列、削峰填谷、架构设计、MQ——它们不是孤立术语而是一组强耦合的工程约束条件。所谓“高并发”不是指QPS破万就叫高并不是数字本身而是瞬时请求到达速率远超后端服务处理能力的持续时间窗口。比如一个订单创建接口平时TPS 200大促峰值瞬间冲到8000但数据库写入能力只有1200 TPS中间6800的差额就是必须被“削”掉的峰而“填谷”是指在业务低谷期比如凌晨2点让积压的消息以可控节奏被消费把系统资源利用率拉起来避免资源闲置浪费。这个过程里MQ不是搬运工它是时空转换器把“时间上不可控的请求洪流”转换成“空间上可调度的任务队列”。适合谁看如果你正面临以下任一场景这篇就是为你写的你刚接到需求要支撑一场预计50万用户同时抢券的活动但DB和库存服务明确告诉你“扛不住瞬时写压力”你发现日志服务总在凌晨崩查监控发现是白天积压的埋点消息在零点集中爆发你被问到“为什么用户点了两次提交按钮订单却生成了两单”而你的回答还停留在“前端防重”层面你在准备架构师面试被追问“如果Kafka集群挂了你的订单流程怎么兜底”。这不是理论推演是我在真实生产环境里用血泪换来的操作手册。下面拆解的每一个环节都对应着一次线上事故、一次压测失败、一次深夜救火。我们不谈“应该怎么做”只说“实际怎么做、为什么这么选、踩过什么坑”。2. 削峰填谷不是功能是架构决策从目标反推技术选型与分层设计2.1 先想清楚你要削的是什么峰填的是什么谷很多团队一上来就争论“用Kafka还是RabbitMQ”结果发现连“峰”的物理形态都没搞清。削峰填谷的本质是解决“请求到达节奏”与“服务处理节奏”之间的异步失配问题。但这个失配有三种完全不同的表现形式决定了底层技术选型的底层逻辑瞬时脉冲型峰值如电商秒杀开始瞬间10万请求在200ms内打进来。特点是时间极短、强度极高、允许少量丢失或延迟。典型场景抢购、抽奖、红包雨。应对策略前端限流 MQ快速接收 后端异步落库允许部分请求排队甚至降级如返回“排队中”。持续高负载型峰值如双11全天订单量是平日的8倍持续8小时。特点是时间长、总量大、不允许丢失、要求最终一致。典型场景订单创建、支付回调、物流更新。应对策略MQ持久化 消费端水平扩容 幂等重试死信处理必须保证每条消息100%处理成功。混合型峰值如春晚红包既有开场瞬间的脉冲抢红包又有后续中奖通知的持续压力发奖。特点是峰谷叠加、链路长、依赖多。应对策略分层MQ脉冲用轻量级队列快速接入后续流转用高可靠队列做状态机驱动并建立跨队列的消息追踪ID体系。提示别被“高并发”三个字带偏。先用监控工具如PrometheusGrafana抓取你系统真实的请求到达曲线Request Arrival Rate标出P999响应时间拐点、DB慢SQL出现时段、线程池满载时刻。这些数据点才是你设计削峰策略的唯一依据。没有数据所有架构设计都是空中楼阁。2.2 MQ选型不是比参数而是比“失控时的底线”市面上主流MQKafka、RocketMQ、RabbitMQ、Pulsar的对比文章汗牛充栋但真正决定选型的往往不是吞吐量数字而是当系统濒临崩溃时它给你留下的逃生通道有多宽。我按实战经验总结出四个关键维度维度KafkaRocketMQRabbitMQPulsar消息丢失容忍度高默认异步刷盘宕机可能丢数中同步刷盘可保不丢但性能下降30%低支持事务镜像队列丢数概率0.001%中分层存储Broker宕机不丢消息堆积能力极强TB级堆积无压力靠磁盘顺序写强GB级堆积需调优依赖PageCache弱内存磁盘混合堆积超100万易OOM极强云原生设计自动分片消息顺序性保障分区级有序同一Partition内严格FIFO全局有序Topic级但吞吐下降50%队列级有序单Queue内有序多Consumer需协调Topic级有序通过Managed Ledger实现运维复杂度高ZooKeeper依赖、磁盘IO敏感、扩容需重平衡中NameServerBroker架构配置项少低单机部署即可用Web管理界面友好高BookKeeperBroker双组件云环境更适配选型决策树如果你的业务绝对不能丢消息如金融交易流水且QPS5000选RabbitMQ它的ACK机制和镜像队列是经过银行级验证的如果你的峰值持续数小时且消息量达亿级如社交Feed推送且团队有Java生态经验选RocketMQ它的事务消息和定时消息对订单场景太友好如果你的系统本身就是大数据平台已有Kafka集群且能接受“最多一次”语义如用户行为日志直接复用别再造轮子如果你正在构建云原生架构且需要跨地域复制如全球用户注册事件同步Pulsar的分层存储和多租户是唯一解。注意别迷信“国产替代”。我见过团队为响应号召强行上RocketMQ结果因不熟悉其消费位点重置机制在一次网络分区后导致百万条消息重复消费回滚数据花了三天。选型第一原则团队熟悉度 参数指标 厂商背景。一个你团队能半夜三点快速定位问题的MQ比一个参数漂亮但文档晦涩的MQ价值高十倍。2.3 架构分层为什么必须把“削峰”和“填谷”拆成两层新手常犯的错误是把所有消息塞进同一个Topic/Exchange认为“MQ自己会调度”。现实是单一队列会把不同SLA要求的消息绑死在同一根绳上。比如订单创建消息要求100%不丢、5秒内处理和用户浏览日志允许丢失、1小时后处理混在一起一旦日志消费卡住订单消息也会被堵死。正确的分层是“三明治”结构接入层Ingress Layer轻量级、高吞吐、低延迟。作用是“快进快出”只做协议转换和基础校验。推荐用RabbitMQ的Direct Exchange或Kafka的单独Topic配置极简关闭消息确认、禁用持久化目标是单机10万 QPS。这里允许丢弃明显非法请求如空参数、恶意刷单IP但绝不处理业务逻辑。核心层Core Layer高可靠、强一致、可追溯。作用是“稳扎稳打”承载核心业务消息。用RocketMQ的事务消息或Kafka的Exactly-Once语义开启同步刷盘、副本数≥3、消息体带全局TraceID。这是你架构的“心脏”所有幂等、重试、死信逻辑都在这一层落地。归档层Archive Layer低成本、大容量、只读。作用是“冷备兜底”存放已处理完成的消息原始快照。用对象存储如S3、OSS或HDFS按天分区保留90天。当业务方质疑“某笔订单没收到通知”你能在5分钟内从归档层捞出原始消息体而不是翻Kafka日志。这种分层不是炫技而是把故障域隔离。去年我们一个支付系统归档层的OSS桶权限配置错误导致无法写入但核心层和接入层完全不受影响业务零感知。这就是分层的价值。3. 核心细节从消息生产到消费的全链路实操要点3.1 生产端不是发消息是“申请排队资格”很多人以为生产者就是producer.send()一行代码其实这是最危险的环节。生产端的健壮性决定了整个削峰系统的生死线。我总结出四个必须落地的硬性规范第一永远启用异步发送 回调处理同步发送会阻塞业务线程一旦MQ响应慢整个HTTP请求就卡死。正确姿势// RocketMQ示例 SendCallback sendCallback new SendCallback() { Override public void onSuccess(SendResult sendResult) { // 记录发送成功日志含msgId和queueOffset log.info(Send success, msgId: {}, offset: {}, sendResult.getMsgId(), sendResult.getQueueOffset()); } Override public void onException(Throwable e) { // 关键这里不能只打日志必须触发降级 if (e instanceof RemotingTimeoutException) { // 网络超时尝试本地缓存定时重发 localCache.put(msg.getKeys(), msg); scheduleRetry(msg); } else if (e instanceof MQClientException) { // 客户端异常大概率是配置错误立即告警 alertService.send(MQ Client Error, e); } } }; producer.sendAsync(message, sendCallback);实操心得回调里的onException是救命稻草。我见过太多团队在这里只写e.printStackTrace()结果MQ集群抖动时生产者线程池被占满整个服务雪崩。必须区分异常类型对网络超时做本地缓存重试对配置错误立即告警。第二消息体必须携带完整上下文而非ID引用常见错误生产者只发{orderId:12345}让消费者自己去DB查订单详情。这会导致消费端成为DB热点且无法做到消息幂等因为DB数据可能被修改。正确做法是发送完整快照{ eventId: evt_20240520_abc123, eventType: ORDER_CREATED, payload: { orderId: 12345, userId: u7890, items: [{skuId:s1,qty:2},{skuId:s2,qty:1}], totalAmount: 299.00, createdAt: 2024-05-20T10:30:45.123Z }, traceId: trc_abcdef1234567890, version: 1.0 }好处消费端无需查库直接处理消息体自带版本号便于未来兼容升级traceId打通全链路监控。第三强制设置消息Key绑定业务实体Kafka/RocketMQ的分区策略默认是Hash(key)%partitionCount。如果你不设Key同一批订单消息会被打散到不同Partition破坏顺序性。正确姿势// RocketMQ用订单ID作为Key确保同一订单的所有消息进同一Queue message.setKeys(ORDER_ order.getId()); // Kafka用订单ID作为Key确保同一订单进同一Partition ProducerRecordString, byte[] record new ProducerRecord(order_topic, order.getId(), messageBytes);注意Key不能是随机UUID必须是业务主键。否则“削峰”就失去了意义——你削的不是流量是业务一致性。第四生产端必须内置熔断器当MQ集群不可用时生产者不能无脑重试。我们采用HystrixSentinel双保险Hystrix控制单个生产者实例的失败率阈值50%10秒窗口Sentinel控制全局QPS如订单Topic每秒最多发5000条超限直接拒绝并返回友好提示熔断触发后自动切换到本地文件队列File Queue每5秒刷盘一次待MQ恢复后自动重发。这套组合拳让我们在去年一次Kafka集群网络分区中业务损失为0用户只看到“系统繁忙请稍后再试”而非“下单失败”。3.2 存储层消息不是存进去就安全了是存得“可追溯、可审计、可回溯”MQ的存储配置是削峰填谷的物理基石。参数调优不是玄学而是基于硬件和业务的精确计算。磁盘选择SSD是底线NVMe是标配机械硬盘HDD在随机写场景下IOPS不足200而Kafka/RocketMQ的刷盘是顺序写但PageCache失效时会触发大量随机IO。实测数据1TB SATA SSD顺序写吞吐 500MB/s随机写 IOPS 50k1TB NVMe SSD顺序写吞吐 3GB/s随机写 IOPS 500k企业级NVMe如Intel Optane延迟稳定在10μs内无长尾。结论只要预算允许一律上NVMe。别省这点钱它决定了你能否扛住峰值后的消息堆积风暴。刷盘策略同步刷盘不是性能杀手是信任基石很多人怕同步刷盘慢改用异步。但异步刷盘在服务器断电时可能丢失数秒消息。我们的折中方案对于核心Topic如订单、支付启用flushDiskTypeSYNC_FLUSH并用UPS保障电源对于非核心Topic如日志、埋点用ASYNC_FLUSH但配置flushIntervalConsume1000每秒强制刷盘一次RocketMQ的commitLog目录必须独立挂载禁止与OS共用磁盘避免IO争抢。副本与ISR宁可慢不可丢Kafka的min.insync.replicas2是底线意味着至少2个副本同步成功才认为写入成功。我们生产环境强制设为3并监控UnderReplicatedPartitions指标。当该值0时自动触发告警并暂停新消息写入直到副本恢复。数据一致性永远优先于吞吐量。消息过期不是越长越好是“够用就好”Kafka默认retention.ms6048000007天但我们的订单Topic只设168h7天日志Topic设72h。原因过期时间越长磁盘占用越大且增加log cleanup的CPU开销。我们通过监控LogStartOffset和LogEndOffset的差值动态调整过期时间——当堆积量超过阈值自动缩短过期时间并告警。实操心得定期执行kafka-log-dirs.sh检查各Topic磁盘占用对长期未消费的Topic如测试环境遗留立即清理。我们曾因一个废弃的Topic占满磁盘导致整个集群不可用教训深刻。3.3 消费端填谷不是“慢慢吃”是“智能调度吃”消费端是削峰填谷的终点也是最容易出问题的环节。填谷失败等于削峰白干。关键在三个动作拉取、处理、确认。拉取策略Pull模式是可控性的根源Kafka/RocketMQ都用Pull模式消费者主动拉取消息而非Push模式MQ推送。这是因为Pull能精确控制消费节奏设置max.poll.records100每次最多拉100条避免单次处理过久设置fetch.max.wait.ms500最长等待500ms确保低峰期也能及时拉取动态调整fetch.min.bytes1024最小拉取1KB避免小消息频繁拉取。处理模型单线程消费 多线程处理 最佳实践消费者线程负责拉取消息并放入内存队列业务线程池从队列取任务处理。这样既保证了消息顺序单线程拉取又提升了吞吐多线程处理// Kafka Consumer线程单线程 while (true) { ConsumerRecordsString, byte[] records consumer.poll(Duration.ofMillis(500)); for (ConsumerRecordString, byte[] record : records) { // 放入内存阻塞队列 processingQueue.put(new ProcessingTask(record)); } } // 业务处理线程池10个线程 ExecutorService processor Executors.newFixedThreadPool(10); processor.submit(() - { while (true) { ProcessingTask task processingQueue.take(); try { processOrder(task.record); // 业务逻辑 task.record.commitSync(); // 手动提交offset } catch (Exception e) { // 处理失败进入死信队列 sendToDLQ(task.record, e); } } });确认机制手动提交是唯一可靠方式自动提交enable.auto.committrue在消费失败时会丢失消息。必须手动提交commitSync()同步阻塞确保提交成功再继续commitAsync()异步非阻塞但需配合回调处理失败我们采用commitSync() 重试三次失败则记录到死信Topic。死信队列DLQ不是兜底是问题探针DLQ不是用来“放着不管”的而是问题分析的第一现场。我们要求DLQ Topic必须开启单独监控消息进入速率、堆积量、平均处理时长每条DLQ消息必须包含原始消息体、失败堆栈、重试次数、首次失败时间自动化脚本每15分钟扫描DLQ对重试3次的消息触发告警并生成诊断报告如“近1小时DLQ中80%为库存不足异常建议扩容库存服务”。4. 实操过程从0到1搭建一个可抗10万QPS的削峰系统4.1 环境准备用最小成本验证架构可行性别一上来就搞K8s集群。我们用三台8C16G的云服务器非SSD盘搭建了一个可抗10万QPS的验证环境成本不到300元/月。步骤如下Step 1部署RocketMQ NameServer Broker单节点# 下载RocketMQ 5.1.4 wget https://archive.apache.org/dist/rocketmq/5.1.4/rocketmq-all-5.1.4-bin-release.zip unzip rocketmq-all-5.1.4-bin-release.zip cd rocketmq-all-5.1.4-bin-release # 修改broker.conf关键配置 brokerClusterName DefaultCluster brokerName broker-a brokerId 0 deleteWhen 04 fileReservedTime 48 brokerRole ASYNC_MASTER flushDiskType SYNC_FLUSH # 开启事务消息 transientStorePoolEnable true # 堆内存调至4G JAVA_OPT${JAVA_OPT} -server -Xms4g -Xmx4g -XX:MetaspaceSize128m -XX:MaxMetaspaceSize320m启动命令nohup sh bin/mqnamesrv # NameServer nohup sh bin/mqbroker -n localhost:9876 -c conf/broker.conf # BrokerStep 2创建Topic与权限# 创建订单Topic8个Queue对应8核CPU sh bin/mqadmin updateTopic -n localhost:9876 -t order_topic -c DefaultCluster -r 8 -w 8 # 创建死信Topic sh bin/mqadmin updateTopic -n localhost:9876 -t order_dlq -c DefaultCluster -r 4 -w 4 # 创建生产者组和消费者组 sh bin/mqadmin updateGroup -n localhost:9876 -g producer_order -c DefaultCluster -m true sh bin/mqadmin updateGroup -n localhost:9876 -g consumer_order -c DefaultCluster -m trueStep 3压测脚本编写JMeter Java Sampler不用JMeter GUI写Java Sampler直连RocketMQpublic class OrderProducerSampler implements JavaSamplerClient { private DefaultMQProducer producer; Override public void setupTest(JavaSamplerContext context) { producer new DefaultMQProducer(producer_order); producer.setNamesrvAddr(localhost:9876); producer.start(); } Override public void runTest(JavaSamplerContext context) { // 构造100条订单消息 for (int i 0; i 100; i) { Message msg new Message(order_topic, ORDER_ UUID.randomUUID().toString(), JSON.toJSONString(buildOrder()).getBytes()); producer.send(msg); // 异步发送 } } }JMeter线程组设为1000线程循环100次总请求数10万。监控Broker的SendMessageTps指标目标值≥10万。Step 4消费端部署与调优消费端用Spring Boot RocketMQ Starterrocketmq: name-server: localhost:9876 producer: group: producer_order consumer: group: consumer_order # 关键拉取批次大小 pull-batch-size: 32 # 关键消费线程数CPU核数*2 consume-thread-max: 16启动后用sh bin/mqadmin statsAll -n localhost:9876查看消费TPS目标值≥8万预留20%缓冲。实测数据三台机器1 NameServer 2 Broker在SSD盘下轻松扛住12万QPS写入消费TPS达10.5万。当换成HDD盘时写入TPS暴跌至3万证明磁盘是瓶颈。4.2 核心参数调优每一处配置都有物理意义Broker端关键参数参数推荐值物理意义调优依据brokerRoleASYNC_MASTER主从异步复制平衡性能与可靠性同步复制会增加20ms延迟异步复制在单机房足够flushDiskTypeSYNC_FLUSH每次写入都刷盘订单场景不容许丢消息transientStorePoolEnabletrue开启堆外内存池提升写入性能RocketMQ 5.x必需否则高并发下GC频繁storePathRootDir/data/rocketmq/store独立挂载点避免与系统盘争IOcommitLogFileSize1073741824 (1G)CommitLog单文件大小太小导致文件碎片太大影响恢复速度Consumer端关键参数参数推荐值物理意义调优依据pullBatchSize32每次拉取最大消息数太大会导致单次处理超时太小增加网络开销consumeThreadMin16消费线程池最小线程数 CPU核数 * 2保证低峰期也有足够线程consumeThreadMax32消费线程池最大线程数峰值时自动扩容但需监控CPU使用率pullInterval5000拉取间隔ms太短浪费CPU太长降低实时性JVM调优Broker# JVM参数放在runbroker.sh中 JAVA_OPT${JAVA_OPT} -server -Xms4g -Xmx4g -XX:MetaspaceSize256m -XX:MaxMetaspaceSize512m JAVA_OPT${JAVA_OPT} -XX:UseG1GC -XX:G1HeapRegionSize2M -XX:MaxGCPauseMillis50 JAVA_OPT${JAVA_OPT} -XX:PrintGCDetails -XX:PrintGCDateStamps -Xloggc:/data/rocketmq/logs/gc.logG1 GC是必须的因为RocketMQ堆内存主要用于Netty缓冲区和消息索引G1能精准控制停顿时间。4.3 全链路监控没有监控的架构等于裸奔我们用Prometheus Grafana ELK搭建监控体系核心指标必须覆盖生产端监控rocketmq_producer_send_tps生产TPS预警阈值目标值*0.8rocketmq_producer_send_fail_rate发送失败率0.1%立即告警rocketmq_producer_local_cache_size本地缓存消息数1000说明MQ不可用。Broker端监控rocketmq_broker_commitlog_flush_time_max刷盘最大耗时500ms说明磁盘IO瓶颈rocketmq_broker_ha_master_diff主从同步延迟1000ms触发主从切换rocketmq_broker_topic_msg_total_todayTopic今日消息总量环比昨日下降50%说明上游异常。消费端监控rocketmq_consumer_offset_lag消费延迟消息堆积量10000告警rocketmq_consumer_process_time_max单条消息处理最长时间5000ms说明业务逻辑有瓶颈rocketmq_consumer_dlq_rate死信率0.01%触发根因分析。告警规则示例Prometheus Alert Rule- alert: RocketMQ_Consumer_Lag_High expr: rocketmq_consumer_offset_lag{topicorder_topic} 5000 for: 2m labels: severity: critical annotations: summary: Order topic consumer lag is high description: Lag is {{ $value }} messages, check consumer health - alert: RocketMQ_Broker_Flush_Slow expr: rocketmq_broker_commitlog_flush_time_max 1000 for: 1m labels: severity: warning annotations: summary: Broker flush time is too slow description: Flush max time is {{ $value }}ms, check disk IO实操心得监控不是装完就完事。我们每周五下午做“监控有效性演练”随机屏蔽一个指标看告警是否触发模拟一条死信消息看诊断报告是否自动生成。只有经过验证的监控才是真正的安全网。5. 常见问题与排查技巧实录那些深夜救火时的真实案例5.1 “消息重复消费”不是Bug是设计缺陷的必然结果几乎所有团队都会遇到这个问题。根本原因不是MQ有问题而是消费端没有实现幂等。我整理了最典型的三类场景及解法场景1消费成功但ACK失败消费者处理完消息调用consumer.commitSync()时网络超时MQ认为消息未处理重新投递。✅ 解法业务主键状态机// 订单创建消息以orderId为主键 String orderId json.get(orderId); // 先查DB如果状态已是CREATED直接return if (orderService.getStatus(orderId).equals(CREATED)) { return; // 幂等退出 } // 否则创建订单 orderService.createOrder(json);场景2消费端重启导致重复拉取消费者进程重启offset未及时提交重启后从旧offset开始拉取。✅ 解法本地缓存布隆过滤器// 内存中维护最近10分钟处理过的orderId布隆过滤器 BloomFilterString bloomFilter BloomFilter.create(Funnels.stringFunnel(Charset.defaultCharset()), 1000000, 0.01); // 处理前检查 if (bloomFilter.mightContain(orderId)) { return; // 可能重复跳过 } bloomFilter.put(orderId); // 处理业务...场景3MQ集群重平衡Kafka消费者组扩容/缩容时Partition重分配同一消息可能被新旧Consumer同时处理。✅ 解法分布式锁 消息指纹// 消息指纹 md5(topic partition offset body) String fingerprint md5(topic partition offset message.getBody()); // 尝试获取分布式锁Redis if (redisLock.tryLock(fingerprint, 30, TimeUnit.SECONDS)) { try { processMessage(message); } finally { redisLock.unlock(fingerprint); } } else { // 锁获取失败说明正在被其他Consumer处理直接return return; }注意别用“数据库唯一索引”做幂等那是最后防线。高频场景下唯一索引冲突会产生大量数据库错误日志拖慢整个系统。5.2 “前端点两次算两条消息吗”——从源头杜绝无效流量这个问题暴露了架构的盲区。答案是如果后端没做防重点两次就是两条消息。但解决方案不在后端而在前端网关协同前端层按钮点击后立即置灰显示“提交中...”3秒内禁止重复点击网关层Nginx/OpenResty做请求指纹去重# 基于URLBody MD5生成指纹 set $fingerprint ; if ($request_method POST) { set $fingerprint ${uri}_${body_md5}; } # 5分钟内相同指纹只放行一次 lua_shared_dict dedupe_dict 10m; location /api/order { access_by_lua_block { local dict ngx.shared.dedupe_dict local ok, err dict:set($fingerprint, 1, 300) -- 5分钟过期 if not ok then ngx.exit(429) -- 返回429 Too Many Requests end } }后端层Token机制防刷用户提交前后端发放一次性Token表单提交时携带后端校验并消耗Token。这套组合拳让我们订单重复提交率从0.3%降至0.002%。5.3 “MQ挂了怎么办”——真正的高可用不是不挂是挂了也不影响业务MQ集群不可能100%可用。我们的SLA承诺是“MQ不可用时核心业务降级为同步模式用户无感知”。实现路径第一层本地消息表Local Message Table业务服务在同一个本地事务中写业务DB 写本地消息表status‘pending’。INSERT INTO local_message (msg_id, topic, payload, status, created_at) VALUES (msg_123, order_topic, {orderId:12345}, pending, NOW());第二层定时任务扫描独立服务每5秒扫描local_message表将statuspending的消息发往MQ成功后update statussent。第三层兜底补偿当MQ恢复后扫描statuspending且created_at5分钟的消息触发人工干预流程如短信通知用户“订单已受理请耐心等待”。这套方案让我们在去年一次RocketMQ集群升级中断期间订单创建成功率保持100%只是用户看到“预计30分钟内处理完成”的提示。5.4 面试题高频陷阱别背答案要懂原理面试官问“MQ如何保证消息不丢失”标准答案是“生产者ACKBroker持久化消费者手动提交”。但真实情况复杂得多生产者ACK级别acks1Leader写入即返回 vsacksallISR全部写入才返回。后者更安全但延迟高20%Broker持久化时机Kafka的log.flush.interval.messages10000表示1万条刷一次盘但若服务器断电这1万条就丢了消费者提交时机auto.commitfalse后commitSync()在处理完
返回列表