MuleSoft中Kafka监听器解析JSON消息出错问题咨询
问题:Kafka消费者中Transform Message无法解析JSON
尝试通过Kafka消息监听器解析主题中的消息,生产者端及Kafka主题中可正常查看消息,但消费者流里的Transform Message组件无法正确解析JSON。
消费者流代码
<?xml version="1.0" encoding="UTF-8"?> <mule xmlns:bigquery="http://www.mulesoft.org/schema/mule/bigquery" xmlns:tls="http://www.mulesoft.org/schema/mule/tls" xmlns:ee="http://www.mulesoft.org/schema/mule/ee/core" xmlns:http="http://www.mulesoft.org/schema/mule/http" xmlns:kafka="http://www.mulesoft.org/schema/mule/kafka" xmlns="http://www.mulesoft.org/schema/mule/core" xmlns:doc="http://www.mulesoft.org/schema/mule/documentation" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://www.mulesoft.org/schema/mule/core http://www.mulesoft.org/schema/mule/core/current/mule.xsd http://www.mulesoft.org/schema/mule/kafka http://www.mulesoft.org/schema/mule/kafka/current/mule-kafka.xsd http://www.mulesoft.org/schema/mule/http http://www.mulesoft.org/schema/mule/http/current/mule-http.xsd http://www.mulesoft.org/schema/mule/ee/core http://www.mulesoft.org/schema/mule/ee/core/current/mule-ee.xsd http://www.mulesoft.org/schema/mule/tls http://www.mulesoft.org/schema/mule/tls/current/mule-tls.xsd http://www.mulesoft.org/schema/mule/bigquery http://www.mulesoft.org/schema/mule/bigquery/current/mule-bigquery.xsd"> <flow name="account-consumer-npFlow" doc:id="0d996208-7cb6-42b9-93f2-77ead77d1883" > <kafka:message-listener doc:name="Message listener" doc:id="c302fcb3-539c-42d9-8d21-347da8328ecd" config-ref="EDH_Consumer_np_configuration"> <repeatable-in-memory-stream /> </kafka:message-listener> <ee:transform doc:name="Transform Message" doc:id="28ca3ce2-2c57-4c86-a845-f00d298a6590" > <ee:message > <ee:set-payload ><![CDATA[%dw 2.0 output application/json --- { out : payload.ChangeEventHeader.changeType }]]></ee:set-payload> </ee:message> </ee:transform> <logger level="INFO" doc:name="Logger" doc:id="c4bc21b6-98d6-466c-a086-1fd5bf48343b" message="#[%dw 2.0 output application/json --- payload]" /> </flow> </mule>
错误日志
INFO 2024-02-06 10:24:43,087 [[MuleRuntime].uber.11: [sf-setotcemailtype-procapi-poc].account-producer-np.CPU_LITE @37e0e7a1] [processor: account-producer-np/processors/0/route/1/processors/0; event: 3bc1f370-c50c-11ee-942a-acde48001122] org.mule.runtime.core.internal.processor.LoggerMessageProcessor: { "ChangeEventHeader": { "commitNumber": 1707236682278494209, "commitUser": "0057f000005zHkEAAU", "sequenceNumber": 1, "entityName": "Account", "changeType": "DELETE", "changedFields": [ ], "changeOrigin": "", "transactionKey": "0000295d-4636-8277-714a-0fbcba7f9c9b", "commitTimestamp": 1707236682000, "recordIds": [ "001O1000007mPzNIAU" ] } } ERROR 2024-02-06 10:24:43,907 [[MuleRuntime].uber.11: [sf-setotcemailtype-procapi-poc].account-consumer-npFlow.CPU_INTENSIVE @3e250b1c] [processor: ; event: 3c5c37a1-c50c-11ee-942a-acde48001122] org.mule.runtime.core.privileged.exception.DefaultExceptionListener: ******************************************************************************** Message : "You called the function 'Value Selector' with these arguments: 1: Binary ("ewogICJpZCI6IFsKICAgICIwMDFPMTAwMDAwN21Qek5JQVUiCiAgXSwKICAiYWN0aXZlX19jIjog...) 2: Name ("ChangeEventHeader") But it expects one of these combinations: (Array, Name) (Array, String) (Date, Name) (DateTime, Name) (LocalDateTime, Name) (LocalTime, Name) (Object, Name) (Object, String) (Period, Name) (Time, Name) 5| out : payload.ChangeEventHeader.changeType ^^^^^^^^^^^^^^^^^^^^^^^^^ Trace: at anonymous::main (line: 5, column: 8)" evaluating expression: "%dw 2.0 output application/json --- { out : payload.ChangeEventHeader.changeType }". Element : account-consumer-npFlow/processors/0 @ sf-setotcemailtype-procapi-poc:account-consumer-np.xml:15 (Transform Message) Element DSL : <ee:transform doc:name="Transform Message" doc:id="28ca3ce2-2c57-4c86-a845-f00d298a6590"> <ee:message> <ee:set-payload><![CDATA[ %dw 2.0 output application/json --- { out : payload.ChangeEventHeader.changeType } ]]></ee:set-payload> </ee:message> </ee:transform> Error type : MULE:EXPRESSION FlowStack : at account-consumer-npFlow(account-consumer-npFlow/processors/0 @ sf-setotcemailtype-procapi-poc:account-consumer-np.xml:15 (Transform Message)) (set debug level logging or '-Dmule.verbose.exceptions=true' for everything) ********************************************************************************
DataWeave Playground测试结果
在DataWeave Playground中可正常解析该消息:
输入JSON
{ "ChangeEventHeader": { "commitNumber": 1707236682278494209, "commitUser": "0057f000005zHkEAAU", "sequenceNumber": 1, "entityName": "Account", "changeType": "DELETE", "changedFields": [ ], "changeOrigin": "", "transactionKey": "0000295d-4636-8277-714a-0fbcba7f9c9b", "commitTimestamp": 1707236682000, "recordIds": [ "001O1000007mPzNIAU" ] } }
DataWeave脚本
%dw 2.0 output application/json --- { out : payload.ChangeEventHeader.changeType }
输出结果
{ "out": "DELETE" }
问题原因及解决方法
原因
错误日志明确指出,当前payload是Binary类型,而DataWeave的对象选择器(.操作)需要的是Object类型。Kafka消息监听器默认将消息以二进制流形式接收,没有自动解析为JSON对象,导致脚本无法访问ChangeEventHeader属性。
解决方法
方法1:在DataWeave脚本中手动解析二进制为JSON
修改Transform Message中的DataWeave脚本,使用read()函数将二进制payload解析为JSON对象:
%dw 2.0 output application/json --- { out : read(payload, "application/json").ChangeEventHeader.changeType }
方法2:添加Parse JSON组件
在Kafka消息监听器之后、Transform Message之前,添加一个Parse JSON组件,将二进制流解析为JSON对象。之后即可使用原来的DataWeave脚本正常访问属性。
内容的提问来源于stack exchange,提问作者Arpit
相关产品推荐
相关产品推荐

