wxyun2 opened a new pull request, #1394:
URL: https://github.com/apache/rocketmq-clients/pull/1394

   ### Which Issue(s) This PR Fixes
   
   Fixes #1392
   
   ### Brief Description
   
   同步 `MessageListener` 将方法返回与消费结果绑定:把处理任务转交其他线程后立即返回 SUCCESS,会提前推进 
ACK;等待异步任务结束再返回,又会占用消费线程。本 PR 新增可选的 
`AsyncMessageListener.consumeAsync(MessageView)`,让返回的 
`CompletionStage<ConsumeResult>` 表示实际处理结果。
   
   **处理与确认链路:**
   
   1. 消息继续通过原 PushConsumer 接收和缓存,受本地缓存条数、字节阈值控制。
   2. 在消费并发限额内调用异步监听器。监听器返回 stage 后释放执行线程,未完成的 stage 继续占据处理名额,避免异步转交导致并发失控。
   3. stage 完成后再执行消费后置 hook、统计和现有 ACK/失败处理。同步抛错、stage 异常、取消、null stage/result 
统一转换为 FAILURE。
   4. 消息缓存保留到终结操作完成;FIFO 同组消息等待当前处理及终结操作完成后继续,保留原有重试链。
   5. 关闭时停止接收准入,等待已接受响应的交接、未完成 stage 及 ACK/NACK/DLQ,再关闭执行器和 RPC 
资源。过滤与损坏消息也纳入排空,避免接收计数归零后的交接窗口漏掉消息。
   
   同步与异步 listener 按最后一次 setter 调用选择。原同步 API 保留,新接口通过公开 Builder 接入,无需额外 Consumer 
子类或反射适配。
   
   **运行期调参与异步消费共用同一套准入逻辑:**
   
   `PushConsumer.updateRuntimeTuning(cacheCount, cacheBytes, 
consumptionConcurrency)` 
在原实例上更新缓存限额和处理并发。扩容允许更多排队任务进入;缩容不打断已有任务,后续任务等待活跃数量降到新限额以下。所有参数先校验,非法参数不会部分生效。平台线程、虚拟线程及其回退模式均支持这一行为。
   
   这两项能力共同解决“任务转交后如何限制未完成处理,并在运行中调整准入”的问题,因此放在同一份 PR 中。此 PR 可独立合并,不依赖内部线程池配置 PR。
   
   **完成与续期边界:**
   
   - stage 完成表示应用已给出处理结果,不表示 ACK RPC 已成功;仍为至少一次投递,需要应用幂等。
   - 沿用现有服务端 `autoRenew=true` 协议及服务端时长限制。
   - `close()` 等待 stage 终结,应用应确保其最终完成,并在 Consumer 关闭后再关闭应用执行器。
   
   包含公开 API、接入示例、README 及 JDK 21 CI 覆盖。
   
   ### How Did You Test This Change?
   
   在此 PR 的独立分支、macOS arm64 上验证:
   
   - JDK 17.0.20.1:完整 `mvn -B package` 成功,364 项测试,0 失败、0 错误,1 
项上游已有测试跳过;Checkstyle、SpotBugs 及普通/shaded JAR 构建通过。
   - JDK 21.0.11:执行本 PR 更新后的 CI 专项命令成功,91 项测试,0 失败、0 错误,1 项条件测试跳过。
   - 生产代码与 API 按上游 Java 8 release 目标编译。
   
   JDK 21 复现命令:
   
   ```sh
   mvn -B -pl test -am -Dspotbugs.skip=true -Dnet.bytebuddy.experimental=true 
-Dtest=ClientConfigurationTest,ExecutorServicesTest,AsyncConsumeTaskTest,AsyncConsumeServiceTest,PushConsumerBuilderImplTest,PushConsumerRuntimeTuningTest,AsyncReceiveLifecycleTest,AsyncPushConsumerIntegrationTest
 -DfailIfNoTests=false -Dsurefire.failIfNoSpecifiedTests=false test
   ```
   
   测试覆盖 stage 延迟完成、成功/失败/异常/取消/null、执行器拒绝、延迟调度、未完成 stage 
的并发限制、活跃任务扩缩容、参数校验、FIFO 顺序、接收交接和确认排空。
   
   新增 3 项本地 gRPC 集成测试,通过公开 Builder 启动客户端,验证未完成前无 ACK、完成后单次 ACK、失败后的 
ChangeInvisibleDuration 重试及 close 等待在途 ACK 响应。服务仅绑定 `127.0.0.1`;未连接真实 
Broker,也不将这些测试解释为生产吞吐或续期上限验证。
   
   JDK 21 使用上游现有 Byte Buddy 的测试兼容选项 `-Dnet.bytebuddy.experimental=true`。
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to