news 2026/9/7 17:53:09

查券返利机器人的异步任务调度:Java XXL-Job+Redis实现海量查券请求的分布式任务分发

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
查券返利机器人的异步任务调度:Java XXL-Job+Redis实现海量查券请求的分布式任务分发

查券返利机器人的异步任务调度:Java XXL-Job+Redis实现海量查券请求的分布式任务分发

大家好,我是 微赚淘客系统3.0 的研发者省赚客!

在高并发场景下,用户通过查券返利机器人发起的优惠券查询请求可能瞬时达到数十万量级。为避免直接冲击核心接口并保障系统稳定性,我们采用“请求入队 + 异步消费”模式,基于 XXL-JOB 作为分布式调度中心,结合 Redis Stream 实现任务的可靠分发与削峰填谷。

整体架构设计

用户请求首先写入 Redis Stream,由多个消费者实例监听流数据;XXL-JOB 定时触发任务拉取器,动态调整消费速率,并支持失败重试、积压告警与人工干预。该方案解耦了请求入口与处理逻辑,提升系统弹性。

Redis Stream 任务队列定义

每个查券请求封装为一个结构化消息,写入名为coupon:query:stream的 Stream:

packagejuwatech.cn.task.model;publicclassCouponQueryTask{privateStringtaskId;// 全局唯一IDprivateStringuserId;// 用户IDprivateStringitemId;// 商品IDprivatelongtimestamp;// 请求时间戳privateintretryCount;// 重试次数// getters and setterspublicStringgetTaskId(){returntaskId;}publicvoidsetTaskId(StringtaskId){this.taskId=taskId;}publicStringgetUserId(){returnuserId;}publicvoidsetUserId(StringuserId){this.userId=userId;}publicStringgetItemId(){returnitemId;}publicvoidsetItemId(StringitemId){this.itemId=itemId;}publiclonggetTimestamp(){returntimestamp;}publicvoidsetTimestamp(longtimestamp){this.timestamp=timestamp;}publicintgetRetryCount(){returnretryCount;}publicvoidsetRetryCount(intretryCount){this.retryCount=retryCount;}}

生产者将任务推入 Redis:

packagejuwatech.cn.task.producer;importcom.fasterxml.jackson.databind.ObjectMapper;importorg.springframework.data.redis.connection.stream.StreamRecords;importorg.springframework.data.redis.core.RedisTemplate;importorg.springframework.stereotype.Component;importjava.util.HashMap;importjava.util.Map;@ComponentpublicclassCouponTaskProducer{privatefinalRedisTemplate<String,Object>redisTemplate;privatefinalObjectMapperobjectMapper=newObjectMapper();publicCouponTaskProducer(RedisTemplate<String,Object>redisTemplate){this.redisTemplate=redisTemplate;}publicvoidsubmitTask(juwatech.cn.task.model.CouponQueryTasktask)throwsException{Map<String,String>payload=newHashMap<>();payload.put("task",objectMapper.writeValueAsString(task));redisTemplate.opsForStream().add(StreamRecords.newRecord().ofObject(payload).withStreamKey("coupon:query:stream"));}}

XXL-JOB 调度任务配置

在 XXL-JOB 控制台注册执行器juwatech-coupon-job,并创建任务coupon-query-consumer,Cron 表达式设为0/5 * * * * ?(每5秒触发一次)。

对应的 JobHandler 实现如下:

packagejuwatech.cn.job;importcom.xxl.job.core.context.XxlJobHelper;importcom.xxl.job.core.handler.annotation.XxlJob;importorg.springframework.stereotype.Component;@ComponentpublicclassCouponQueryConsumerJob{privatefinaljuwatech.cn.task.consumer.CouponTaskConsumerconsumer;publicCouponQueryConsumerJob(juwatech.cn.task.consumer.CouponTaskConsumerconsumer){this.consumer=consumer;}@XxlJob("couponQueryConsumer")publicvoidexecute(){try{intconsumed=consumer.consumeBatch(100);// 每次最多消费100条XxlJobHelper.handleSuccess("Consumed "+consumed+" tasks");}catch(Exceptione){XxlJobHelper.handleFail(e.getMessage());}}}

Redis Stream 消费逻辑

消费者从 Stream 读取待处理任务,并调用查券服务:

packagejuwatech.cn.task.consumer;importcom.fasterxml.jackson.databind.ObjectMapper;importorg.springframework.data.redis.connection.stream.*;importorg.springframework.data.redis.core.RedisTemplate;importorg.springframework.stereotype.Service;importjava.time.Duration;importjava.util.Collections;importjava.util.List;importjava.util.Map;@ServicepublicclassCouponTaskConsumer{privatefinalRedisTemplate<String,Object>redisTemplate;privatefinaljuwatech.cn.service.CouponQueryServicecouponQueryService;privatefinalObjectMapperobjectMapper=newObjectMapper();publicCouponTaskConsumer(RedisTemplate<String,Object>redisTemplate,juwatech.cn.service.CouponQueryServicecouponQueryService){this.redisTemplate=redisTemplate;this.couponQueryService=couponQueryService;}publicintconsumeBatch(intbatchSize)throwsException{StreamReadOptionsoptions=StreamReadOptions.empty().count(batchSize).block(Duration.ofSeconds(2));// 阻塞等待新消息List<MapRecord<String,Object,Object>>records=redisTemplate.opsForStream().read(options,StreamOffset.create("coupon:query:stream",ReadOffset.lastConsumed()));if(records==null||records.isEmpty()){return0;}for(MapRecord<String,Object,Object>record:records){try{Map<Object,Object>value=record.getValue();StringtaskJson=(String)value.get("task");juwatech.cn.task.model.CouponQueryTasktask=objectMapper.readValue(taskJson,juwatech.cn.task.model.CouponQueryTask.class);// 执行查券couponQueryService.queryAndNotify(task.getUserId(),task.getItemId());// 确认消费(删除消息)redisTemplate.opsForStream().acknowledge("coupon:query:stream","coupon-group",record.getId());}catch(Exceptione){handleConsumeFailure(record,e);}}returnrecords.size();}privatevoidhandleConsumeFailure(MapRecord<String,Object,Object>record,Exceptione){// 可选:将失败任务写入重试队列或 DLQjuwatech.cn.util.AsyncLogger.logAsync("Failed to process task "+record.getId()+": "+e.getMessage());}}

消费者组与消息可靠性

初始化消费者组以支持多实例并行消费:

XGROUP CREATE coupon:query:stream coupon-group $ MKSTREAM

在应用启动时自动创建(可选):

packagejuwatech.cn.config;importorg.springframework.data.redis.core.RedisTemplate;importorg.springframework.stereotype.Component;importjavax.annotation.PostConstruct;@ComponentpublicclassRedisStreamInit{privatefinalRedisTemplate<String,String>redisTemplate;publicRedisStreamInit(RedisTemplate<String,String>redisTemplate){this.redisTemplate=redisTemplate;}@PostConstructpublicvoidinitConsumerGroup(){try{redisTemplate.execute((connection)->{connection.streamCommands().xGroupCreate("coupon:query:stream".getBytes(),"coupon-group".getBytes(),org.springframework.data.redis.connection.stream.ReadOffset.from("0-0").getOffset().getBytes(),true);returnnull;});}catch(Exceptionignored){// Group already exists}}}

积压监控与弹性扩缩容

通过 Redis 命令XINFO GROUPS coupon:query:stream获取 pending 消息数,在 XXL-JOB 中上报指标,结合 Prometheus + Alertmanager 实现积压告警。当 pending > 10000 时,自动扩容消费者实例。

本文著作权归 微赚淘客系统3.0 研发团队,转载请注明出处!

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/4 10:32:39

【小程序毕设源码分享】基于springboot+小程序的武设专业解读武术兴趣班小程序的设计与实现(程序+文档+代码讲解+一条龙定制)

博主介绍&#xff1a;✌️码农一枚 &#xff0c;专注于大学生项目实战开发、讲解和毕业&#x1f6a2;文撰写修改等。全栈领域优质创作者&#xff0c;博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围&#xff1a;&am…

作者头像 李华
网站建设 2026/8/25 8:44:48

如何科学地“设计”SFT 数据?一次关于 ODA 的完整平台级验证

在大模型后训练阶段&#xff0c;SFT&#xff08;监督微调&#xff09;数据的构建至关重要。然而&#xff0c;长期以来&#xff0c;这一过程业界的通行做法往往依赖“直觉”或“试错”&#xff0c;即多收一点、再筛一轮、训一次模型、看下效果&#xff0c;然后再调整。这个过程不…

作者头像 李华
网站建设 2026/9/2 23:00:51

黑客攻击MongoDB实例删除数据库并植入勒索信息

威胁行为者正通过大规模自动化勒索软件活动&#xff0c;持续攻击暴露在互联网上的MongoDB实例。攻击模式高度一致&#xff1a;攻击者扫描公网可访问的未受保护MongoDB数据库&#xff0c;删除存储数据后植入比特币勒索信息。 MongoDB实例遭入侵分析 最新证据显示&#xff0c;尽…

作者头像 李华
网站建设 2026/9/4 6:10:46

基于SSM的高校旧书交易系统的设计与实现(毕业论文)

摘 要 随着教育资源的日益丰富和高等教育的普及&#xff0c;大学生群体在学习和科研过程中产生了大量的书籍需求。然而&#xff0c;由于课程结束或毕业离校等原因&#xff0c;许多书籍在使用一段时间后便被闲置&#xff0c;这造成了大量资源的浪费。基于此&#xff0c;本文基于…

作者头像 李华
网站建设 2026/9/2 22:17:37

Vue与Web Components的集成:技术原理、实践方案与生态协同

Vue与Web Components的集成&#xff1a;技术原理、实践方案与生态协同 一、技术演进背景与核心价值 Web Components作为W3C标准化的浏览器原生组件技术&#xff0c;由Custom Elements、Shadow DOM和HTML Templates三大核心规范构成。其设计初衷在于解决Web开发中的组件复用难题…

作者头像 李华
网站建设 2026/9/5 17:23:25

GitHub项目上传、删除与协议设置:新手到高手的完整指南

GitHub项目上传、删除与协议设置&#xff1a;新手到高手的完整指南 引言 对于每一位开发者而言&#xff0c;GitHub不仅是代码的托管平台&#xff0c;更是个人技术履历和协作开发的核心。然而&#xff0c;从如何将第一个项目成功推送&#xff0c;到管理项目生命周期&#xff0…

作者头像 李华