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

资讯详情

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

Java实现无限循环任务:从while(true)到可治理的常驻消费模型

Java实现无限循环任务:从while(true)到可治理的常驻消费模型 拿到“以防你不知道活动关 EX-7 可以无限循环”这个题目时很容易联想到游戏活动关卡的攻略分享。但从工程视角看这句话里真正值得展开的并不是某个具体关卡怎么打而是一个在任务系统里反复出现的通用问题如果一项任务需要长期运行、循环消费、持续处理我们应该怎么设计才能做到既能无限循环又不出现死循环、不丢数据、不重复消费、不把人困在日志里排查到深夜这篇文章不讨论具体游戏的关闭规则也不教如何通过脚本重复刷取奖励。它把“EX-7”当成一个活动任务编号把“无限循环”当成一种任务消费模式用一个 Java 实现的常驻消费任务来演示循环任务的核心机制、停止条件、幂等设计、异常恢复、监控排错以及生产环境和学习环境的差异。读完之后你可以把这套思路迁移到定时对账、消息消费、状态轮询、奖励发放等常见业务场景。1. 先把“EX-7 无限循环”转化为可设计的任务模型1.1 无限循环场景并不少只是大家常叫它“常驻任务”实际开发中有很多任务不是跑一次就结束的。活动期间持续发放奖励、订单状态轮询、延迟任务扫描、消息队列消费、库存超时释放都是同一个抽象某个消费者进程要一直运行不断从队列或数据库里取出待处理数据执行完再继续取下一条。这种任务通常被叫“常驻任务”“守护任务”或“消费者任务”。它的“无限循环”不是代码里的while(true)这么简单而是一套包含获取任务、处理、状态回写、失败重试、退出恢复的完整机制。学习环境里很多人为了快速跑通会写下这样一个循环while (true) { Task task getNextTask(); handle(task); }这段代码能运行但它没有停止条件、没有异常处理、没有空轮询抑制、没有状态回写。一旦handle(task)抛出 RuntimeException整个线程直接死掉一旦没有任务CPU 会被空转打满。所谓“无限循环”不是无限空转而是无限地、有节奏地处理有效任务。1.2 把一句话标题拆解成任务字段如果把“活动关 EX-7 可以无限循环”当成一条业务需求它可以拆成下面这个任务模型字段示例值说明taskNoEX-7任务编号标记是哪一类活动任务bizKeyACTIVITY_2025_EX_7_USER_888888业务幂等键标识同一件事不能被重复处理statusPENDING / PROCESSING / SUCCESS / FAILED当前处理状态payload用户ID、活动ID、发放数量等处理时需要的业务数据retryCount0已重试次数nextRetryTime2025-06-01 10:00:00下次可重试时间这里最关键的是bizKey。因为循环任务会一直接收数据如果不能区分“这条数据是否已经被处理过”重复消费就无法避免。无限循环的“无限”指的是消费者进程的运行时长而bizKey负责保证业务上的“每个事件只处理一次”。1.3 “无限”是运行时长“循环”是处理方式二者都要有约束对技术人来说标题里的“可以无限循环”很容易让人兴奋但落到生产环境任何无限运行的程序都需要约束运行时间要可管理能优雅停止停止后已处理任务不重复、未处理任务不丢失。循环速度要可控制没有任务时需要有退避或阻塞等待不能用空轮询打满 CPU。失败行为要可预期单个任务失败不能拖死整个消费者。资源占用要可观测至少能看到队列长度、处理速度、错误数量。一句话无限循环是系统能力不是代码写法。没有约束的while(true)是事故有边界的循环消费模型才是工程方案。2. 实现前先定边界停止条件、幂等键与状态流转2.1 无限循环任务不是没有退出条件进程重启、发布部署、资源不足、人工排查问题时都需要让循环停下来。如果循环逻辑里完全没有“是否继续运行”的判断就只能靠kill -9强制结束这样会带来两个问题正在处理的任务可能只处理了一半没有回写状态。内存队列中已经拉取但还没处理完的任务直接丢失。正确做法是给循环设置一个运行开关。在 Java 中常见的是volatile boolean running主线程通过开关控制消费者线程退出。停止时先让消费者停止接收新任务再把正在处理的任务处理完最后做状态落盘。public class TaskConsumer implements Runnable { private final BlockingQueueEx7Task queue; private final Ex7TaskHandler handler; private final Ex7TaskStore store; private final AtomicBoolean running new AtomicBoolean(true); public void shutdown() { running.set(false); } Override public void run() { while (running.get()) { try { Ex7Task task queue.poll(500, TimeUnit.MILLISECONDS); if (task ! null) { executeSafely(task); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } } }注意这里使用了queue.poll(500, TimeUnit.MILLISECONDS)而不是queue.take()。take()在没有任务时会一直阻塞优点是省 CPU缺点是一旦想要退出循环外部很难打断poll在超时后会返回循环体就能趁机检查running标志。这是循环任务最常见的退出实现方式之一。2.2 幂等键是循环任务的生命线无限循环最大的风险不是“跑得太久”而是同一条数据被处理了很多次。比如奖励发放任务如果消费者在处理完任务后、回写数据库之前发生了重启重启后它扫描到这条任务还在 PENDING 状态就会再发一次奖励。避免重复消费的常用手段是数据库唯一约束加业务幂等键。任务表设计成这样CREATE TABLE activity_task_record ( id BIGINT PRIMARY KEY AUTO_INCREMENT, task_no VARCHAR(32) NOT NULL, biz_key VARCHAR(128) NOT NULL, payload VARCHAR(512) NOT NULL, status VARCHAR(16) NOT NULL, retry_count INT NOT NULL DEFAULT 0, last_error VARCHAR(1024) DEFAULT NULL, create_time DATETIME NOT NULL, update_time DATETIME NOT NULL, UNIQUE KEY uk_biz_key (biz_key) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;uk_biz_key是最后一道防线。即使消费者程序里判断状态判断错了数据库也会拒绝重复插入同一条biz_key从而避免重复发放。注意幂等不能只靠代码判断代码判断存在时区和并发问题。要把唯一约束建在数据库层代码层做业务校验两边同时存在才可靠。2.3 任务状态至少要有三种以上只有“处理中”和“完成”两个状态是不够的。循环任务执行过程中任务可能在任意一步失败消费者可能随时重启。状态字段至少要能表达PENDING待处理。PROCESSING已被消费者拉取正在处理。SUCCESS处理成功。FAILED处理失败等待重试。为什么需要 PROCESSING 状态因为消费者从任务表里捞数据时不能只捞 PENDING。如果两个消费者实例同时运行它们可能同时把同一条 PENDING 任务捞走。用 PROCESSING 状态加更新时间配合乐观锁或SELECT ... FOR UPDATE SKIP LOCKED可以显著降低并发消费冲突。3. 用 Java 实现一个最小可运行的 EX-7 循环消费模型3.1 环境准备和项目结构为了把上面的概念落成可运行代码下面用 Java 17 加 Spring Boot 3.x 做一个最小模型。这个模型不会依赖 MQ、Redis 等中间件只用 JDK 自带能力和一个数据库表方便本地跑通。准备环境依赖版本建议说明JDK17 或 21现代 Java 项目常用长期支持版本Spring Boot3.2.xWeb 不是必须但方便做健康检查接口MySQL8.x用于任务状态持久化Maven3.8项目构建工具工程结构建议如下ex7-task-consumer ├── pom.xml └── src/main/java/com/example/ex7 ├── Ex7Application.java ├── task │ ├── Ex7Task.java │ ├── Ex7TaskStore.java │ ├── Ex7TaskHandler.java │ ├── Ex7TaskConsumer.java │ └── Ex7ConsumerRunner.java └── web └── TaskController.java实际项目可以根据自己的包名调整核心是区分任务实体、任务存储、任务处理器、消费者循环、启动器五个模块。3.2 定义任务实体和状态库任务实体放到Ex7Task中public class Ex7Task { private Long id; private String taskNo; private String bizKey; private String payload; private String status; private int retryCount; private String lastError; private LocalDateTime nextRetryTime; // 省略 getter / setter }任务状态库实现下一步“从数据库取一批任务、更新任务状态”两个动作。注意取任务时要限定taskNo EX-7避免把其他任务的记录捞走。Repository public class Ex7TaskStore { private final JdbcTemplate jdbcTemplate; public Ex7TaskStore(JdbcTemplate jdbcTemplate) { this.jdbcTemplate jdbcTemplate; } public ListEx7Task fetchPendingTasks(int limit) { String sql SELECT id, task_no, biz_key, payload, status, retry_count, last_error, next_retry_time FROM activity_task_record WHERE task_no ? AND status PENDING AND next_retry_time NOW() ORDER BY id LIMIT ? ; return jdbcTemplate.query(sql, new BeanPropertyRowMapper(Ex7Task.class), EX-7, limit); } public int markProcessing(Long id) { return jdbcTemplate.update( UPDATE activity_task_record SET status PROCESSING, update_time NOW() WHERE id ?, id ); } }这里用LIMIT限制了每次拉取数量防止一次性把所有待处理任务全加载到内存。无限循环任务要尽量“小批量、多频次”而不是“大批量、少频次”。3.3 核心循环消费者阻塞队列 超时轮询 优雅停止消费者循环是整个模型的核心。它做的事情是从内存队列中取任务取到就调用 handler 处理处理成功则标记 SUCCESS处理失败则判断是否重试。Component public class Ex7TaskConsumer implements Runnable { private static final Logger log LoggerFactory.getLogger(Ex7TaskConsumer.class); private final BlockingQueueEx7Task queue; private final Ex7TaskStore store; private final Ex7TaskHandler handler; private final AtomicBoolean running new AtomicBoolean(true); public Ex7TaskConsumer(BlockingQueueEx7Task queue, Ex7TaskStore store, Ex7TaskHandler handler) { this.queue queue; this.store store; this.handler handler; } public void shutdown() { log.info(receive shutdown signal, consumer will stop after current task); running.set(false); } Override public void run() { log.info(EX-7 consumer started); while (running.get()) { try { Ex7Task task queue.poll(500, TimeUnit.MILLISECONDS); if (task null) { continue; } handleTask(task); } catch (InterruptedException e) { Thread.currentThread().interrupt(); log.warn(consumer interrupted, exit loop); break; } } log.info(EX-7 consumer stopped); } private void handleTask(Ex7Task task) { try { handler.process(task); store.markSuccess(task.getId()); log.info(task success, id{}, bizKey{}, task.getId(), task.getBizKey()); } catch (Exception e) { store.markFailed(task.getId(), e.getMessage()); log.error(task failed, id{}, bizKey{}, error{}, task.getId(), task.getBizKey(), e.getMessage(), e); } } }这里有几个容易被忽略的设计点running使用AtomicBoolean或volatile保证多线程可见性。queue.poll(500, TimeUnit.MILLISECONDS)既避免了空轮询打满 CPU又保留了退出循环的机会。InterruptedException被单独捕获并且重新设置线程中断标志这是 Java 并发编程中容易被忽略的习惯。单个任务失败不会导致循环退出因为异常在handleTask内部被捕获。3.4 处理器怎么设计才不会拖垮循环循环主体只负责取任务和调度真正的业务逻辑在Ex7TaskHandler。这里的核心原则是处理器不应该做没有超时的远程调用也不应该执行一个可能跑几个小时的大任务。Component public class Ex7TaskHandler { public void process(Ex7Task task) { // 模拟根据 payload 执行业务处理 String bizKey task.getBizKey(); if (bizKey null || bizKey.isBlank()) { throw new IllegalArgumentException(bizKey cannot be blank); } // 实际项目里这里可能是调用发放接口、更新用户权益、生成对账文件等 // 推荐给远程调用设置连接和读取超时避免线程一直卡住 System.out.println(handle EX-7 task: bizKey); } }如果业务处理本身很重不推荐把它直接放在消费者线程里而是应该拆成“读取任务、派发子线程、子线程异步处理、异步回写状态”。消费者线程只做轻量的分配和状态流转避免一个慢任务阻塞后续所有任务。3.5 启动一个长期运行的消费线程Spring Boot 启动后需要一个地方把消费者线程跑起来。这里用ApplicationRunner完成启动并注册关闭钩子实现优雅停止。Component public class Ex7ConsumerRunner implements ApplicationRunner, DisposableBean { private final Ex7TaskConsumer consumer; private final ExecutorService executorService; public Ex7ConsumerRunner(Ex7TaskConsumer consumer) { this.consumer consumer; this.executorService Executors.newSingleThreadExecutor(r - { Thread t new Thread(r, ex7-task-consumer); t.setDaemon(false); return t; }); } Override public void run(ApplicationArguments args) { executorService.submit(consumer); } Override public void destroy() { consumer.shutdown(); executorService.shutdown(); } }这里把消费者线程设置为非守护线程保证应用运行期间消费者不会被 JVM 退出逻辑提前终止。4. 验证一个无限循环任务不能只看“能跑”4.1 单元测试验证消费、失败和停止验证循环任务时不要只看“程序启动了”要验证三件事处理成功、处理失败不阻断循环、停止后循环能退出。Test void consumerShouldStopWhenShutdownCalled() throws Exception { BlockingQueueEx7Task queue new LinkedBlockingQueue(); Ex7Task task new Ex7Task(); task.setId(1L); task.setTaskNo(EX-7); task.setBizKey(TEST_1001); queue.add(task); Ex7TaskConsumer consumer new Ex7TaskConsumer(queue, store, handler); Thread thread new Thread(consumer); thread.start(); Thread.sleep(200); consumer.shutdown(); thread.join(3000); assertFalse(thread.isAlive()); }停止验证很关键。很多循环任务在生产环境停不下来就是因为消费者阻塞在take()或久等远程响应上导致发布流程卡死。4.2 压测验证吞吐和堆积把 1 万条测试任务写入数据库观察消费者每分钟处理的记录数以及队列积压趋势。至少要看两个指标TPS每秒处理任务数判断处理性能是否满足业务预期。堆积趋势队列长度随时间变化如果只增不减说明消费速度低于生产速度。简单统计可以在消费者里增加计数private final AtomicLong successCount new AtomicLong(); private final AtomicLong failCount new AtomicLong();再通过一个 HTTP 接口暴露处理数量GetMapping(/metrics/ex7) public MapString, Object metrics() { return Map.of( successCount, consumer.getSuccessCount(), failCount, consumer.getFailCount(), running, consumer.isRunning() ); }生产环境不建议自己造监控直接接入 Prometheus、Micrometer 或公司已有的监控平台即可。4.3 从日志和数据库状态确认结果验证是否成功的最终依据是数据库。处理完成后任务对应的状态应该变成 SUCCESS。如果出现大量 FAILED需要看last_error字段。建议日志格式包含三个信息任务 id、bizKey、处理结果。没有 bizKey 的日志在排查重复消费时会非常痛苦。4.4 学习环境与生产环境的差异对比关注点学习环境生产环境任务存储内存队列MySQL、MQ 或 Redis Stream循环退出手动停止测试接入发布系统优雅停机日志控制台打印结构化日志采集监控无指标、告警、大盘并发消费单线程多线程或分布式多实例幂等数据库唯一约束幂等表、分布式锁、对账任务学习环境的目标是快速理解原理生产环境的目标是不丢不重、可观测、可回滚。两者差异很大千万不要把学习环境的内存队列方案直接搬上线。5. 常见问题排查循环不跑、重复消费、堆积、假死、停不下来5.1 一张表定位五类问题问题现象常见原因检查点处理建议循环不跑消费者线程没启动启动日志、线程名检查 ApplicationRunner 是否生效任务长期不消费拉取条件不满足SQL 条件、任务状态检查 next_retry_time 和 status重复消费缺少幂等键数据库唯一约束、日志建唯一索引修复状态回写队列堆积消费速度低于生产速度TPS、数据库慢查询批量处理、增加消费线程循环假死远程调用无超时线程 dump、慢日志为远程调用设置超时时间停不下来消费者阻塞在原生锁或 IOshutdown 日志使用 poll 超时机制5.2 场景一线程启动了但任务不消费现象日志里输出了EX-7 consumer started但数据库里 PENDING 状态的任务一直不变。排查顺序检查 SQL 查询条件是否命中task_no是否写成了EX7。检查next_retry_time是否小于等于当前时间若任务初始时间写错了会一直查询不到。检查任务状态是不是 PENDING如果是 PROCESSING 或 FAILED查询条件不会命中。检查fetchPendingTasks每次取的数量LIMIT 0也会导致不消费。这类问题最常见的根因是 SQL 条件写错尤其是时间字段比较方向和字符串状态写错。5.3 场景二任务被重复消费现象同一个 bizKey 的奖励发放了多次数据库里能看到多条记录。可能原因循环里在成功后没有回写 SUCCESS任务保持在 PENDING。消费者在markSuccess之前重启导致任务重新进入待处理。并发消费时多个线程同时捞到同一条任务。解决方案数据库bizKey建唯一索引。拉取任务时使用FOR UPDATE SKIP LOCKED让并发消费者不会捞到同一条。对账任务周期性扫描 PROCESSING 超时任务确认是否真的处理成功。5.4 场景三队列一直堆积现象任务生产速度正常消费者也在运行但队列长度持续上涨。排查路径查看消费者 TPS判断是否处理太慢。查看是否存在慢 SQL尤其是状态回写语句。查看处理器里是否做了大量同步远程调用导致每个任务耗时过长。查看是否消费者线程数只有 1而任务量远超单线程处理能力。处理建议先给每个任务增加耗时统计找到最耗时的环节。再把远程调用改成批量接口或异步处理最后再考虑增加消费者线程。5.5 场景四进程看着活着循环已经假死现象进程还在健康检查也正常但任务不再处理日志长时间没有输出。假死通常不是因为循环逻辑退出而是线程卡在某个操作上。常见的坑是没有设置超时的 HTTP 调用、数据库连接池耗尽、磁盘写满导致日志阻塞。排查方式执行jstack pid抓线程栈看消费者线程停在哪个方法。检查数据库连接池活跃数。检查磁盘空间。查看 GC 日志确认是否频繁 Full GC 导致停顿。循环任务的心跳非常有用。在消费者循环里每处理 N 条任务或每 30 秒输出一条心跳日志可以让假死问题更快被感知。建议处理日志打在每个任务的关键阶段心跳日志打在循环空闲或达到阈值时。两者要分开否则任务量大时心跳会淹没在业务日志里。5.6 场景五停机后任务丢失现象发布重启后部分任务既不是 SUCCESS也不是 FAILED一直停留在 PROCESSING。原因消费者把任务标记为 PROCESSING 后还没处理完成就收到停服信号。如果代码里没有恢复机制这条任务就永远卡在 PROCESSING。解决方案在拉取任务时把 PROCESSING 状态且更新时间超过阈值的任务重新拉回来。在markSuccess前先判断任务状态防止重复回写。发布前先停止接收新任务再等待正在处理的任务完成最后退出。SQL 可以参考下面这个更新UPDATE activity_task_record SET status PENDING, retry_count retry_count 1, update_time NOW() WHERE status PROCESSING AND update_time NOW() - INTERVAL 5 MINUTE;6. 把循环任务放进生产环境监控、分布式与合规边界6.1 循环任务开发检查清单在写任何长期运行的循环任务前建议先过一遍这个清单有运行开关支持优雅停止。停止后再次启动不会重复处理已成功任务。数据库层有唯一索引兜底幂等。每次拉取任务有数量上限避免一次性加载过多。没有任务时不会空转打满 CPU。单个任务异常不会导致消费者线程退出。远程调用都设置了连接超时和读取超时。记录了任务成功、失败、心跳三类日志。暴露了处理量、失败量、队列长度等指标。对长时间处于 PROCESSING 的任务有超时回收机制。这套清单可以直接用于代码评审。如果评审时发现哪一条不满足建议先补齐再上线。6.2 从单机循环到分布式调度单机循环只能在一个进程内运行。生产环境如果有多台机器同时启动相同任务就会导致重复消费。这时必须引入分布式约束。常见做法有三类方案优点缺点适用场景数据库SELECT FOR UPDATE SKIP LOCKED实现简单无额外组件并发高时数据库压力大任务量千万以下Redis 分布式锁性能好控制灵活需要额外维护 Redis任务量较大MQ 消费者组天然分布式支持重平衡需要引入消息中间件处理逻辑适合事件驱动从单机循环切换到分布式时核心原则不变幂等键、状态流转、失败重试、停止机制全部保留只是“谁负责执行”从本机线程变成了分布式消费组。不要为了追求分布式而直接换掉存储层。先确认任务量确实需要多实例消费再考虑引入中间件。6.3 如果“无限循环”出现在奖励领取场景它就是一个需要治理的风险点文章开头提到的“活动关 EX-7 可以无限循环”在游戏玩家语境下往往意味着可以反复刷取奖励。这类“无限循环”如果真实存在通常属于业务规则漏洞而不是可以推广的技术方案。真正对的做法是在业务规则上限制同一用户同一活动的参与次数和领取次数。在数据层用用户ID加活动ID加奖励批次做唯一约束。在接口层加入频控和服务端校验不能只依赖前端按钮置灰。对账任务定期检查奖励发放记录和活动参与记录发现异常要告警。作为技术博客这篇文章要强调的是无限循环用于任务消费是合理的用于无限制领取业务资源是需要治理的风险点。两者不能混为一谈。6.4 下一步练习方向如果想继续深入可以从三个方向练手一是把内存队列替换成 Redis Stream 或 RabbitMQ模拟生产环境的异步消息模型。二是给这个循环任务增加多线程消费验证并发拉取、幂等回写和线程池退出是否依然可靠。三是引入 Micrometer 指标把处理量和失败量接入监控平台再模拟一次消费者假死练习用线程 dump 定位问题。这三步做完你对“无限循环任务”的理解会从while(true)提升到生产级任务调度的完整认知。以后再遇到“某个任务可以反复跑、需要一直跑、跑挂了还能自动恢复”这类需求你已经知道该从哪里下手了。
返回列表