لماذا المراسلة غير المتزامنة؟ مقارنة بين HTTP المتزامن وMQ
في تكامل HTTP المتزامن، ينتظر كل طلب استجابة المستهلك. عند وصول 1000 طلب متزامن، ينضب تجمع الخيوط (thread pool) وتنتشر حالات timeout بشكل متسلسل. يحل Anypoint MQ هذه المشكلة بفصل المنتج (producer) عن المستهلك (consumer).
| المعيار | HTTP المتزامن | Anypoint MQ (غير متزامن) |
|---|---|---|
| السلوك تحت الحمل | Timeout & خطأ | تتراكم الرسائل في قائمة الانتظار |
| الاعتماد بين المنتج والمستهلك | اقتران محكم (tight coupling) | اقتران مرن (loose coupling) |
| فقدان البيانات عند الخطأ | خطر مرتفع | محمية بـ DLQ |
| التوسع | يتطلب توسعًا رأسيًا | يُزاد عدد المستهلكين |
| التأخر (latency) | منخفض (بالملي ثانية) | متوسط (بالثانية) |
متى تستخدم MQ؟ إذا لم يكن الحصول على الاستجابة فوريًا أمرًا حيويًا، أو كان المستهلك بطيئًا أو كان الحمل غير متوازن — فـ MQ هو الخيار الصحيح. في السيناريوهات التي تتطلب استجابة فورية للمستخدم، يمكن تفضيل المزامن.
تدفق المنتج (Producer Flow): نشر الرسائل
في جانب المنتج، يُستخدم مكوّن Anypoint MQ Publish. يُرسَل نص الرسالة وmessage ID والخصائص مثل correlation ID والأولوية عبر headers.
<!-- تدفق المنتج: 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 zorunludur"/> <!-- النشر إلى 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 فورًا. مهما كانت سرعة معالجة المستهلك، لا ينتظر العميل ولا تحدث حالة timeout.تدفق المستهلك (Consumer Flow): استهلاك الرسائل والتكرارية
في جانب المستهلك، يُستخدم مكوّن 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="Duplicate mesaj atlandı: #[attributes.messageId]"/> <anypoint-mq:ack ackToken=#[attributes.ackToken]/> </otherwise> </choice> </flow>
- ✓maxConcurrency="10" — تُعالَج 10 رسائل بالتوازي في آنٍ واحد، ويُوزَّع الحمل
- ✓ackMode="MANUAL" — يُرسَل ACK عند نجاح العملية، وعند الخطأ يعود إلى قائمة الانتظار بـ NACK
- ✓Object Store TTL — يُحذف سجل التكرارية بعد 24 ساعة، فلا تنتفخ الذاكرة
إدارة الأخطاء: قائمة الرسائل الميتة (DLQ)
في حركة المرور الكثيفة، قد لا تتمكن بعض الرسائل من المعالجة — انقطاع الاتصال أو خطأ في البيانات أو عطل في النظام الدني. تلتقط DLQ هذه الرسائل دون فقدانها وتُعيد معالجتها عبر آلية إعادة المحاولة.
<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="Geçici hata, 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 هو أنظف طريقة لإجراء هذا التحويل داخل تدفق المستهلك.
// MQ'dan gelen sipariş → SAP formatına dönüşüm %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" } }) } }
- ✓يُصل إلى قيم header عبر attributes.properties في MQ — لا يُلوّث payload
- ✓يوفّر عامل default ربطًا آمنًا من null — لا يحدث NullPointerException في الرسائل الخاطئة
- ✓تنسيق تاريخ SAP yyyyMMddHHmmss — يُحوَّل في DataWeave بـ
as String {format: ...}
تحسين الأداء: أي الإعدادات تُحدث الفرق؟
يؤثر التكوين الصحيح في مشترك Anypoint MQ تأثيرًا دراماتيكيًا على الإنتاجية. المعاملات التالية حرجة لبيئة الإنتاج.
<anypoint-mq:subscriber config-ref="Anypoint_MQ_Config" destination="orders-queue" <!-- عدد الرسائل التي تُعالَج بالتوازي --> maxConcurrency="20" <!-- عدد الرسائل المسحوبة في استطلاع واحد (1-10) --> fetchSize="10" <!-- المدة المنتظَرة إذا كانت قائمة الانتظار فارغة (ms) --> pollingTime="1000" <!-- مهلة معالجة الرسالة --> acknowledgementTimeout="60000" <!-- وضع ACK اليدوي --> ackMode="MANUAL"/>
| المعامل | حركة مرور منخفضة | حركة مرور عالية | الوصف |
|---|---|---|---|
| maxConcurrency | 5 | 20–50 | يُضبط وفقًا لـ CloudHub worker vCPU |
| fetchSize | 1–3 | 10 | عدد الرسائل المسحوبة في كل مرة |
| pollingTime | 5000 ms | 500–1000 ms | التخفيض يقلل التأخر لكن يزيد التكلفة |
| ackTimeout | 30 ثانية | 60–120 ثانية | أضف مدة المعالجة + مخزن مؤقت |




