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

Spring Integration中通过SPEL/XPath提取XML字段生成Kafka消息键咨询

解决方案:从XML Payload提取字段生成Kafka消息键

我来帮你搞定这个问题,你之前尝试的SPEL和XPath两种思路都可行,只是路径表达式和调用方式有问题,咱们分别修正:

方法一:用SPEL正确提取字段

你之前的SPEL代码有两个核心问题:一是没有把当前消息对象传入expression.getValue(),导致SPEL无法获取到payload上下文;二是XML有命名空间且层级结构没对应上,直接写payload.Body.ACIFlight.fltNbr找不到节点。

假设你的payload已经被Spring Integration转换为org.w3c.dom.Document对象,修正后的代码如下:

return IntegrationFlows.from(Kafka.messageDrivenChannelAdapter(kafkaListenerContainer))
    .wireTap(ACARS_WIRE_TAP_CHNL) // 记录原始消息
    .enrichHeaders(h -> h.headerFunction(KafkaHeaders.MESSAGE_KEY, message -> {
        SpelExpressionParser parser = new SpelExpressionParser();
        // 利用命名空间查找ACIFlight节点,再提取子字段
        Expression flightNbrExpr = parser.parseExpression(
            "payload.documentElement.getElementsByTagNameNS('http://ual.com/cep/aero/ACIFlight', 'ACIFlight').item(0).getElementsByTagName('fltNbr').item(0).textContent"
        );
        Expression depDateExpr = parser.parseExpression(
            "payload.documentElement.getElementsByTagNameNS('http://ual.com/cep/aero/ACIFlight', 'ACIFlight').item(0).getElementsByTagName('fltLastLegDepDt').item(0).textContent"
        );
        // 传入message作为SPEL的上下文
        String flightNbr = flightNbrExpr.getValue(message, String.class);
        String depDate = depDateExpr.getValue(message, String.class);
        return flightNbr + depDate;
    }))
    .get();

如果你的payload还是原始XML字符串,可以先通过XmlPayloadTransformer把它转换成Document对象,再用上面的SPEL表达式。

方法二:用XPath精准提取(更简洁)

你之前的XPath表达式路径错误,/*[local-name()='fltNbr']是匹配根节点为fltNbr的元素,但你的XML根节点是Envelope,所以需要调整表达式来匹配正确的层级。

方式A:自定义工具类调用

修正XPath表达式后的代码:

return IntegrationFlows.from(Kafka.messageDrivenChannelAdapter(kafkaListenerContainer))
    .wireTap(ACARS_WIRE_TAP_CHNL) // 记录原始消息
    .transform(message -> {
        Object payload = message.getPayload();
        // 用//递归查找ACIFlight节点,再匹配子字段,local-name规避命名空间问题
        String flightNbr = XPathUtils.evaluate(payload, 
            "//*[local-name()='ACIFlight']/*[local-name()='fltNbr']/text()", 
            XPathUtils.STRING);
        String depDate = XPathUtils.evaluate(payload, 
            "//*[local-name()='ACIFlight']/*[local-name()='fltLastLegDepDt']/text()", 
            XPathUtils.STRING);
        return MessageBuilder.fromMessage(message)
            .setHeader(KafkaHeaders.MESSAGE_KEY, flightNbr + depDate)
            .build();
    })
    .get();

方式B:用Spring Integration内置的XPath支持(推荐)

Spring Integration提供了内置的XPath表达式支持,不用自己写工具类,代码更简洁:

return IntegrationFlows.from(Kafka.messageDrivenChannelAdapter(kafkaListenerContainer))
    .wireTap(ACARS_WIRE_TAP_CHNL) // 记录原始消息
    .enrichHeaders(h -> h
        .headerExpression(KafkaHeaders.MESSAGE_KEY, 
            "#xpath(payload, '//*[local-name()=\"ACIFlight\"]/*[local-name()=\"fltNbr\"]/text()') + #xpath(payload, '//*[local-name()=\"ACIFlight\"]/*[local-name()=\"fltLastLegDepDt\"]/text()')")
    )
    .get();

这种方式直接在headerExpression里用#xpath()函数,自动处理XML payload的解析,非常方便。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 10:22:45