HTTP 返回以后谁唤醒 outbox

第 23 篇的订单与 outbox 已在单库事务中创建,第 30 篇已有内嵌持久队列。当前的受管任务是显式触发一次:HTTP POST /procurement/api/lab/broker/schedule?tenant=... 在取得单在途令牌后将一次 sendNext(tenant,false) 提交给 ManagedExecutorService,立即返回 202。把 new Thread(...) 放到 /order 方法末尾依然错误:线程绕开容器生命周期,内存中待执行任务随部署消失,原订单事务也不会跨线程继续运行。当前实现没有周期唤醒、数据库租约或自动扫描,不能把一次 execute() 写成可持久恢复的定时发布器。

代码入口是 BrokerLabResource.schedule:通过 @Resource(lookup="java:comp/DefaultManagedExecutorService") 注入执行器;Semaphore(1) 限制单实例同时在途任务,忙时响应 429,拒绝提交时响应 503;任务异常写日志,finally 释放令牌。配置使用消息服务器副本和已执行的 002/003 迁移;33 场景脚本等待 SQL 从 PENDING 到 PUBLISHED,再同步收消息、查一次收据。它不证明应用启动后无需请求即可自动轮询。

容器提供线程,应用定义负载上限

Jakarta Concurrency 3.1 的 ManagedExecutorService 负责受管环境执行任务;当前 execute() 一次只投递一项任务,不提供固定触发间隔。Semaphore(1) 控制的是本 JVM 的在途数,不是服务器线程池容量或集群级配额;测试只按一次手工触发成功收集证据。FUTURE/NOT_RUN:周期触发、单轮 50、双实例协调、ManagedScheduledExecutorService、应用重启后无人触发的补扫,都须独立实现/验收。

1
2
3
HTTP schedule ──> tryAcquire(1) ──> execute(sendNext) ──> 内嵌 JMS 提交 ──> SQL PUBLISHED
│ │
202 finally 释放令牌;异常仅写日志,需查 DB 终态

受管不意味着自动继承请求事务或受信租户。当前 schedule 把未经认证的 tenant 查询参数捕获到 lambda,worker 直接调用同一资源实例的 sendNext,而非经过受管业务服务;尚未记录后台线程的认证主体、请求作用域或事务传播。订单写入与后台 JDBC/JMS 操作有独立的执行时刻,sendNext 不会因为自调用就变成有容器事务拦截的服务。未来需从持久事件重新读取并核实受信租户,分别测量上下文,不可把演示参数当作已验证身份。

当前代码没有保存 Future、取消接口或优雅停机后的补扫逻辑;availableWorkers.release() 仅释放本地令牌,不标记数据库任务完成。worker 如果在发消息后、写 PUBLISHED 前失败,event 仍可被下一次人工请求再次发送,依第 30 篇的收据去重;任务若在 JVM 队列里尚未开始就停机,也不会自动重排。FUTURE/NOT_RUN:周期补扫、取消/关闭、租约超时回收及多实例协调;调用 Future.cancel(true) 也不可能撤回已经提交的 JMS 消息。

任务在库,唤醒尚靠请求

当前订单事务写入 outbox;只有请求 /schedule 才提交一次 worker,该 worker 在执行时查某租户一条 PENDING。HTTP 客户端提交 202 后若进程退出,尚未运行的 lambda 不会自动恢复;订单和 outbox 仍在数据库,但需要再次人工触发才能发布。PUBLISHED 是发送进度,不是消费收据,更不是邮件送达。当前审计应分别查 notification_outbox 的 PENDING/PUBLISHED 和 notification_receipt 的 event ID,而不是检查已经返回的 202 或 JVM 中的 Future。

202 甚至不承诺 worker 在 HTTP 应答后才开始:执行器可能立即运行,也可能排队;应答只说明 execute() 当时接受了提交。出错后的处理不能只靠响应码:worker 在 sendNext 里抛出的异常被记录进服务器日志,调用者仍已拿到 202。第一次 SQL 查询看见 PENDING,可能是任务未开始,也可能是消息已提交而 markPublished 尚未执行或已失败。重复请求 /schedule 前先保留此时服务器日志与数据库状态;再次调用可能发出同一 event ID 的另一条消息。本次只有一次无故障 worker 的原始记录,没有证明发生 202 后的这些错误分支。

