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

Quarkus中Apache Camel缓存场景下Http响应码空指针问题

问题

我在Quarkus环境中用Apache Camel开发微服务,identifyOrigin微服务会调用getActivateSessionByIp微服务,现在想给这个调用加缓存。当前逻辑里,首次请求未命中缓存时,能正常捕获getActivateSessionByIp的Http响应码,响应码200就取字段,否则抛异常处理;但命中缓存时,服务识别不了httpResponseCode变量,抛出空指针异常。首次非缓存请求能看到httpResponseCode的响应值,缓存请求时这个变量是空的。

相关代码如下:

RestRoute类

@ApplicationScoped
public class RestRoute extends RouteBuilder {
    @ConfigProperty(name = "client.url.userProfile")
    String urlUserProfile;

    @ConfigProperty(name = "client.utl.getActiveSessionByIp")
    String urlGetActivateSessionByIp;

    @ConfigProperty(name = "descripcion.servicio")
    String descriptionService;

    private ConfigureSsl configureSsl;

    private static final String SALIDA_BSS_EXCEPTION = "Salida Microservicio IdentifyOrigin ${body}";
    private static final String MSG_EXCEPTION = "Descripcion de la Exception: ${exception.message}";

    private static final String DATE_LOG = "[${bean:BeanDate.getCurrentDateTime()}] ";

    public RestRoute() {
        TimeZone.setDefault(TimeZone.getTimeZone("GMT-4"));
        configureSsl= new ConfigureSsl();
    }
    @Override
    public void configure() throws Exception {

        BeanDate beanDate= new BeanDate();
        getContext().getRegistry().bind("BeanDate", beanDate);

        restConfiguration().bindingMode(RestBindingMode.json).dataFormatProperty("json.in.disableFeatures","FAIL_ON_UNKNOWN_PROPERTIES")
                .clientRequestValidation(true)
                .apiProperty("api.title","IdentifyOrigin")
                .apiProperty("api.description",descriptionService)
                .apiProperty("api-version","1.0.0")
                .apiProperty("cors","true");
        rest("/")
                .produces("application/json")
                .consumes("application/json")
                .post("/identifyOrigin/v1/info")
                .type(Request.class)
                .outType(ResponseSuccess.class)
                .param().name("Response").type(RestParamType.body).description("Parametros de Salidas")
                .endParam().to("direct:pipeline");

                from("direct:pipeline")
                        .doTry()
                            .log("["+"${bean:BeanDate.getCurrentDateTime()}"+"] "+ "Valor Header x-correlator: ${header.x-correlator}")
                            .to("bean-validator:validateRequest")
                            .log("["+"${bean:BeanDate.getCurrentDateTime()}"+"] "+ "Datos de Entrada del MS: ${body}")
                             /*Processor classes to create the Microservice Request to call GetActiveSession and Ip */
                            .process(new GetActivateSessionByIpReqProcessor())
                            .log("\n[${bean:BeanDate.getCurrentDateTime()}] " +"Entrada Microservicio GetActivateSessionByIp: ${exchangeProperty[getActivateSessionByIpRequest]}")
                             /* Configure the parameters for the cache*/
                            .setHeader(CaffeineConstants.ACTION, constant(CaffeineConstants.ACTION_GET))
                            .setHeader(CaffeineConstants.KEY).exchangeProperty("ipAddress")
                            .toF("caffeine-cache://%s", "GetActiveSessionByIpCache")
                            .log("Hay Resultado en Cache de la consulta asociado a la IP Consultada: ${exchangeProperty[ipAddress]} ${header.CamelCaffeineActionHasResult}")
                            .log("CamelCaffeineActionSucceeded: ${header.CamelCaffeineActionSucceeded}")
                            .choice()
                                /* It is queried if there is a result in the cache for the sent id*/
                                .when(header(CaffeineConstants.ACTION_HAS_RESULT).isEqualTo(Boolean.FALSE))
                                    .to(configureSsl.setupSSLContext(getCamelContext(), urlGetActivateSessionByIp))
                                    /* A process class is used to validate the response and take different actions depending on the type of response*/
                                    .process(new GetActivateSessionByIpResProcessor())
                                    .setHeader(CaffeineConstants.ACTION, constant(CaffeineConstants.ACTION_PUT))
                                    .setHeader(CaffeineConstants.KEY).exchangeProperty("ipAddress")
                                    .toF("caffeine-cache://%s", "GetActiveSessionByIpCache")
                                .otherwise()
                                        .log("Cache is working")
                                         .process(new GetActivateSessionByIpResProcessor())
                                .endChoice()
                            .end()
                            /* It is required to validate if the http response type is 200, 400, 404, so a log will be used to see it.*/
                             .log("Resultado header: ${exchangeProperty[httpResponseCode]}")
                        .endDoTry()
    }
}

