Зачем асинхронный обмен сообщениями? Сравнение синхронного подхода и MQ
В синхронной HTTP-интеграции каждый запрос ждёт ответа от consumer'а. При 1000 одновременных запросах thread pool исчерпывается, таймауты распространяются цепочкой. Anypoint MQ решает эту проблему разделением producer'а и consumer'а.
| Критерий | Синхронный HTTP | Anypoint MQ (Асинхронный) |
|---|---|---|
| Поведение под нагрузкой | Timeout & ошибки | Сообщения накапливаются в очереди |
| Зависимость producer-consumer | Тесная связь (tight coupling) | Слабая связь (loose coupling) |
| Потеря данных при ошибке | Высокий риск | Защита через DLQ |
| Масштабирование | Требуется вертикальное масштабирование | Увеличивается количество consumer'ов |
| Задержка (latency) | Низкая (мс) | Средняя (порядка секунд) |
Когда использовать MQ? Если немедленный ответ не критичен, consumer работает медленно или нагрузка неравномерна — MQ является правильным выбором. В сценариях, требующих мгновенного ответа пользователю, предпочтительнее синхронный подход.
Producer Flow: публикация сообщений
На стороне producer'а используется компонент Anypoint MQ Publish. Тело сообщения, correlation ID и другие свойства, такие как приоритет, передаются через headers.
<!-- Producer Flow: HTTP → Anypoint MQ --> <flow name="order-producer-flow"> <http:listener path="/orders" allowedMethods="POST" config-ref="HTTP_Listener_config"/> <!-- Валидировать сообщение --> <validation:is-not-null value="#[payload.orderId]" message="orderId обязателен"/> <!-- Опубликовать в MQ --> <anypoint-mq:publish config-ref="Anypoint_MQ_Config" destination="orders-exchange" messageId=#[payload.orderId]> <anypoint-mq:properties> <anypoint-mq:property key="priority" value=#[payload.priority default 'NORMAL']/> <anypoint-mq:property key="correlationId" value=#[correlationId]/> </anypoint-mq:properties> </anypoint-mq:publish> <!-- Немедленно вернуть 202 Accepted --> <set-payload value=#[output application/json --- {status: "queued", messageId: payload.orderId}]/> <http:response statusCode="202"/> </flow>
202 Accepted. Как бы медленно consumer ни обрабатывал, клиент не ждёт и не получает таймаут.Consumer Flow: потребление сообщений и идемпотентность
На стороне consumer'а используется компонент anypoint-mq:subscriber. В сценариях с высокой нагрузкой контроль идемпотентности через Object Store обязателен, чтобы одно и то же сообщение не обрабатывалось дважды.
<flow name="order-consumer-flow"> <anypoint-mq:subscriber config-ref="Anypoint_MQ_Config" destination="orders-queue" maxConcurrency="10" ackMode="MANUAL"/> <!-- Идемпотентность: не обрабатывать одно сообщение дважды --> <os:retrieve config-ref="ObjectStore_Config" key=#[attributes.messageId] target="alreadyProcessed" defaultValue="false"/> <choice> <when expression=#[vars.alreadyProcessed == false]> <flow-ref name="process-order-subflow"/> <os:store config-ref="ObjectStore_Config" key=#[attributes.messageId] value="true" ttl="86400" ttlUnit="SECONDS"/> <anypoint-mq:ack ackToken=#[attributes.ackToken]/> </when> <otherwise> <!-- Дубликат: только ACK, без обработки --> <logger message="Дубликат сообщения пропущен: #[attributes.messageId]"/> <anypoint-mq:ack ackToken=#[attributes.ackToken]/> </otherwise> </choice> </flow>
- ✓maxConcurrency="10" — одновременно обрабатываются 10 сообщений параллельно, нагрузка балансируется
- ✓ackMode="MANUAL" — ACK отправляется при успешной обработке, при ошибке NACK возвращает сообщение в очередь
- ✓Object Store TTL — через 24 часа запись идемпотентности удаляется, память не засоряется
Обработка ошибок: Dead Letter Queue (DLQ)
При высокой нагрузке некоторые сообщения не удаётся обработать — из-за разрыва соединения, ошибки данных или сбоя downstream-системы. DLQ перехватывает эти сообщения без потери и передаёт их на повторную обработку через механизм retry.
<flow name="order-consumer-flow"> <anypoint-mq:subscriber destination="orders-queue" ackMode="MANUAL"/> <try> <flow-ref name="process-order-subflow"/> <anypoint-mq:ack ackToken=#[attributes.ackToken]/> <error-handler> <!-- Временная ошибка: NACK, возврат в очередь --> <on-error-continue type="CONNECTIVITY, TIMEOUT"> <logger message="Временная ошибка, NACK: #[error.description]" level="WARN"/> <anypoint-mq:nack ackToken=#[attributes.ackToken]/> </on-error-continue> <!-- Постоянная ошибка: отправить напрямую в DLQ --> <on-error-continue type="VALIDATION, TRANSFORMATION"> <anypoint-mq:publish destination="orders-dlq"> <anypoint-mq:properties> <anypoint-mq:property key="errorReason" value=#[error.description]/> <anypoint-mq:property key="originalMessageId" value=#[attributes.messageId]/> </anypoint-mq:properties> </anypoint-mq:publish> <anypoint-mq:ack ackToken=#[attributes.ackToken]/> </on-error-continue> </error-handler> </try> </flow>
Трансформация сообщений с помощью DataWeave
Сообщения, поступающие через MQ, зачастую требуют преобразования формата между разными системами. DataWeave — самый чистый способ выполнить эту трансформацию внутри consumer flow.
// Заказ из MQ → трансформация в формат SAP %dw 2.0 output application/xml var priorityMap = { "HIGH": "01", "NORMAL": "02", "LOW": "03" } --- SALESORDER: { HEADER: { ORDER_ID: payload.orderId, CUSTOMER: payload.customerId, PRIORITY: priorityMap[attributes.properties."priority" default "NORMAL"], CREATED_AT: now() as String { format: "yyyyMMddHHmmss" }, EXT_REF: attributes.properties."correlationId" default "" }, ITEMS: { (payload.items map (item, idx) -> { ITEM: { LINE_NO: (idx + 1) * 10, MATERIAL: item.sku, QTY: item.quantity, UNIT: item.unit default "EA" } }) } }
- ✓Доступ к значениям заголовков через MQ attributes.properties — payload не загрязняется
- ✓Оператор default обеспечивает null-safe маппинг — нет NullPointerException для некорректных сообщений
- ✓Формат даты SAP yyyyMMddHHmmss — преобразуется в DataWeave через
as String {format: ...}
Оптимизация производительности: какие настройки имеют значение?
Правильная конфигурация subscriber'а в Anypoint MQ кардинально влияет на пропускную способность. Следующие параметры критически важны для производственной среды.
<anypoint-mq:subscriber config-ref="Anypoint_MQ_Config" destination="orders-queue" <!-- Сколько сообщений обрабатывается параллельно --> maxConcurrency="20" <!-- Сколько сообщений забирать за один polling (1-10) --> fetchSize="10" <!-- Сколько ждать при пустой очереди (мс) --> pollingTime="1000" <!-- Таймаут обработки сообщения --> acknowledgementTimeout="60000" <!-- Режим ручного ACK --> ackMode="MANUAL"/>
| Параметр | Низкая нагрузка | Высокая нагрузка | Описание |
|---|---|---|---|
| maxConcurrency | 5 | 20–50 | Настраивается по количеству vCPU CloudHub worker |
| fetchSize | 1–3 | 10 | Количество сообщений, забираемых за один раз |
| pollingTime | 5000 мс | 500–1000 мс | Низкое значение снижает latency, но увеличивает стоимость |
| ackTimeout | 30 сек | 60–120 сек | Время обработки + добавьте буфер |




