Eltay Yazılım
منصة ANYPOINT

كيف تديرون حركة المرور الكثيفة مع Anypoint MQ؟

منع فقدان البيانات عبر إدارة الأخطاء وقوائم DLQ: نصائح لتوسيع حركة الرسائل عالية الحجم باستخدام Anypoint MQ.

Hakan ÇelikHakan Çelik
٢٤ أبريل ٢٠٢٦ · 6 دقيقة قراءة
كيف تديرون حركة المرور الكثيفة مع Anypoint MQ؟
في الأنظمة ذات حركة المرور العالية، تصبح استدعاءات HTTP المتزامنة غير كافية في نقطة ما — تتضاعف حالات timeout، وتنهار المستهلكون (consumers)، وتُفقد البيانات. يحل Anypoint MQ، بوصفه خدمة المراسلة المُدارة من MuleSoft، هذه المشكلة بالانتقال إلى معمارية غير متزامنة. في هذا المقال نفحص كيف يعمل Anypoint MQ تحت حركة المرور الكثيفة من خلال أمثلة كود حقيقية ومخططات معمارية.
1

لماذا المراسلة غير المتزامنة؟ مقارنة بين HTTP المتزامن وMQ

في تكامل HTTP المتزامن، ينتظر كل طلب استجابة المستهلك. عند وصول 1000 طلب متزامن، ينضب تجمع الخيوط (thread pool) وتنتشر حالات timeout بشكل متسلسل. يحل Anypoint MQ هذه المشكلة بفصل المنتج (producer) عن المستهلك (consumer).

المعيارHTTP المتزامنAnypoint MQ (غير متزامن)
السلوك تحت الحملTimeout & خطأتتراكم الرسائل في قائمة الانتظار
الاعتماد بين المنتج والمستهلكاقتران محكم (tight coupling)اقتران مرن (loose coupling)
فقدان البيانات عند الخطأخطر مرتفعمحمية بـ DLQ
التوسعيتطلب توسعًا رأسيًايُزاد عدد المستهلكين
التأخر (latency)منخفض (بالملي ثانية)متوسط (بالثانية)
متى تستخدم MQ؟ إذا لم يكن الحصول على الاستجابة فوريًا أمرًا حيويًا، أو كان المستهلك بطيئًا أو كان الحمل غير متوازن — فـ MQ هو الخيار الصحيح. في السيناريوهات التي تتطلب استجابة فورية للمستخدم، يمكن تفضيل المزامن.
2

تدفق المنتج (Producer Flow): نشر الرسائل

في جانب المنتج، يُستخدم مكوّن Anypoint MQ Publish. يُرسَل نص الرسالة وmessage ID والخصائص مثل correlation ID والأولوية عبر headers.

Mule 4 — XML Flowproducer-flow.xml
<!-- تدفق المنتج: 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.
3

تدفق المستهلك (Consumer Flow): استهلاك الرسائل والتكرارية

في جانب المستهلك، يُستخدم مكوّن anypoint-mq:subscriber. في سيناريوهات حركة المرور العالية، فحص التكرارية مع Object Store ضروري لمنع معالجة نفس الرسالة مرتين.

Mule 4 — XML Flowconsumer-flow.xml
<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 ساعة، فلا تنتفخ الذاكرة
4

إدارة الأخطاء: قائمة الرسائل الميتة (DLQ)

في حركة المرور الكثيفة، قد لا تتمكن بعض الرسائل من المعالجة — انقطاع الاتصال أو خطأ في البيانات أو عطل في النظام الدني. تلتقط DLQ هذه الرسائل دون فقدانها وتُعيد معالجتها عبر آلية إعادة المحاولة.

Mule 4 — إدارة الأخطاءerror-handler.xml
<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>
تنبيه: راقب الرسائل في DLQ بانتظام. إعداد تنبيه لعمق DLQ في Anypoint Monitoring هو المؤشر الأساسي للنضج التشغيلي.
5

تحويل الرسائل بـ DataWeave

تتطلب الرسائل الواردة عبر MQ في الغالب تحويل التنسيق بين أنظمة مختلفة. DataWeave هو أنظف طريقة لإجراء هذا التحويل داخل تدفق المستهلك.

DataWeave 2.0order-transform.dwl
// 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: ...}
6

تحسين الأداء: أي الإعدادات تُحدث الفرق؟

يؤثر التكوين الصحيح في مشترك Anypoint MQ تأثيرًا دراماتيكيًا على الإنتاجية. المعاملات التالية حرجة لبيئة الإنتاج.

Mule 4 — تكوين المشتركoptimized-subscriber.xml
<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"/>
المعاملحركة مرور منخفضةحركة مرور عاليةالوصف
maxConcurrency520–50يُضبط وفقًا لـ CloudHub worker vCPU
fetchSize1–310عدد الرسائل المسحوبة في كل مرة
pollingTime5000 ms500–1000 msالتخفيض يقلل التأخر لكن يزيد التكلفة
ackTimeout30 ثانية60–120 ثانيةأضف مدة المعالجة + مخزن مؤقت
تجربة Eltay: في أحد عملائنا في قطاع التصنيع، رفعنا maxConcurrency من 5 إلى 20 وجعلنا fetchSize يساوي 10، فزادت الإنتاجية 4x على نفس CloudHub worker. لم تكن هناك حاجة لترقية worker.
مشاركة

ابدأوا رحلة MuleSoft مع الشريك المناسب

لنقيّم معًا احتياجاتكم في الترخيص والاستشارات والترحيل والتدريب والخدمات المُدارة. من خلال تحليل احتياجات مجاني، نُعدّ خارطة طريق MuleSoft الأنسب لمؤسستكم.