FUTURE/NOT_RUN:若改为每 5 秒触发一次,单轮耗时 12 秒就会重入;需在领取前限制在途数,容量满时跳轮并保留数据库里的 PENDING。双实例不能共用 JVM 信号量;须在新增迁移中定义租约字段并以数据库条件领取,实例 A 长暂停、B 接管后,A 必须校验 token 才能回写。两次发送仍可能发生,须靠稳定 event ID 与收据限制本地重复记账。002-notification-outbox.sql 没有 lease_until、lease_token、next_attempt_at,当前不提供假装可执行的租约 SQL。

即使补上租约,JMS 提交与 outbox 标记仍非一笔可假定的 XA 事务;领取失败后的原因与退避时间也必须持久记录。现有 attempts 列不能替代租约或退避时间。

哪些上下文能过线程边界

受管执行器提供在服务器控制之下运行任务的机制;哪些上下文可用须按 Jakarta Concurrency 3.1 和实际配置区分。当前 lambda 捕获不可信演示查询参数 tenant,而 sendNext 查询数据库时按这个字符串过滤 PENDING。没有认证主体、请求作用域或事务上下文的传播测量,不能从 @Resource 注入成功推断其安全性。FUTURE/NOT_RUN:执行时从持久任务读取并核验受信租户与订单绑定,审计用户身份需定义允许保存的主体标识与重新授权规则,不能把演示头当认证身份。

事务上下文尤其不能假定跨线程传播。ProcurementUseCases.order() 的 @Transactional 并不意味着后来提交的后台任务获得原事务或 JDBC 连接。当前 sendNext 内的 JMS 事务与两个 JDBC 操作独立:JMS commit 后、SQL PUBLISHED 前失败仍留 PENDING。FUTURE:若迁至受管事务服务,应走容器可拦截入口而非同类直接自调用;不要跨线程共享 Connection 或 JMSContext,还要分别测试事务边界与关闭时机。

FUTURE/NOT_RUN:取消或停机前后分别记录 JMS 接受和数据库终态;Future.cancel(true) 只是中断请求,不能撤回已提交的消息。若增加自动恢复,要验证强停时清理钩子未运行、新实例如何领取未完成项,而不是仅靠关闭回调释放本地令牌。

怎样控制容量而不饿死业务

当前容量上界只在应用代码中体现:Semaphore(1) 使同一资源实例同时只能提交一项,tryAcquire() 失败响应 429,execute() 拒绝响应 503;这并非执行器线程池或容器队列容量的测量,也不说明两个实例的总并发。FUTURE/NOT_RUN:“单轮 50、在途 2、间隔 5 秒”是待实现的容量假设。若每秒新增十项、每任务发送耗时三秒,两名 worker 最多处理不到一项/秒,积压必增;要测量待处理年龄、发送耗时、资源池借还与实际峰值,再设批次和退避,不能靠盲目增大线程数解决消息服务故障。

信号量先于任务提交获取,在 worker 的 finally 中释放;同步的 RejectedExecutionException 路径也显式释放。这样可以限制当前资源实例中排队加运行的 worker 总数,却没有给数据库行上锁。/send-next 仍能绕过 /schedule 直接被调用,两个 HTTP 请求若指向同一个 PENDING 行,信号量不会阻止重复发送。遇到 429 只能判断本实例此刻没有空余许可,不能推断任何 event 已由另一个线程完成;遇到 503 说明执行器拒绝这次提交,应先确认数据库里依旧存在待处理任务。两条分支在 33 场景中都未触发,不能报告其实际恢复耗时。若在 202 响应后进程停止而任务尚在容器队列中,信号量的内存状态也随进程消失;重新启动应先按 outbox 行查询是否仍 PENDING,再决定重试,不能凭旧响应补写 PUBLISHED。现有脚本没有在这个窗口强停实例,恢复断言仍待实测。

