Skip to main content

Message Queues in System Design: Asynchronous Processing, Reliability and Scaling

Calculating read time…

Your fraud detection AI just flagged a suspicious transaction. 🚨 Now five things need to happen simultaneously:

  • 📧 Send an email alert to the customer
  • 📱 Push a notification to their mobile app
  • 🔒 Temporarily block the card
  • 📊 Update the fraud dashboard
  • 📋 Write an audit log entry

If your system does all five synchronously — waiting for each one to finish before the next — the fraud check takes 4+ seconds and the whole pipeline stalls. One slow email server brings everything down. 💥

💡 Did You Know?
Every major system you use daily runs on message queues. When you order on Amazon — a queue coordinates warehousing, payment, delivery, and email. When you post on Instagram — a queue handles notification fanout to millions of followers. Uber, Netflix, WhatsApp — all built on message queues at their core. 🌍



What is a Message Queue?

Imagine you walk into a very busy pizza restaurant. 🍕 You place your order with the cashier. The cashier does NOT run to the kitchen, make the pizza, bring it to you, and then come back to take the next customer's order. That would be insanely slow!

Instead, the cashier writes your order on a ticket and puts it on a spinning rack in the kitchen window. 🎰 Then immediately turns to serve the next customer.

The chefs in the kitchen pick up tickets one at a time, make the food, and put completed orders on the counter. Everyone works at their own pace. Nobody is waiting for anybody else.

A Message Queue is exactly this spinning ticket rack — for computers.


✅ Simple Definition:
A Message Queue = A temporary holding area where one system (the Producer) drops messages, and another system (the Consumer) picks them up independently — completely decoupled from each other. Neither has to wait for the other. Messages are stored safely in the queue until a consumer is ready to process them.

🏦 Our Real-World Example: AnomalyAI Fraud Detection

When AnomalyAI detects a fraudulent transaction, it publishes one event to a message queue. Five completely independent services pick it up and process it simultaneously — at their own speed, in their own time.

  WITHOUT MESSAGE QUEUE (synchronous — slow & fragile):

  AnomalyAI detects fraud
       ├──► 1. Send email (2000ms) ── wait ──►
       ├──► 2. Push notification (500ms) ── wait ──►
       ├──► 3. Block card (300ms) ── wait ──►
       ├──► 4. Update dashboard (400ms) ── wait ──►
       └──► 5. Write audit log (200ms) ── wait ──►
  Total: 3400ms 😴  AND if email server is down → EVERYTHING fails! 💥

  WITH MESSAGE QUEUE (asynchronous — fast & resilient):

  AnomalyAI detects fraud
       └──► Publish "fraud_detected" event to queue (5ms) ✅ DONE!

  Queue fans out simultaneously:
  Consumer 1: Email Service    — picks it up, sends email (no one waits)
  Consumer 2: Push Notifier    — picks it up, sends notification
  Consumer 3: Card Blocker     — picks it up, blocks card
  Consumer 4: Dashboard        — picks it up, updates UI
  Consumer 5: Audit Logger     — picks it up, writes log
  Total: 5ms for AnomalyAI 🚀  (Each service runs in parallel, independently)

🔑 The Three Key Players

📮 Point-to-Point vs Pub/Sub — Two Flavours of Queuing

📮 Pattern 1 — Point-to-Point (One Producer, One Consumer)

Like sending a letter to one specific person. 📬 One message is picked up by exactly one consumer. Once consumed — it is gone from the queue. This is used when you want only one system to process each job.

  POINT-TO-POINT — ANOMALYAI USE CASE:

  AnomalyAI → Queue → [Card Blocker Service]

  If there are 3 Card Blocker instances running,
  each message goes to ONLY ONE of them (load balanced):

  MSG #1 → Card Blocker Instance #1 (processes it)
  MSG #2 → Card Blocker Instance #2 (processes it)
  MSG #3 → Card Blocker Instance #3 (processes it)
  MSG #4 → Card Blocker Instance #1 (round-robin)

  ✅ Each card is blocked ONCE — not three times!
  ✅ Load automatically balanced across instances

📡 Pattern 2 — Pub/Sub (One Producer, Many Consumers)

Like a radio station broadcast. 📻 One message is received by ALL subscribed consumers simultaneously. Each consumer gets their own copy. This is used when many systems need the same event.

☁️ OCI Streaming — Kafka-Compatible Message Queue on Oracle Cloud

The most powerful message queue service on Oracle Cloud is OCI Streaming. It is fully compatible with Apache Kafka — the gold standard for real-time data streaming.

  • ⚡ Millions of messages per second — handles peak loads effortlessly
  • 🔒 7-day message retention — replayed anytime if a consumer fails
  • 📊 Multiple consumer groups — Email, Dashboard, Auditor all read independently
  • 🌍 Geo-redundant — data replicated across Oracle Availability Domains
  • 💰 Pay per message — zero idle cost when nothing is streaming
