You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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事件循环,影响应用性能。

修复方案

  1. 替换自定义线程池,使用Quarkus阻塞注解:用@Blocking标记过滤器方法,让Quarkus将任务调度到Worker线程池执行,既避免阻塞EventLoop,又保证同步处理请求上下文。
  2. 缓存并重置请求体流:读取请求体后保存为字节数组,确保后续可重复读取;校验通过后将处理后的payload重新封装为输入流替换原流。
  3. 同步处理校验逻辑:所有签名校验、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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.09 16:25:15