此外,即使一个进程的信号量始终只有一个许可,连续两次 HTTP 请求仍可能先后挑中同一行:第一次 JMS commit 后若 markPublished 抛错,许可被释放,第二次又读到同一个 PENDING。这与两个线程同时执行是不同的重复来源;限制并行数不能替代持久进度或消费去重。若同一运行时还部署了其他应用,默认受管执行器可能承载别的工作,代码里的本地许可能控制本应用提交数,却不能用来推断整个容器线程池的空闲数量。测试池拒绝或公平性必须另有容器配置与并行负载证据。

FUTURE/NOT_RUN:周期轮询要限制空轮数据库访问、单次领取数量和拒绝后的再次扫描;不要无限排队。在同一小池混跑批任务或慢 HTTP 也会造成饥饿,需要结合资源预算隔离。现有 33 脚本既没轮询 51 项,也没模拟池满;它证明的是一次受管提交经独立数据库连接观察到进度,并由另一个 HTTP 请求同步消费成功。

当前配置可复核的边界

按照工程说明,使用 JDK 21、Open Liberty 26.0.0.5、PostgreSQL 16.15/pgJDBC 42.7.7、独立 javaee_lab,先按顺序执行 001、002、003 迁移各一次,构建 WAR 并部署到消息配置 deploy/server-messaging.xml(不是基础 server.xml)。需要本地 curl、jq、psql;JAVAEE_DEMO_MODE=true 及 HTTP/JMS 回环监听只用于隔离教学,口令不得入库。现行表可查的是 notification_outbox.status 与 notification_receipt.event_id,没有租约 token。

1
2
3
4
5
6
7
cd examples/javaee-enterprise
../hibernate-lab/mvnw -B -ntp -f "$PWD/pom.xml" clean verify
rg -n 'jakartaee-11.0|messagingEngine|jdbc/Procurement' deploy/server-messaging.xml
PGPASSWORD="$JAVAEE_LAB_PASSWORD" psql -X -v ON_ERROR_STOP=1 \
-h 127.0.0.1 -U javaee_lab -d javaee_lab \
-Atc "SELECT current_database(), to_regclass('public.notification_outbox'), to_regclass('public.notification_receipt')"
/path/to/wlp/bin/productInfo version

查询应显示 javaee_lab|notification_outbox|notification_receipt,否则先排查库或迁移;productInfo 路径须换成本机安装位置。这些静态检查不能替代运行脚本,更不能从 jakartaee-11.0 推出轮询器已经部署。若还没复制 WAR/配置并启动服务器,则以下运行步骤均为 NOT_RUN。

如需在脚本之外读取同一条 event 的终态,可从脚本打印的 event_id 取实际数值,针对隔离库只读查询;不要把原始记录中的 6 直接当成复跑的 ID:

1
2
3
4
5
6
7
PGPASSWORD="$JAVAEE_LAB_PASSWORD" psql -X -v ON_ERROR_STOP=1 \
-h 127.0.0.1 -U javaee_lab -d javaee_lab -v event_id=6 -At <<'SQL'
SELECT o.id, o.status, o.attempts, COUNT(r.event_id)
FROM notification_outbox o LEFT JOIN notification_receipt r ON r.event_id = o.id
WHERE o.id = :event_id
GROUP BY o.id, o.status, o.attempts;
SQL

预期在事件处理结束后得到 id|PUBLISHED|attempts|1;现行 sendNext 不递增 attempts,不能以 attempts 数字断言 JMS 只发了一次。若读不到行,先核对脚本真正生成的 ID 与连接的库;若仍是 PENDING,检查 worker 异常与 JMS 提交/SQL 回写之间的断点,不能直接把状态更新成 PUBLISHED。若 PUBLISHED 但收据为零,先检查显式 /receive-one 请求和队列配置。这个查询只能核实数据库,不会显示服务器线程的存活状态。

正常与故障检查

在上述环境已启动且迁移完成后,复跑33 场景脚本:

1
2
3
4
cd examples/javaee-enterprise
../hibernate-lab/mvnw -B -ntp -f "$PWD/pom.xml" clean verify
JAVAEE_PORT=9085 JAVAEE_DEMO_MODE=true JAVAEE_LAB_USER=javaee_lab \
JAVAEE_LAB_PASSWORD='<本地专用口令>' bash scenarios/33-managed-outbox.sh