📝 What the code below does (in simple words):
This is the AnomalyAI Producer. When fraud is detected, it sends a message to OCI Streaming (like dropping a letter in a postbox). It does NOT wait for anyone to read it. It just drops it and moves on instantly. The message contains all the fraud details as a JSON object. 📬
# ── FILE: fraud_producer.py ────────────────────────────────────
# PURPOSE: When AnomalyAI detects fraud, this code sends a message
#          to OCI Streaming instantly (fire-and-forget).
#          It does NOT wait for email/notification/card block.
#          All those happen independently after receiving the message.
# ───────────────────────────────────────────────────────────────

import oci
import json
import base64
from datetime import datetime

# Connect to OCI Streaming service
config         = oci.config.from_file("~/.oci/config")
stream_client  = oci.streaming.StreamClient(
    config,
    service_endpoint = "https://streaming.us-ashburn-1.oci.oraclecloud.com"
)

# The Stream (Queue) OCID — created once in OCI Console
STREAM_ID = "ocid1.stream.oc1.iad.anomalyai-fraud-events-xxxxx"


def publish_fraud_event(
    transaction_id : str,
    customer_id    : str,
    amount         : float,
    risk_score     : float,
    reason         : str
):
    """
    Publish a fraud detection event to OCI Streaming.
    This is the ONLY thing AnomalyAI does after detecting fraud.
    Everything else (email, block, notify) is handled by subscribers.
    """

    # Step 1: Build the fraud event payload as a dictionary
    fraud_event = {
        "event_type"    : "FRAUD_DETECTED",
        "event_id"      : f"EVT-{transaction_id}-{int(datetime.now().timestamp())}",
        "transaction_id": transaction_id,
        "customer_id"   : customer_id,
        "amount"        : amount,
        "risk_score"    : risk_score,
        "reason"        : reason,
        "detected_at"   : datetime.utcnow().isoformat() + "Z",
        "severity"      : "CRITICAL" if risk_score > 0.9 else "HIGH"
    }

    # Step 2: Convert dict to JSON string, then base64-encode it
    # OCI Streaming requires message data to be base64-encoded
    message_json    = json.dumps(fraud_event)
    message_encoded = base64.b64encode(message_json.encode()).decode()

    # Step 3: Create the message object for OCI Streaming
    # The "key" is used for partitioning — same customer always goes to same partition
    put_message = oci.streaming.models.PutMessagesDetailsEntry(
        key   = customer_id.encode(),    # Route same customer to same partition
        value = message_encoded          # The encoded fraud event
    )

    # Step 4: Publish to OCI Streaming (takes about 5ms)
    response = stream_client.put_messages(
        stream_id         = STREAM_ID,
        put_messages_details = oci.streaming.models.PutMessagesDetails(
            messages = [put_message]
        )
    )

    # Step 5: Check if it was accepted
    if response.data.failures == 0:
        print(f"✅ Fraud event published in ~5ms!")
        print(f"   Transaction: {transaction_id}")
        print(f"   Risk Score: {risk_score:.2%}")
        print(f"   AnomalyAI is now FREE to process next transaction! ⚡")
    else:
        print(f"⚠️ Some messages failed: {response.data.failures}")


# ── USAGE ───────────────────────────────────────────────────────
publish_fraud_event(
    transaction_id = "TXN-9923",
    customer_id    = "CUST-42",
    amount         = 48750.00,
    risk_score     = 0.97,
    reason         = "Cross-border transaction, 10× average spend, unusual hour"
)
📝 What the code below does (in simple words):
This is the Email Service Consumer. It constantly watches the queue (like a chef watching the ticket rack). When a fraud event arrives, it reads the message and sends an email to the customer. This runs completely independently from AnomalyAI — in its own service, on its own server. 📧
# ── FILE: email_consumer.py ─────────────────────────────────────
# PURPOSE: A separate service that watches the OCI Streaming queue.
#          When it sees a "fraud_detected" event, it sends an email.
#          This runs independently — AnomalyAI does NOT know this exists.
# ───────────────────────────────────────────────────────────────

import oci
import json
import base64
import time

config        = oci.config.from_file("~/.oci/config")
stream_client = oci.streaming.StreamClient(
    config,
    service_endpoint="https://streaming.us-ashburn-1.oci.oraclecloud.com"
)

STREAM_ID        = "ocid1.stream.oc1.iad.anomalyai-fraud-events-xxxxx"
CONSUMER_GROUP   = "email-service-group"    # Unique name for Email Service consumers
INSTANCE_NAME    = "email-consumer-1"       # This consumer's name