GetActivateSessionByIpReqProcessor类

@ApplicationScoped
public class GetActivateSessionByIpReqProcessor implements Processor {

    private final ObjectMapper objectMapper;
    private final IGetActivateSessionByIpMappingReq getActivateSessionByIpMappingReq;

    public GetActivateSessionByIpReqProcessor() {
        objectMapper = new ObjectMapper();
        getActivateSessionByIpMappingReq=new GetActivateSessionByIpMappingReqImpl();
    }

    @Override
    public void process(Exchange exchange) throws Exception {

        Request request=exchange.getIn().getBody(Request.class);
        String correlatorId = Optional.ofNullable(exchange.getIn().getHeader("x-correlator")).map(Object::toString).orElse("c8ff2443-1b49-4f5f-8412-3fecbb10b836");
        String ipAddress=request.getIpAddress();
        String port= String.valueOf(request.getPort());
        RequestGetActivateSessionByIp requestGetActivateSessionByIp=getActivateSessionByIpMappingReq.toRequest(ipAddress,port);
        String jsonGetActivateSessionByIpMappingReq = objectMapper.writeValueAsString(requestGetActivateSessionByIp);
        exchange.setProperty("getActivateSessionByIpRq", requestGetActivateSessionByIp);
        exchange.setProperty("getActivateSessionByIpRequest", jsonGetActivateSessionByIpMappingReq);
        exchange.setProperty("correlatorId", correlatorId);
        exchange.setProperty("ipAddress", ipAddress);
        exchange.getOut().setHeader(Exchange.CONTENT_TYPE, "application/json");
        exchange.getOut().setHeader(Exchange.HTTP_METHOD, "POST");
        exchange.getOut().setHeader(Exchange.HTTP_PATH, "/api/v1/getActivateSessionByIp");
        exchange.getOut().setHeader("x-transaction-id", correlatorId);
        exchange.getOut().setHeader("x-user", "UsuarioKernel");
        exchange.getOut().setHeader("x-token", "bc4a129867486c6ee7436fa6111c2e");
        exchange.getOut().setBody(jsonGetActivateSessionByIpMappingReq);

    }
}

GetActivateSessionByIpResProcessor类

@ApplicationScoped
public class GetActivateSessionByIpResProcessor implements Processor {

    private final ObjectMapper objectMapper;

    public GetActivateSessionByIpResProcessor() {
        objectMapper = new ObjectMapper();
    }

    @Override
    public void process(Exchange exchange) throws Exception {

        String bodyResponse=exchange.getIn().getBody(String.class);
        String httpResponseCode=exchange.getIn().getHeader("CamelHttpResponseCode", String.class);
        ResponseGetActivateSessionByIp responseGetActivateSessionByIp= objectMapper.readValue(bodyResponse,ResponseGetActivateSessionByIp.class);
        String jsonGetActivateSessionByIpRs = objectMapper.writeValueAsString(responseGetActivateSessionByIp);
        String phoneNumberByIp=responseGetActivateSessionByIp.getId();
        exchange.setProperty("phoneNumberByIp", phoneNumberByIp);
        exchange.setProperty("httpResponseCode", httpResponseCode);

    }
}

原因分析与解决方案

原因

你当前只缓存了getActivateSessionByIp的响应体,没有缓存对应的CamelHttpResponseCode头信息。当命中缓存时,Camel只会把缓存的响应体放回Exchange的In消息中,CamelHttpResponseCode这个头根本不存在,所以处理器里获取这个头时得到null,后续使用httpResponseCode变量就会触发空指针异常。

另外,缓存命中分支直接调用GetActivateSessionByIpResProcessor,但此时Exchange里没有Http调用产生的头信息,处理器自然拿不到响应码。

解决方案

需要把响应体和响应码一起缓存,封装成一个包含这两个信息的DTO类,缓存这个DTO对象而非单独的响应体。

1. 创建缓存用的DTO类

public class CachedSessionResponse {
    private String responseBody;
    private String httpResponseCode;

    public CachedSessionResponse(String responseBody, String httpResponseCode) {
        this.responseBody = responseBody;
        this.httpResponseCode = httpResponseCode;
    }