脚本先创建订单及一项 PENDING,请求 /schedule 预期 202,然后最多等十次、每次一秒查 PUBLISHED;显式调用 /receive-one 预期 200、newReceipt=true,SQL 收据数为 1。任何超时/状态不符都退出非零;202 不保证后续完成。现存原始记录在 verification/20261004T-managed-outbox/:event 6、accepted-not-completed=HTTP202、6|PUBLISHED、consume-ledger=HTTP200、receipt_rows=1,场景退出码 0;构建及部署 WAR SHA-256 同为 691a0acb9811be0a1d3ee1c9ede81ed33e5f5a2ce1e0f7b6c452a6f67969a0d4。这只对应当时的 WAR、单次显式触发与同步消费,不替代未来复跑的构建/部署对照。对应合同见 执行器实验合同。

这里的 SQL 轮询是测试脚本在客户端做的等待,不是服务器后台定时扫描;脚本轮询最多约十秒,超时意味着该次测试没有看到发布完成,不能断言消息一定没进入队列。可另查 SQL 进度和服务器日志来定位失败点,并对队列中同 ID 消息做可核对的隔离操作;继续盲发一次 /schedule 会把未知结果当成失败结果。若环境变更,先比较构建 WAR 与实际部署 WAR SHA,再核查配置是否确有 messagingEngine 与 jms/ProcurementCF。Java 线程已经执行、消息确实落盘、消费收据出现三件事必须各自取证。

对调度本身做故障注入

FUTURE/NOT_RUN:扩成周期 51 项验收时,第一轮领取至多 50、单实例同时发布至多 2,余下一项可由后续轮次领取;记录峰值和每个 event 的 broker/收据终态。失败注入应在领取后停机,重启先核对 PENDING 集合,再等新增租约到期并验证恢复;拒绝提交时必须保留可再领的持久任务。不同设计可选“先取令牌再领租约”或“失败后释放租约”,只有实现和新增迁移到位才能测试,不能拿现行 002 表直接跑。

FUTURE/NOT_RUN:模拟 JMS 暂时不可用,记录 worker 异常和 PENDING 的持久状态;补齐退避字段后,再验下一次尝试时间。双实例租约过期后可能两次发送同一 ID,收据最多一条;这仍不说明真实通知只发送一次。现有脚本没有断开 broker、强停受管任务、重启自动补扫或取消注入;这些均不能用第 30 篇同步发送后的失败注入代替。

FUTURE/NOT_RUN:上下文测试须分别观测请求端与 worker 的主体、租户来源和事务状态;取消与停机测试须保存接受响应、实际进度与重启后数据。演示入口没有身份认证,不能拿 X-Lab-Tenant 当合法主体;实现后还应记录拒绝次数、池容量、空轮耗时和租约年龄。这些指标当前没有原始记录,不能由一条成功的 202 推算。

两道练习

练习一:单机把池大小固定为 2,就能保证双实例最多两条任务同时发吗?解:不能。每个实例可各有两个线程;应以持久租约和数据库条件领取防止同一任务并发处理,整体吞吐上限还要有跨实例治理。对重复投递保留消费去重。

练习二:Future.cancel(true) 返回 true,是否意味着 broker 不可能收到消息?解:不意味着。中断是对线程的取消请求,网络发送可能已经完成;任务状态需对账并重试/去重,不能把 cancel 的返回值当作 broker 的撤销凭据。

FUTURE/NOT_RUN:滚动升级可能让新旧实例同时扫描;需要单独测试 lease token、防止旧实例回写及新旧消息格式兼容。当前既没有这些租约字段,也未做跨实例容量实验,单机一次受管发送不能证明分布式安全。

适用范围与版本资料

一次执行的证据只覆盖本地单在途发布、数据库进度和同步消费;51 项、上下文传播、周期轮询、取消/停机恢复、多实例租约均 NOT_RUN。规范:Jakarta Concurrency 3.1:ManagedExecutorService 与事务管理、Jakarta Transactions 2.0;应用平台:Jakarta EE 11;具体配置参照 Open Liberty 26.0.0.5 jakartaee-11.0 功能。