def get_or_create_cursor() -> str:
    """
    A 'cursor' is like a bookmark — it tells OCI Streaming
    WHERE in the queue this consumer last read from.
    TRIM_HORIZON = start from the very beginning (all unread messages).
    LATEST = start from right now (only new messages).
    """
    cursor_response = stream_client.create_group_cursor(
        stream_id = STREAM_ID,
        create_group_cursor_details = oci.streaming.models.CreateGroupCursorDetails(
            group_name   = CONSUMER_GROUP,
            instance_name= INSTANCE_NAME,
            type         = oci.streaming.models.CreateGroupCursorDetails.TYPE_TRIM_HORIZON,
            commit_on_get= True    # Auto-acknowledge after reading
        )
    )
    return cursor_response.data.value


def send_fraud_email(fraud_event: dict):
    """
    Send a fraud alert email to the affected customer.
    (Simplified — in production, use OCI Email Delivery or SendGrid.)
    """
    customer_id = fraud_event.get("customer_id")
    amount      = fraud_event.get("amount", 0)
    txn_id      = fraud_event.get("transaction_id")
    risk_score  = fraud_event.get("risk_score", 0)

    # Build email content
    subject = f"🚨 Fraud Alert — Transaction {txn_id}"
    body    = f"""
    Dear Customer {customer_id},

    We detected a suspicious transaction on your account:
    • Transaction ID : {txn_id}
    • Amount         : ${amount:,.2f}
    • Risk Score     : {risk_score:.1%}
    • Detected At    : {fraud_event.get('detected_at')}

    If this was NOT you, your card has been temporarily blocked.
    Please call our fraud line: 1-800-FRAUD-99

    Stay safe,
    AnomalyAI Fraud Protection Team
    """
    print(f"📧 Email sent to customer {customer_id}!")
    print(f"   Subject: {subject}")
    # oci_email_client.send_email(...) ← real implementation here


def run_consumer_loop():
    """
    Main loop — constantly polls OCI Streaming for new fraud events.
    Like a chef who never stops checking the ticket rack! 🍕
    """
    print(f"🚀 Email Service Consumer started!")
    print(f"   Watching stream: {STREAM_ID}")
    print(f"   Consumer group : {CONSUMER_GROUP}")
    print(f"   Ready for fraud events... 👀")

    cursor = get_or_create_cursor()

    while True:
        # Step 1: Poll OCI Streaming for new messages (every 1 second)
        get_response = stream_client.get_messages(
            stream_id = STREAM_ID,
            cursor    = cursor,
            limit     = 10    # Read up to 10 messages per poll
        )

        messages = get_response.data

        if messages:
            print(f"\n📬 Received {len(messages)} message(s)!")
            for message in messages:
                # Step 2: Decode the base64-encoded message
                raw_json    = base64.b64decode(message.value).decode()
                fraud_event = json.loads(raw_json)

                print(f"📨 Processing: {fraud_event.get('event_id')}")
                print(f"   Customer  : {fraud_event.get('customer_id')}")
                print(f"   Risk Score: {fraud_event.get('risk_score'):.1%}")

                # Step 3: Only process fraud events (ignore others)
                if fraud_event.get("event_type") == "FRAUD_DETECTED":
                    send_fraud_email(fraud_event)

            # Step 4: Update cursor to mark these messages as processed
            cursor = get_response.headers.get("opc-next-cursor", cursor)

        else:
            # No new messages — wait 1 second before polling again
            time.sleep(1)


# Start the consumer
if __name__ == "__main__":
    run_consumer_loop()

💀 Dead Letter Queue — What Happens When Processing Fails?

What if the Email Service crashes halfway through sending an email? What if a malformed message arrives that crashes the consumer?

Without protection — the message is lost forever. 😱 With a Dead Letter Queue (DLQ), failed messages are automatically moved to a separate queue where engineers can inspect them, fix the issue, and retry. 🔧


⚖️ Message Ordering and Partitions

In AnomalyAI, the order of fraud events for the same customer matters. If "Card Blocked" arrives before "Fraud Detected" — the timeline is wrong. Partitioning ensures messages for the same customer always go to the same partition, preserving order.


