Quarkus+Reactor异步过滤器中重写数据流的问题排查与修复
Quarkus Webhook签名校验过滤器问题排查与修复方案
问题背景
对接第三方Webhook服务,需在Quarkus+Reactor环境下实现签名校验过滤器,要求通过请求头与请求体完成校验,且第三方类WebhookHandlerUtility的同步方法verifySignature(signature, body)无法修改。当前控制器接收参数为null,无法进入控制器方法,同时出现多类报错。
错误分析
结合代码与报错信息,核心问题如下:
- 异步操作请求上下文非法:自定义线程池异步执行校验逻辑,主线程已完成请求流程进入响应阶段,异步线程再修改
ContainerRequestContext触发IllegalStateException: Cannot be called from response filter。 - 请求体流未缓存重置:直接读取
entityStream后未做缓存,导致后续控制器无法读取请求体,参数为null。 - EventLoop线程阻塞风险:第三方同步校验方法直接在EventLoop线程执行,会阻塞Vert.x事件循环,影响应用性能。
修复方案
- 替换自定义线程池,使用Quarkus阻塞注解:用
@Blocking标记过滤器方法,让Quarkus将任务调度到Worker线程池执行,既避免阻塞EventLoop,又保证同步处理请求上下文。 - 缓存并重置请求体流:读取请求体后保存为字节数组,确保后续可重复读取;校验通过后将处理后的payload重新封装为输入流替换原流。
- 同步处理校验逻辑:所有签名校验、payload解析和流替换操作在同一请求线程中完成,避免异步操作请求上下文引发的异常。
完整修复代码
修复后的签名过滤器
import jakarta.ws.rs.container.ContainerRequestContext import jakarta.ws.rs.container.ContainerRequestFilter import jakarta.ws.rs.core.MediaType import jakarta.ws.rs.core.Response import jakarta.ws.rs.ext.Provider import org.jboss.resteasy.reactive.server.ServerExceptionMapper import org.jboss.resteasy.reactive.server.annotations.Blocking import java.io.ByteArrayInputStream import java.io.ByteArrayOutputStream @Provider class SignatureFilter( private val objectMapper: ObjectMapper, private val webhookHandlerUtility: WebhookHandlerUtility ) : ContainerRequestFilter { @Blocking // 标记为阻塞操作,Quarkus自动调度到Worker线程池 override fun filter(requestContext: ContainerRequestContext) { // 仅处理JSON类型请求 val contentType = requestContext.headerString("Content-Type").orEmpty() if (!contentType.contains(MediaType.APPLICATION_JSON)) { return } // 读取并缓存请求体字节数组 val byteArrayOutputStream = ByteArrayOutputStream() requestContext.entityStream.copyTo(byteArrayOutputStream) val bodyBytes = byteArrayOutputStream.toByteArray() val body = String(bodyBytes) // 签名校验 val signature = requestContext.getHeaderString("signature") if (!webhookHandlerUtility.verifySignature(signature, body)) { throw TerraSignatureVerificationException() } // 解析payload并替换请求体流 val payload = webhookHandlerUtility.parseWebhookPayload(body) ?: throw TerraWebHookPayloadParseException() val processedBodyBytes = objectMapper.writeValueAsBytes(payload.raw) requestContext.entityStream = ByteArrayInputStream(processedBodyBytes) } } // 自定义异常全局处理器,统一返回错误响应 @ServerExceptionMapper fun handleSignatureException(e: TerraSignatureVerificationException): Response { return Response.status(Response.Status.UNAUTHORIZED).entity("签名校验失败").build() } @ServerExceptionMapper fun handlePayloadParseException(e: TerraWebHookPayloadParseException): Response { return Response.status(Response.Status.BAD_REQUEST).entity("Payload解析失败").build() }
控制器代码(无需修改)
import jakarta.enterprise.context.ApplicationScoped import jakarta.ws.rs.Consumes import jakarta.ws.rs.POST import jakarta.ws.rs.Path import jakarta.ws.rs.core.MediaType import jakarta.ws.rs.core.Response import io.smallrye.mutiny.Uni import org.slf4j.LoggerFactory @ApplicationScoped @Path("/webHook") class WebHookController( private val omronService: OmronService ) { private val log = LoggerFactory.getLogger(javaClass) @POST @Path("/device") @Consumes(MediaType.APPLICATION_JSON) fun omronWebHook(req: HookDto): Uni<Response> { log.info("Received data: $req") return omronService.save(req) .map { Response.noContent().build() } } }
关键说明
- @Blocking注解:Quarkus Reactive环境中,阻塞操作必须标记此注解,避免阻塞EventLoop线程,保障应用性能与响应性。
- 请求体缓存:通过
ByteArrayOutputStream缓存请求体字节数组,确保流可重复读取,解决控制器参数为null的问题。 - 同步处理逻辑:所有校验与流修改操作在同一线程完成,避免异步操作请求上下文引发的非法状态异常。
内容的提问来源于stack exchange,提问作者Pavel
相关产品推荐
相关产品推荐

