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
相关产品推荐
相关产品推荐