🏗️ Complete AnomalyAI Architecture with Message Queue

  ANOMALYAI — COMPLETE MESSAGE QUEUE ARCHITECTURE (OCI Streaming)

  ┌───────────────────────────────────────────────────────────────────────────┐
  │                                                                           │
  │  TRANSACTIONS IN (500K/min)                                               │
  │         │                                                                 │
  │         ▼                                                                 │
  │  ┌─────────────────────────────────────────────────────────────────────┐  │
  │  │  AnomalyAI Detection Service (the PRODUCER)                         │  │
  │  │  Detects fraud → publishes to OCI Streaming in 5ms → moves on ⚡   │  │
  │  └─────────────────────────────────┬───────────────────────────────────┘  │
  │                                    │                                      │
  │                                    ▼                                      │
  │  ┌─────────────────────────────────────────────────────────────────────┐  │
  │  │  OCI Streaming — "fraud-events" Topic                               │  │
  │  │  ├── Partition 0 (Customer IDs: 1–1M)    ── Consumer Group A       │  │
  │  │  ├── Partition 1 (Customer IDs: 1M–2M)   ── Consumer Group B       │  │
  │  │  └── Partition 2 (Customer IDs: 2M–3M)   ── Consumer Group C       │  │
  │  │  7-day retention | Exactly-once delivery | Kafka-compatible         │  │
  │  └───────────────────────────────────────────────────────┬─────────────┘  │
  │                                                          │                │
  │         ┌──────────────────────────┬───────────────────┐│                │
  │         │                          │                   ││                │
  │         ▼                          ▼                   ▼▼                │
  │  ┌──────────────┐  ┌─────────────────────┐  ┌──────────────────────┐    │
  │  │ 📧 Email     │  │ 📱 Push Notifier     │  │ 🔒 Card Blocker      │    │
  │  │  Service     │  │  (Consumer Group 2) │  │  (Consumer Group 3) │    │
  │  │  Consumer    │  │                     │  │                      │    │
  │  │  Group 1     │  │ 500ms              │  │ 300ms               │    │
  │  │  DLQ: ✅      │  │ DLQ: ✅             │  │ DLQ: ✅              │    │
  │  └──────────────┘  └─────────────────────┘  └──────────────────────┘    │
  │                                                                           │
  │  ALL CONSUMERS: Independent · Scalable · Self-healing · DLQ protected    │
  │                                                                           │
  └───────────────────────────────────────────────────────────────────────────┘

⚙️ Message Queue vs Direct API Call — When to Use Which

  USE DIRECT API CALL WHEN:          USE MESSAGE QUEUE WHEN:
  ──────────────────────────────     ────────────────────────────────
  ✅ You need an immediate answer    ✅ You don't need an instant answer
     "Is this card valid?"              "Send a notification later"

  ✅ Simple 1-to-1 communication     ✅ One event → many consumers
     "Charge this card"                 "Fraud detected → email + block + audit"

  ✅ Sub-100ms response needed       ✅ Tasks take long time (video encode, etc.)
     "Accept or deny this txn?"         "Generate fraud report PDF"

  ✅ Both services always online     ✅ Consumer can be offline temporarily
     Real-time decision needed           Retry later when it comes back

  ANOMALYAI USES BOTH:
  Direct API:     AnomalyAI → Oracle 23ai DB (get customer profile NOW)
  Message Queue:  AnomalyAI → OCI Streaming → Email/Block/Notify (later is fine)

⚠️ Common Message Queue Mistakes

🚫 DON'T #1 — Assume a message is processed exactly once.
Networks fail. Consumers crash mid-processing. OCI Streaming may redeliver a message. Always design consumers to be idempotent — processing the same message twice must give the same result as processing it once. Use a processed-message ID tracker. 🔄
🚫 DON'T #2 — Put large binary files (PDFs, images) in the queue.
Message queues are for small JSON messages (under 1 MB). For large files: upload to OCI Object Storage first, then put only the file reference URL in the queue message. 📄
🚫 DON'T #3 — Ignore the Dead Letter Queue.
DLQ messages are your system's cry for help. 🆘 Set up monitoring alerts whenever a message lands in the DLQ. A full DLQ means fraud alerts are not being processed! 💸
🚫 DON'T #4 — Use a message queue for real-time request-response.
If your user clicks "Check Balance" and expects an answer in under 200ms — a message queue is the wrong tool. Use a direct synchronous API call instead. Queues are for decoupled, asynchronous work. ⏱️
✅ DO #1 — Use a unique event_id in every message. This lets consumers check: "Have I already processed this exact event?" Store processed IDs in a Redis set with TTL. Idempotency = bulletproof consumers. 🛡️
✅ DO #2 — Monitor queue depth (backlog size) continuously. If the queue is growing faster than consumers can process — you need more consumer instances. A rising queue = a future outage warning. 📊
✅ DO #3 — Use consumer groups for independent processing. Email Service and Card Blocker should have separate consumer groups in OCI Streaming. Each group gets its own independent read cursor — they never interfere with each other. 🎯

🗺️ Your Message Queue Learning Roadmap


  • 🍕 The analogy — pizza restaurant ticket rack = message queue
  • ⚙️ Three players — Producer, Queue, Consumer (decoupled forever)
  • 📡 Pub/Sub vs Point-to-Point — broadcast vs direct delivery
  • ☁️ OCI Streaming — Kafka-compatible, full Python producer/consumer code
  • 💀 Dead Letter Queue — the safety net that never loses a message
  • 🗂️ Partitioning — how same-customer messages always stay ordered
  • 🔄 Idempotency — why processing the same message twice must be safe

Happy queuing! 📬✨

Comments