    public String getResponseBody() {
        return responseBody;
    }

    public String getHttpResponseCode() {
        return httpResponseCode;
    }
}

2. 修改GetActivateSessionByIpResProcessor,生成缓存DTO

@ApplicationScoped
public class GetActivateSessionByIpResProcessor implements Processor {

    private final ObjectMapper objectMapper;

    public GetActivateSessionByIpResProcessor() {
        objectMapper = new ObjectMapper();
    }

    @Override
    public void process(Exchange exchange) throws Exception {
        String bodyResponse=exchange.getIn().getBody(String.class);
        String httpResponseCode=exchange.getIn().getHeader("CamelHttpResponseCode", String.class);
        // 解析响应体并设置业务属性
        ResponseGetActivateSessionByIp responseGetActivateSessionByIp= objectMapper.readValue(bodyResponse,ResponseGetActivateSessionByIp.class);
        String phoneNumberByIp=responseGetActivateSessionByIp.getId();
        exchange.setProperty("phoneNumberByIp", phoneNumberByIp);
        exchange.setProperty("httpResponseCode", httpResponseCode);
        // 创建缓存DTO并存入Exchange属性
        CachedSessionResponse cachedResponse = new CachedSessionResponse(bodyResponse, httpResponseCode);
        exchange.setProperty("cachedSessionResponse", cachedResponse);
    }
}

3. 修改RestRoute的缓存逻辑,缓存DTO并在命中时解析

from("direct:pipeline")
        .doTry()
            .log("["+"${bean:BeanDate.getCurrentDateTime()}"+"] "+ "Valor Header x-correlator: ${header.x-correlator}")
            .to("bean-validator:validateRequest")
            .log("["+"${bean:BeanDate.getCurrentDateTime()}"+"] "+ "Datos de Entrada del MS: ${body}")
            .process(new GetActivateSessionByIpReqProcessor())
            .log("\n[${bean:BeanDate.getCurrentDateTime()}] " +"Entrada Microservicio GetActivateSessionByIp: ${exchangeProperty[getActivateSessionByIpRequest]}")
            /* Configure the parameters for the cache*/
            .setHeader(CaffeineConstants.ACTION, constant(CaffeineConstants.ACTION_GET))
            .setHeader(CaffeineConstants.KEY).exchangeProperty("ipAddress")
            .toF("caffeine-cache://%s", "GetActiveSessionByIpCache")
            .log("Hay Resultado en Cache de la consulta asociado a la IP Consultada: ${exchangeProperty[ipAddress]} ${header.CamelCaffeineActionHasResult}")
            .log("CamelCaffeineActionSucceeded: ${header.CamelCaffeineActionSucceeded}")
            .choice()
                .when(header(CaffeineConstants.ACTION_HAS_RESULT).isEqualTo(Boolean.FALSE))
                    .to(configureSsl.setupSSLContext(getCamelContext(), urlGetActivateSessionByIp))
                    .process(new GetActivateSessionByIpResProcessor())
                    // 缓存DTO对象
                    .setHeader(CaffeineConstants.ACTION, constant(CaffeineConstants.ACTION_PUT))
                    .setHeader(CaffeineConstants.KEY).exchangeProperty("ipAddress")
                    .setBody(exchangeProperty("cachedSessionResponse"))
                    .toF("caffeine-cache://%s", "GetActiveSessionByIpCache")
                .otherwise()
                        .log("Cache is working")
                        // 从缓存DTO中取出响应体和响应码,设置到Exchange中
                        .process(exchange -> {
                            CachedSessionResponse cachedResponse = exchange.getIn().getBody(CachedSessionResponse.class);
                            if (cachedResponse != null) {
                                exchange.getIn().setBody(cachedResponse.getResponseBody());
                                exchange.setProperty("httpResponseCode", cachedResponse.getHttpResponseCode());
                                // 解析响应体设置业务属性
                                ResponseGetActivateSessionByIp response = new ObjectMapper().readValue(cachedResponse.getResponseBody(), ResponseGetActivateSessionByIp.class);
                                exchange.setProperty("phoneNumberByIp", response.getId());
                            }
                        })
                .endChoice()
            .end()
            .log("Resultado header: ${exchangeProperty[httpResponseCode]}")
        .endDoTry()

修改后,缓存命中时就能从缓存的DTO中拿到响应码和响应体,避免空指针异常。

内容的提问来源于stack exchange,提问作者Cesar Justo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 19:12:04