¿Por Qué Mensajería Asíncrona? Comparativa Síncrono vs MQ
En la integración HTTP síncrona, cada solicitud espera la respuesta del consumer. Cuando llegan 1000 solicitudes simultáneas, el thread pool se agota y los timeouts se propagan en cadena. Anypoint MQ resuelve este problema con el desacoplamiento producer-consumer.
| Criterio | HTTP Síncrono | Anypoint MQ (Asíncrono) |
|---|---|---|
| Comportamiento bajo carga | Timeout & error | Los mensajes se acumulan en la cola |
| Dependencia Producer-Consumer | Fuertemente acoplado (tight coupling) | Débilmente acoplado (loose coupling) |
| Pérdida de datos en caso de error | Alto riesgo | Protegido con DLQ |
| Escalabilidad | Requiere escala vertical | Se incrementa el número de consumers |
| Latencia | Baja (ms) | Media (orden de segundos) |
¿Cuándo usar MQ? Si la respuesta inmediata no es crítica, si el consumer es lento o si la carga es desequilibrada — MQ es la elección correcta. En escenarios donde se requiere respuesta instantánea al usuario, se puede preferir el modo síncrono.
Producer Flow: Publicación de Mensajes
En el lado del producer se utiliza el componente Anypoint MQ Publish. El cuerpo del mensaje, el correlation ID y propiedades como la prioridad se transmiten mediante headers.
<!-- Producer Flow: HTTP → Anypoint MQ --> <flow name="order-producer-flow"> <http:listener path="/orders" allowedMethods="POST" config-ref="HTTP_Listener_config"/> <!-- Validar el mensaje --> <validation:is-not-null value="#[payload.orderId]" message="orderId zorunludur"/> <!-- Publicar en 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> <!-- Retornar 202 Accepted inmediatamente --> <set-payload value=#[output application/json --- {status: "queued", messageId: payload.orderId}]/> <http:response statusCode="202"/> </flow>
202 Accepted. Sin importar qué tan lento procese el consumer, el cliente no espera y no experimenta timeout.Consumer Flow: Consumo de Mensajes e Idempotencia
En el lado del consumer se utiliza el componente anypoint-mq:subscriber. En escenarios de alto tráfico, el control de idempotencia con Object Store es imprescindible para evitar que el mismo mensaje se procese dos veces.
<flow name="order-consumer-flow"> <anypoint-mq:subscriber config-ref="Anypoint_MQ_Config" destination="orders-queue" maxConcurrency="10" ackMode="MANUAL"/> <!-- Idempotencia: no procesar el mismo mensaje dos veces --> <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> <!-- Duplicado: solo ACK, sin procesar --> <logger message="Duplicate mesaj atlandı: #[attributes.messageId]"/> <anypoint-mq:ack ackToken=#[attributes.ackToken]/> </otherwise> </choice> </flow>
- ✓maxConcurrency="10" — se procesan 10 mensajes en paralelo a la vez, equilibrando la carga
- ✓ackMode="MANUAL" — se envía ACK cuando la operación es exitosa; en caso de error, regresa a la cola con NACK
- ✓Object Store TTL — el registro de idempotencia se elimina después de 24 horas, evitando el crecimiento excesivo de memoria
Gestión de Errores: Dead Letter Queue (DLQ)
En tráfico intenso, algunos mensajes no pueden procesarse — corte de conexión, error de datos o fallo del sistema downstream. La DLQ captura estos mensajes sin pérdida y los retoma con el mecanismo de reintento.
<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> <!-- Error temporal: NACK, regresa a la cola --> <on-error-continue type="CONNECTIVITY, TIMEOUT"> <logger message="Geçici hata, NACK: #[error.description]" level="WARN"/> <anypoint-mq:nack ackToken=#[attributes.ackToken]/> </on-error-continue> <!-- Error permanente: enviar directamente a 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>
Transformación de Mensajes con DataWeave
Los mensajes que llegan por MQ generalmente requieren transformación de formato entre diferentes sistemas. DataWeave es la forma más limpia de realizar esta transformación dentro del consumer flow.
// Pedido recibido por MQ → transformación al formato 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" } }) } }
- ✓Se accede a los valores de header mediante attributes.properties de MQ — no contamina el payload
- ✓El operador default proporciona mapeo null-safe — no habrá NullPointerException en mensajes defectuosos
- ✓El formato de fecha SAP es yyyyMMddHHmmss — se convierte en DataWeave con
as String {format: ...}
Optimización de Rendimiento: ¿Qué Configuraciones Marcan la Diferencia?
La configuración correcta en el subscriber de Anypoint MQ impacta dramáticamente el throughput. Los siguientes parámetros son críticos para el entorno de producción.
<anypoint-mq:subscriber config-ref="Anypoint_MQ_Config" destination="orders-queue" <!-- Cuántos mensajes se procesan en paralelo --> maxConcurrency="20" <!-- Cuántos mensajes se obtienen en un solo polling (1-10) --> fetchSize="10" <!-- Cuánto esperar si la cola está vacía (ms) --> pollingTime="1000" <!-- Timeout de procesamiento de mensajes --> acknowledgementTimeout="60000" <!-- Modo ACK manual --> ackMode="MANUAL"/>
| Parámetro | Tráfico Bajo | Tráfico Alto | Descripción |
|---|---|---|---|
| maxConcurrency | 5 | 20–50 | Se ajusta según los vCPU del worker de CloudHub |
| fetchSize | 1–3 | 10 | Número de mensajes obtenidos de una sola vez |
| pollingTime | 5000 ms | 500–1000 ms | Mantenerlo bajo reduce la latencia, aumenta el costo |
| ackTimeout | 30 s | 60–120 s | Añadir tiempo de procesamiento + buffer |




