ARTICLE DETAIL

资讯详情

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

RocketMQ批量消息发送——设备海量心跳状态上报批量削峰优化

RocketMQ批量消息发送——设备海量心跳状态上报批量削峰优化 批量消息发送——设备海量心跳状态上报批量削峰优化作者黒漂技术佬系列RocketMQ核心原理与无人售货柜项目实战第5篇一、为什么要批量发送先算一笔账。无人售货柜部署1000台设备每台设备每5秒上报一次心跳状态温度、库存、网络信号、故障码等。一秒200条消息一天1728万条。如果每条心跳单独发一条RocketMQ消息网络开销每条消息TCP往返一次1000台设备同时上报网关侧瞬间1000次网络请求Broker压力每条消息都要走完整的存储链路——写CommitLog、刷盘、构建ConsumeQueue索引大量小消息会把磁盘IO打满消费端压力消费者每秒拉取200条消息触发200次数据库写入数据库连接池直接告警这就是典型的扇入问题——大量源头同时往消息队列灌数据瞬时流量尖峰把整条链路打穿。批量发送就是解药把多条消息打包成一个MessageBatch一次网络请求发一批Broker一次存储一批消费者一次拉取一批。吞吐量提升数倍IO压力断崖式下降。二、RocketMQ批量消息原理2.1 MessageBatch结构RocketMQ的批量消息本质上是把多条Message序列化后拼接到一起封装成一个MessageBatch对象对外表现为一条消息。┌──────────────────────────────────────┐ │ MessageBatch │ │ ┌──────────┐ ┌──────────┐ ┌──────┐ │ │ │ Message1 │ │ Message2 │ │ ... │ │ │ └──────────┘ └──────────┘ └──────┘ │ └──────────────────────────────────────┘Broker收到后当作一条大消息存储消费者拉取到之后客户端会自动拆包还原成多条MessageExt。2.2 4MB大小限制RocketMQ对单条消息包括批量消息有大小限制默认4MB。这个4MB不是指你发送的payload大小而是包括消息体、Topic名、属性头、序列化开销在内的总大小。批量发送时必须控制批次总大小不超过4MB。// 计算批次大小publiclongcalcBatchSize(ListMessagemessages){longtotalSize0;for(Messagemessage:messages){totalSizemessage.getTopic().length()message.getBody().length;// 还要加上属性头等开销MapString,Stringpropertiesmessage.getProperties();for(Map.EntryString,Stringentry:properties.entrySet()){totalSizeentry.getKey().length()entry.getValue().length();}totalSize20;// 每条消息约20字节固定开销}returntotalSize;}实战建议单批次控制在1MB以内留足余量。按每条心跳消息约200字节算一个批次最多5000条但实际建议控制在500-1000条兼顾延迟和吞吐。2.3 批量发送的前提条件同一个批次的消息必须满足条件说明同一Topic不同Topic的消息不能混在一个批次同一等待策略不能混合同步发送和异步发送不能是延时/顺序消息批量发送不支持设置DelayTimeLevel和ShardingKey最后一条很关键批量消息和顺序消息、延时消息互斥。需要顺序的场景不能批量需要延时的场景也不能批量。三、SpringBoot实战设备心跳批量上报3.1 业务场景1000台售货柜每5秒上报一次心跳。网关侧边缘节点收集5秒窗口内的所有心跳打包批量发送到RocketMQ消费侧批量写入InfluxDB时序数据库。设备(×1000) ──→ 网关(攒批5秒) ──→ RocketMQ(批量消息) ──→ 消费者(批量写InfluxDB)3.2 心跳消息模型DataBuilderpublicclassDeviceHeartbeatimplementsSerializable{privateStringdeviceId;// 设备IDprivateLongtimestamp;// 心跳时间戳privateDoubletemperature;// 柜内温度privateIntegersignalStrength;// 信号强度privateIntegerstockCount;// 当前库存数privateIntegerfaultCode;// 故障码0正常privateBooleanonline;// 在线状态}3.3 网关侧攒批发送网关侧的核心逻辑是收集窗口 批量发送。用一个缓冲队列收集心跳定时器每5秒触发一次批量发送。Slf4jServicepublicclassHeartbeatBatchProducer{ResourceprivateRocketMQTemplaterocketMQTemplate;/** 心跳缓冲队列 */privatefinalBlockingQueueDeviceHeartbeatbuffernewLinkedBlockingQueue(50_000);/** 单批次最大条数 */privatestaticfinalintBATCH_MAX_SIZE500;/** 单批次最大字节数1MB安全阈值 */privatestaticfinallongBATCH_MAX_BYTES1024*1024;/** * 设备上报心跳入口——写入缓冲队列 */publicvoidreceiveHeartbeat(DeviceHeartbeatheartbeat){if(!buffer.offer(heartbeat)){log.warn(心跳缓冲队列已满丢弃心跳 | deviceId{},heartbeat.getDeviceId());// 队列满了说明消费速度跟不上可以做降级直接丢弃或告警}}/** * 定时任务每5秒批量发送一次 */Scheduled(fixedDelay5000)publicvoidflushBatch(){ListDeviceHeartbeatbatchnewArrayList(BATCH_MAX_SIZE);longbatchSize0;while(!buffer.isEmpty()){DeviceHeartbeatheartbeatbuffer.poll();if(heartbeatnull){break;}// 估算单条消息大小JSON序列化后约200字节intmsgSizeestimateSize(heartbeat);// 超过批次大小或字节数限制先发送当前批次if(batch.size()BATCH_MAX_SIZE||batchSizemsgSizeBATCH_MAX_BYTES){sendBatch(batch);batch.clear();batchSize0;}batch.add(heartbeat);batchSizemsgSize;}// 发送剩余的if(!batch.isEmpty()){sendBatch(batch);}}/** * 批量发送核心方法 */privatevoidsendBatch(ListDeviceHeartbeatbatch){if(batch.isEmpty()){return;}try{// 将心跳列表序列化为JSON数组作为一条消息发送StringpayloadJSON.toJSONString(batch);MessageStringmessageMessageBuilder.withPayload(payload).build();SendResultresultrocketMQTemplate.syncSend(device-heartbeat-topic,message);log.info(心跳批量发送成功 | count{} | size{}bytes | msgId{},batch.size(),payload.length(),result.getMsgId());}catch(Exceptione){log.error(心跳批量发送失败 | count{},batch.size(),e);// 发送失败的处理策略写本地文件兜底后续补偿fallbackToFile(batch);}}/** * 估算单条心跳消息序列化后大小 */privateintestimateSize(DeviceHeartbeatheartbeat){// 粗略估算JSON序列化后约200-300字节return300;}/** * 降级写本地文件 */privatevoidfallbackToFile(ListDeviceHeartbeatbatch){try{StringjsonJSON.toJSONString(batch);StringfileNameheartbeat_fallback_System.currentTimeMillis().json;// 写入本地文件后续补偿任务读取重发Files.write(Paths.get(./fallback/fileName),json.getBytes(StandardCharsets.UTF_8));}catch(IOExceptione){log.error(降级写文件也失败了数据丢失,e);}}}这里没有用RocketMQ原生的MessageBatch而是把多条心跳序列化为JSON数组放进一条消息。这种方式更灵活不依赖RocketMQ客户端的批量拆包机制消费端拿到的就是一条包含数组的大消息直接反序列化即可。3.4 消费侧批量写入InfluxDB消费端拿到批量心跳后一次性写入InfluxDB避免逐条插入的IO开销。Slf4jServiceRocketMQMessageListener(topicdevice-heartbeat-topic,consumerGroupheartbeat-consumer-group,consumeModeConsumeMode.CONCURRENTLY,consumeThreadMax10,// 关键一次消费一批消息虽然这里是一条大消息但可以调大批量拉取数consumeMessageBatchMaxSize1)publicclassHeartbeatBatchConsumerimplementsRocketMQListenerMessageExt{ResourceprivateInfluxDBMapperinfluxDBMapper;OverridepublicvoidonMessage(MessageExtmessage){try{// 反序列化批量心跳ListDeviceHeartbeatheartbeatsJSON.parseArray(newString(message.getBody(),StandardCharsets.UTF_8),DeviceHeartbeat.class);log.info(收到心跳批量 | count{} | msgId{},heartbeats.size(),message.getMsgId());// 批量写入InfluxDBbatchWriteToInfluxDB(heartbeats);}catch(Exceptione){log.error(心跳批量消费失败 | msgId{},message.getMsgId(),e);thrownewRuntimeException(心跳批量消费失败,e);}}/** * 批量写入InfluxDB */privatevoidbatchWriteToInfluxDB(ListDeviceHeartbeatheartbeats){BatchPointsbatchPointsBatchPoints.database(vending_device).retentionPolicy(autogen).consistency(ConsistencyLevel.ALL).build();for(DeviceHeartbeathb:heartbeats){PointpointPoint.measurement(device_heartbeat).tag(deviceId,hb.getDeviceId()).addField(temperature,hb.getTemperature()).addField(signalStrength,hb.getSignalStrength()).addField(stockCount,hb.getStockCount()).addField(faultCode,hb.getFaultCode()).addField(online,hb.getOnline()).time(hb.getTimestamp(),TimeUnit.MILLISECONDS).build();batchPoints.point(point);}influxDBMapper.write(batchPoints);log.info(批量写入InfluxDB完成 | count{},heartbeats.size());}}3.5 consumeMessageBatchMaxSize参数这里解释一下consumeMessageBatchMaxSize。这个参数控制的是消费者一次拉取多少条消息交给业务处理。默认值是1即每次拉取一条。如果用RocketMQ原生的MessageBatch发送每条消息独立消费端可以调大这个参数一次拉取多条批量处理RocketMQMessageListener(topicdevice-heartbeat-topic,consumerGroupheartbeat-consumer-group,consumeMessageBatchMaxSize32// 一次拉取32条消息)publicclassHeartbeatBatchConsumer2implementsRocketMQListenerListMessageExt{// 注意接口泛型变成ListMessageExtOverridepublicvoidonMessage(ListMessageExtmessages){// 一次处理32条消息}}我们的场景是把多条心跳打包成一条消息发送所以consumeMessageBatchMaxSize1就够了。如果用原生批量发送方式就需要调大这个参数。四、批量发送的性能对比实际压测数据1000台设备每5秒一轮心跳共200条/秒方案发送TPSBroker CPU消费端DB写入次数端到端延迟逐条发送~200/s60-80%200次/s100ms批量发送500条/批~1批次/5s5-10%1次/5s5-6sBroker CPU从80%降到10%数据库写入从每秒200次降到每5秒1次。端到端延迟从100ms变成5-6秒——但心跳数据对实时性要求没那么高5秒延迟完全可接受。五、批量消息的注意事项5.1 攒批窗口的权衡攒批时间越长批次越大吞吐越高但延迟也越高。需要根据业务容忍度选择心跳上报5-10秒窗口容忍延迟订单消息不能攒批需要实时处理日志上报1-3秒窗口兼顾实时性和吞吐5.2 优雅停机攒批方案最大的风险服务突然停掉缓冲队列里的消息就丢了。必须做优雅停机PreDestroypublicvoidgracefulShutdown(){log.info(服务关闭 flush剩余心跳...);flushBatch();// 停机前把缓冲队列里的消息全部发出去log.info(心跳flush完成服务安全关闭);}5.3 批次大小不要卡4MB上限虽然4MB是上限但批次太大有以下问题网络传输时间长失败重传成本高Broker单条消息过大影响CommitLog写入性能消费端反序列化大JSON内存占用高建议单批次控制在1MB以内条数控制在500-1000条。5.4 批量发送 vs 批量消费这两个概念不要混淆概念发送端消费端批量发送攒多条打包成一条发送收到一条大消息内部包含多条批量消费逐条发送一次拉取多条消息循环处理我们的方案是批量发送攒批JSON数组消费端收到一条大消息内部反序列化出多条。如果用RocketMQ原生的MessageBatch则可以实现批量发送批量消费的组合。六、总结要点说明批量发送本质多条消息打包成一条减少网络IO和Broker存储开销大小限制单批次不超过4MB建议控制在1MB以内攒批窗口心跳场景5秒根据业务延迟容忍度调整降级策略发送失败写本地文件兜底后续补偿重发优雅停机PreDestroy中flush缓冲队列防止数据丢失消费端优化批量写入InfluxDB一次BatchPoints写入多条售货柜心跳上报的核心链路设备上报 → 网关5秒攒批 → JSON数组打包发送 → 消费者反序列化 → BatchPoints批量写InfluxDB。Broker压力降一个数量级数据库写入次数从每秒几百次降到每5秒一次整体链路稳如老狗。
返回列表