Anypoint MQ消息长期处于飞行状态未处理及并发异常排查
Anypoint MQ 订阅者问题排查
环境信息
- MuleRuntime版本:4.4
- MQ Connector版本:3.2.0
- 使用普通非FIFO队列,需求为每次处理1条消息;因任务从
jobUpload状态变为JobCompleted状态最长需15分钟,故设置确认超时时间为15分钟
订阅者配置
<anypoint-mq:subscriber doc:name="Subscribering Bulk Query Job Details" config-ref="Anypoint_MQ_Config" destination="${anyPointMq.name}" acknowledgementTimeout="15" acknowledgementTimeoutUnit="MINUTES"> <anypoint-mq:subscriber-type > <anypoint-mq:prefetch maxLocalMessages="1" /> </anypoint-mq:subscriber-type> </anypoint-mq:subscriber>
Anypoint MQ 连接器配置
<anypoint-mq:config name="Anypoint_MQ_Config" doc:name="Anypoint MQ Config" doc:id="ce3aaed9-dcba-41bc-8c68-037c5b1420e2"> <anypoint-mq:connection clientId="${secure::anyPointMq.clientId}" clientSecret="${secure::anyPointMq.clientSecret}" url="${anyPointMq.url}"> <reconnection> <reconnect frequency="3000" count="3" /> </reconnection> <anypoint-mq:tcp-client-socket-properties connectionTimeout="30000" /> </anypoint-mq:connection> </anypoint-mq:config>
订阅者流配置
<flow name="sfdc-bulk-query-job-subscription" doc:id="7e1e23d0-d7f1-45ed-a609-0fb35dd23e6a" maxConcurrency="1"> <anypoint-mq:subscriber doc:name="Subscribering Bulk Query Job Details" doc:id="98b8b25e-3141-4bd7-a9ab-86548902196a" config-ref="Anypoint_MQ_Config" destination="${anyPointMq.sfPartnerEds.name}" acknowledgementTimeout="${anyPointMq.ackTimeout}" acknowledgementTimeoutUnit="MINUTES"> <anypoint-mq:subscriber-type > <anypoint-mq:prefetch maxLocalMessages="${anyPointMq.prefecth.maxLocalMsg}" /> </anypoint-mq:subscriber-type> </anypoint-mq:subscriber> <json-logger:logger doc:name="INFO - Bulk Job Details have been fetched" doc:id="b25c3850-8185-42be-a293-659ebff546d7" config-ref="JSON_Logger_Config" message='#["Bulk Job Details have been fetched for " ++ payload.object default ""]'> <json-logger:content ><![CDATA[#[output application/json --- payload]]]></json-logger:content> </json-logger:logger> <set-variable value="#[p('serviceName.sfdcToEds')]" doc:name="ServiceName" doc:id="f1ece944-0ed8-4c0e-94f2-3152956a2736" variableName="ServiceName"/> <set-variable value="#[payload.object]" doc:name="sfObject" doc:id="2857c8d9-fe8d-46fa-8774-0eed91e3a3a6" variableName="sfObject" /> <set-variable value="#[message.attributes.properties.key]" doc:name="key" doc:id="57028932-04ab-44c0-bd15-befc850946ec" variableName="key" /> <flow-ref doc:name="bulk-job-status-check" doc:id="c6b9cd40-4674-47b8-afaa-0f789ccff657" name="bulk-job-status-check" /> <json-logger:logger doc:name="INFO - subscribed bulk job id has been processed successfully" doc:id="7e469f92-2aff-4bf4-84d0-76577d44479a" config-ref="JSON_Logger_Config" message='#["subscribed bulk job id has been processed successfully for salesforce " ++ vars.sfObject default "" ++ " object"]' tracePoint="END"/> </flow>
业务流程补充说明
- 订阅消息后,在
Until Successful作用域内以1分钟为间隔检查任务状态,最多尝试5次。通常会耗尽所有尝试次数,之后重新订阅并重复流程直到任务完成,单个任务会多次耗尽尝试次数。 - 任务状态变为
jobComplete后,获取结果并通过MuleSoft系统API发送至AWS S3存储桶,此处配置了重试逻辑,首次调用常出现以下错误,第二次重试可成功:HTTP POST on resource 'https://****//dlb.lb.anypointdns.net:443/api/sys/aws/s3/databricks/object' failed: Remotely closed.
主要问题
- 消息长期处于飞行状态,数天后仍未被订阅者拾取,曾出现7条消息滞留在飞行状态的情况。
- 已将
maxConcurrency和maxPrefetchLocalMsg设置为1,但仍有超过1条消息被取出队列,需排查原因。
内容的提问来源于stack exchange,提问作者joono
相关产品推荐
相关产品推荐

