BelajarKoding Logobelajarkoding

Platform belajar web development Indonesia. Artikel, cheat sheets, roadmap, dan code challenges untuk developer Indonesia.

Navigasi

  • Artikel
  • Cheat Sheets
  • Roadmap
  • Challenges
  • Pricing
  • Search

Produk Lain

  • JagoHermes
  • KelasClaude
  • KilatKoding
  • BelajarVibeCoding
  • JualanKoding

Support

  • Privacy Policy
  • Terms of Service
  • Email

© 2026 BelajarKoding. All rights reserved.

Galih PratamaBagian dari ekosistem Galih Pratama
belajarkoding LogobyGalih Pratama
RoadmapArtikelCheat SheetsChallengesUpgrade
belajarkoding LogobyGalih Pratama
RoadmapArtikelCheat SheetsChallengesUpgrade
belajarkoding LogobyGalih Pratama
RoadmapArtikelCheat SheetsChallengesUpgrade

Daftar Isi

Konsep DasarApa Itu Message QueueKomponen UtamaArsitektur DasarPola KomunikasiPoint-to-PointPublish/Subscribe (Pub/Sub)Work Queue (Competing Consumers)Request/Reply (RPC)RabbitMQ ExchangesTipe ExchangeTopic ExchangeHeader ExchangeRedis StreamsApa Itu Redis StreamsOperasi DasarConsumer GroupPending Entries dan ClaimPerbandingan Redis Streams vs RabbitMQAmazon SQS dan SNSSQS (Simple Queue Service)Visibility TimeoutSNS (Simple Notification Service)Fanout PatternDead Letter Queue (DLQ)Apa Itu DLQDLQ di RabbitMQDLQ di SQSRetry dan BackoffStrategi RetryImplementasi Exponential Backoff dengan JitterRetry Queue di RabbitMQIdempotencyApa Itu IdempotencyStrategi IdempotencyContoh ImplementasiIdempotency Key dengan HTTP APIAt-Least-Once vs At-Most-Once vs Exactly-OnceGlossary
MessagingDistributed SystemsPython

Message Queue Patterns Cheat Sheet

Referensi cepat message queue patterns. Point-to-point, pub/sub, work queue, dead letter queue, retry backoff, idempotency. Perfect buat developer yang bangun sistem async dan event-driven.

Python10 min read1.914 kata
Silakan login atau daftar untuk membaca cheat sheet ini.

Cheat sheet ini membahas pola-pola message queue yang umum dipakai di sistem terdistribusi, lengkap dengan contoh kode untuk RabbitMQ, Redis Streams, dan Amazon SQS/SNS. Memahami pola ini penting agar sistem kamu scalable, reliable, dan fault-tolerant.

#Konsep Dasar

#Apa Itu Message Queue

Message queue adalah sistem komunikasi async antar komponen aplikasi. Producer mengirim pesan ke queue, consumer mengambil dan memproses pesan tersebut. Queue memutus (decouple) producer dari consumer sehingga keduanya bisa bekerja dengan kecepatan berbeda.

#Komponen Utama

KomponenFungsi
Producer / PublisherMengirim pesan ke queue atau exchange
Consumer / SubscriberMenerima dan memproses pesan
QueuePenampung pesan sementara
Exchange / TopicRouter yang mendistribusikan pesan (RabbitMQ)
BrokerServer yang menyimpan dan mengelola pesan

#Arsitektur Dasar

plaintext
Producer ----> [Broker: Queue/Exchange] ----> Consumer
                    |
                    +---> Dead Letter Queue (jika gagal)

#Pola Komunikasi

#Point-to-Point

Satu producer, satu consumer. Setiap pesan diproses oleh tepat satu consumer.

plaintext
Producer ----> [Queue] ----> Consumer

Contoh RabbitMQ:

python
import pika
 
# Producer
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
 
channel.queue_declare(queue='task_queue', durable=True)
 
channel.basic_publish(
    exchange='',
    routing_key='task_queue',
    body='Process this task',
    properties=pika.BasicProperties(
        delivery_mode=2  # Pesan persistent
    )
)
print(" [x] Sent 'Process this task'")
connection.close()
 
# Consumer
def callback(ch, method, properties, body):
    print(f" [x] Received {body}")
    ch.basic_ack(delivery_tag=method.delivery_tag)
 
channel.basic_consume(queue='task_queue', on_message_callback=callback)
channel.start_consuming()

#Publish/Subscribe (Pub/Sub)

Satu producer, banyak consumer. Setiap consumer menerima copy pesan yang sama.

plaintext
                    +---> Consumer A
Producer ----> [Exchange]
                    +---> Consumer B
                    +---> Consumer C

Contoh RabbitMQ dengan fanout exchange:

python
import pika
 
# Publisher
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
 
channel.exchange_declare(exchange='logs', exchange_type='fanout')
 
message = "INFO: User logged in"
channel.basic_publish(exchange='logs', routing_key='', body=message)
print(f" [x] Sent: {message}")
connection.close()
 
# Subscriber
def callback(ch, method, properties, body):
    print(f" [x] {body}")
 
result = channel.queue_declare(queue='', exclusive=True)
queue_name = result.method.queue
 
channel.queue_bind(exchange='logs', queue=queue_name)
channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True)
channel.start_consuming()

#Work Queue (Competing Consumers)

Satu queue, banyak consumer. Setiap pesan diproses oleh tepat satu consumer, tapi beban tersebar.

plaintext
                    +---> Worker 1
Producer ----> [Queue]
                    +---> Worker 2
                    +---> Worker 3

Contoh dengan fair dispatch di RabbitMQ:

python
# Pastikan prefetch_count = 1 untuk fair dispatch
channel.basic_qos(prefetch_count=1)
 
def callback(ch, method, properties, body):
    import time
    print(f" [x] Processing {body}")
    time.sleep(body.count(b'.'))
    print(" [x] Done")
    ch.basic_ack(delivery_tag=method.delivery_tag)
 
channel.basic_consume(queue='task_queue', on_message_callback=callback)

#Request/Reply (RPC)

Producer mengirim pesan dan menunggu balasan dari consumer. Dipakai untuk pola RPC async.

python
import pika
import uuid
import json
 
class RpcClient:
    def __init__(self):
        self.connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
        self.channel = self.connection.channel()
        result = self.channel.queue_declare(queue='', exclusive=True)
        self.callback_queue = result.method.queue
        self.channel.basic_consume(
            queue=self.callback_queue,
            on_message_callback=self.on_response,
            auto_ack=True
        )
        self.response = None
        self.corr_id = None
 
    def on_response(self, ch, method, props, body):
        if self.corr_id == props.correlation_id:
            self.response = body
 
    def call(self, n):
        self.response = None
        self.corr_id = str(uuid.uuid4())
        self.channel.basic_publish(
            exchange='',
            routing_key='rpc_queue',
            properties=pika.BasicProperties(
                reply_to=self.callback_queue,
                correlation_id=self.corr_id,
            ),
            body=str(n)
        )
        while self.response is None:
            self.connection.process_data_events(time_limit=None)
        return int(self.response)
 
client = RpcClient()
result = client.call(30)
print(f" [.] Got {result}")

#RabbitMQ Exchanges

#Tipe Exchange

TipeRoutingContoh Kasus
directRouting key harus sama persisFilter by severity
fanoutBroadcast ke semua queuePub/sub, logs
topicRouting key dengan pattern wildcardHierarchical routing
headersBerdasarkan header message, bukan routing keyComplex routing

#Topic Exchange

Topic exchange memakai routing key dengan pattern yang dipisahkan titik. Wildcard: * (satu kata), # (nol atau lebih kata).

python
# Publisher
channel.exchange_declare(exchange='topic_logs', exchange_type='topic')
 
routing_key = 'kern.critical'
channel.basic_publish(
    exchange='topic_logs',
    routing_key=routing_key,
    body='Kernel critical error'
)
 
# Subscriber: dengarkan semua kern.*
channel.queue_bind(exchange='topic_logs', queue=queue_name, routing_key='kern.*')
 
# Subscriber: dengarkan semua critical dari apa pun
channel.queue_bind(exchange='topic_logs', queue=queue_name, routing_key='*.critical')
 
# Subscriber: dengarkan semuanya
channel.queue_bind(exchange='topic_logs', queue=queue_name, routing_key='#')

#Header Exchange

python
# Bind queue berdasarkan header
channel.queue_bind(
    exchange='headers_logs',
    queue=queue_name,
    routing_key='',
    arguments={
        'format': 'pdf',
        'type': 'report',
        'x-match': 'all'  # atau 'any'
    }
)

#Redis Streams

#Apa Itu Redis Streams

Redis Streams adalah struktur data log append-only yang mendukung consumer group. Cocok untuk message queue ringan tanpa broker terpisah.

#Operasi Dasar

python
import redis
 
r = redis.Redis(host='localhost', port=6379, db=0)
 
# Producer: tambah entry ke stream
entry_id = r.xadd('mystream', {
    'sensor_id': 'temp_01',
    'temperature': '24.5',
    'timestamp': '2024-01-01T10:00:00'
})
print(f"Entry ID: {entry_id}")
 
# Consumer: baca entry
entries = r.xread({'mystream': '0'}, count=2, block=5000)
for stream, messages in entries:
    for msg_id, data in messages:
        print(f"ID: {msg_id}, Data: {data}")

#Consumer Group

Consumer group memungkinkan multiple consumer membagi stream secara paralel, mirip work queue.

python
# Buat consumer group
r.xgroup_create('mystream', 'mygroup', id='0', mkstream=True)
 
# Consumer ambil pesan dari group
messages = r.xreadgroup(
    groupname='mygroup',
    consumername='consumer_1',
    streams={'mystream': '>'},
    count=10
)
for stream, msgs in messages:
    for msg_id, data in msgs:
        print(f"Processing {msg_id}: {data}")
        # Acknowledge setelah diproses
        r.xack('mystream', 'mygroup', msg_id)

#Pending Entries dan Claim

Pesan yang belum di ack tetap berada di pending entry list (PEL). Kalau consumer mati, consumer lain bisa claim pesannya.

python
# Lihat pending entries
pending = r.xpending('mystream', 'mygroup')
print(pending)
 
# Claim pesan yang idle terlalu lama
claimed = r.xautoclaim(
    'mystream', 'mygroup', 'consumer_2',
    min_idle_time=30000,
    start_id='0-0',
    count=10
)

#Perbandingan Redis Streams vs RabbitMQ

AspekRedis StreamsRabbitMQ
KompleksitasRendah (in-memory)Sedang (broker terpisah)
PersistenceOptional (RDB/AOF)Default persistent
RoutingLinear, consumer groupExchange (direct, fanout, topic, header)
ThroughputSangat tinggiTinggi
Fitur lanjutanTerbatasRich (DLX, priority, TTL)
Cocok untukStreaming, logKompleks routing, enterprise

#Amazon SQS dan SNS

#SQS (Simple Queue Service)

SQS adalah managed queue service dari AWS. Tipe: Standard (at-least-once, tidak terurut) dan FIFO (exactly-once, terurut).

python
import boto3
 
sqs = boto3.client('sqs', region_name='us-east-1')
 
# Buat queue
response = sqs.create_queue(QueueName='my-queue')
queue_url = response['QueueUrl']
 
# Kirim pesan
sqs.send_message(
    QueueUrl=queue_url,
    MessageBody='Process this order',
    DelaySeconds=10,
    MessageAttributes={
        'order_type': {
            'DataType': 'String',
            'StringValue': 'premium'
        }
    }
)
 
# Terima pesan
response = sqs.receive_message(
    QueueUrl=queue_url,
    MaxNumberOfMessages=10,
    WaitTimeSeconds=20,  # Long polling
    MessageAttributeNames=['All']
)
 
for message in response['Messages']:
    print(f"Body: {message['Body']}")
    # Hapus setelah diproses
    sqs.delete_message(
        QueueUrl=queue_url,
        ReceiptHandle=message['ReceiptHandle']
    )

#Visibility Timeout

Saat consumer menerima pesan, pesan jadi invisible untuk consumer lain selama visibility timeout. Kalau tidak di ack (delete), pesan akan muncul lagi.

python
# Set visibility timeout per queue
sqs.set_queue_attributes(
    QueueUrl=queue_url,
    Attributes={'VisibilityTimeout': '120'}
)
 
# Change visibility timeout per message
sqs.change_message_visibility(
    QueueUrl=queue_url,
    ReceiptHandle=receipt_handle,
    VisibilityTimeout=300
)

#SNS (Simple Notification Service)

SNS adalah pub/sub messaging. Topic mengirim ke multiple subscriber (SQS, Lambda, HTTP, email).

python
sns = boto3.client('sns', region_name='us-east-1')
 
# Buat topic
topic = sns.create_topic(Name='order-events')
topic_arn = topic['TopicArn']
 
# Subscribe SQS queue ke SNS topic
sns.subscribe(
    TopicArn=topic_arn,
    Protocol='sqs',
    Endpoint=sqs_queue_arn
)
 
# Publish ke topic
sns.publish(
    TopicArn=topic_arn,
    Message='{"event": "order_created", "order_id": 12345}',
    Subject='New Order'
)

#Fanout Pattern

Fanout adalah pola di mana satu SNS topic mengirim ke multiple SQS queue, setiap queue mendapat copy pesan.

plaintext
                    +---> SQS Queue A (email service)
SNS Topic ---->     +---> SQS Queue B (inventory service)
                    +---> SQS Queue C (analytics)

#Dead Letter Queue (DLQ)

#Apa Itu DLQ

DLQ adalah queue tempat pesan yang gagal diproses dipindahkan setelah mencapai batas retry. Mencegah pesan beracun (poison pill) mengulang terus.

#DLQ di RabbitMQ

python
# Deklarasi DLQ
channel.queue_declare(queue='main_queue', arguments={
    'x-dead-letter-exchange': '',
    'x-dead-letter-routing-key': 'dlq',
    'x-max-retries': 3
})
channel.queue_declare(queue='dlq', durable=True)
 
# Pesan yang gagal 3x otomatis pindah ke dlq

#DLQ di SQS

python
# Buat DLQ
dlq = sqs.create_queue(QueueName='my-dlq')
dlq_arn = sqs.get_queue_attributes(
    QueueUrl=dlq['QueueUrl'],
    AttributeNames=['QueueArn']
)['Attributes']['QueueArn']
 
# Set redrive policy pada main queue
sqs.set_queue_attributes(
    QueueUrl=main_queue_url,
    Attributes={
        'RedrivePolicy': json.dumps({
            'deadLetterTargetArn': dlq_arn,
            'maxReceiveCount': '3'
        })
    }
)

#Retry dan Backoff

#Strategi Retry

StrategiCara KerjaCocok Untuk
Fixed delayTunggu waktu tetap (mis. 5 detik) antar retrySederhana, error sementara
Linear backoffTunggu makin lama secara linearBeban moderat
Exponential backoffTunggu 2x lebih lama setiap retrySistem sibuk, menghindari overload
Exponential + jitterExponential backoff dengan random noiseDistribusi beban, mencegah thundering herd

#Implementasi Exponential Backoff dengan Jitter

python
import time
import random
 
def retry_with_backoff(func, max_retries=5, base_delay=1.0, max_delay=60.0):
    for attempt in range(max_retries):
        try:
            return func()
        except Exception as e:
            if attempt == max_retries - 1:
                raise
            # Exponential backoff dengan full jitter
            delay = min(max_delay, base_delay * (2 ** attempt))
            jitter = random.uniform(0, delay)
            print(f"Attempt {attempt + 1} failed: {e}. Retrying in {jitter:.2f}s")
            time.sleep(jitter)
 
# Penggunaan
result = retry_with_backoff(lambda: process_message(msg))

#Retry Queue di RabbitMQ

python
# Pattern: TTL queue untuk delayed retry
# 1. Main queue dengan DLX
channel.queue_declare(queue='main', arguments={
    'x-dead-letter-exchange': 'retry_exchange'
})
 
# 2. Retry queue dengan TTL, DLX kembali ke main
channel.queue_declare(queue='retry_30s', arguments={
    'x-message-ttl': 30000,
    'x-dead-letter-exchange': '',
    'x-dead-letter-routing-key': 'main'
})
 
# 3. Bind retry exchange ke retry queue
channel.exchange_declare(exchange='retry_exchange', exchange_type='direct')
channel.queue_bind(exchange='retry_exchange', queue='retry_30s', routing_key='')

#Idempotency

#Apa Itu Idempotency

Idempotency berarti memproses pesan yang sama berkali-kali menghasilkan efek yang sama dengan memproses sekali. Penting karena message queue bisa mengirim pesan duplikat (at-least-once delivery).

#Strategi Idempotency

  1. Unique ID: Cek apakah pesan dengan ID ini sudah diproses
  2. Database constraint: UNIQUE constraint pada kolom idempotency key
  3. State check: Cek state saat ini sebelum memproses
  4. Optimistic locking: Versi atau timestamp untuk deteksi duplikat

#Contoh Implementasi

python
import redis
import json
 
r = redis.Redis(host='localhost', port=6379, db=0)
 
def process_message_idempotent(message):
    message_id = message['id']
 
    # Cek apakah sudah diproses
    if r.exists(f"processed:{message_id}"):
        print(f"Message {message_id} already processed, skipping")
        return
 
    # Proses pesan
    result = do_something_with(message['data'])
 
    # Tandai sebagai sudah diproses dengan TTL 24 jam
    r.setex(f"processed:{message_id}", 86400, "1")
 
    return result
 
# Versi dengan database
def process_with_db(message, db_connection):
    cursor = db_connection.cursor()
    try:
        cursor.execute(
            "INSERT INTO processed_messages (message_id, result) VALUES (?, ?)",
            (message['id'], result)
        )
        db_connection.commit()
    except IntegrityError:
        print(f"Duplicate message {message['id']}, skipping")

#Idempotency Key dengan HTTP API

python
from fastapi import FastAPI, Request, HTTPException
 
app = FastAPI()
processed = {}
 
@app.post("/api/order")
async def create_order(request: Request):
    idempotency_key = request.headers.get("Idempotency-Key")
    if not idempotency_key:
        raise HTTPException(status_code=400, detail="Idempotency-Key required")
 
    if idempotency_key in processed:
        return processed[idempotency_key]
 
    body = await request.json()
    result = create_order_in_db(body)
    processed[idempotency_key] = result
    return result

#At-Least-Once vs At-Most-Once vs Exactly-Once

GaransiArtiTrade-off
At-most-oncePesan mungkin hilang, tidak duplikatCepat tapi bisa kehilangan data
At-least-oncePesan tidak hilang, mungkin duplikatButuh idempotency
Exactly-onceTidak hilang, tidak duplikatSulit, biasanya simulasi dengan at-least-once + idempotency

#Glossary

  • Producer / Publisher: Komponen yang mengirim pesan ke queue atau exchange.
  • Consumer / Subscriber: Komponen yang menerima dan memproses pesan.
  • Queue: Penampung pesan sementara yang berurutan.
  • Exchange: Router di RabbitMQ yang mendistribusikan pesan ke queue berdasarkan tipe dan routing key.
  • Binding: Hubungan antara exchange dan queue dengan routing rule.
  • Routing Key: String yang dipakai exchange untuk menentukan queue tujuan.
  • Consumer Group: Sekelompok consumer yang berbagi stream di Redis Streams.
  • Dead Letter Queue (DLQ): Queue untuk pesan yang gagal diproses setelah batas retry.
  • Retry / Backoff: Strategi mencoba ulang pemrosesan pesan yang gagal dengan jeda makin lama.
  • Jitter: Random noise yang ditambahkan ke backoff untuk menghindari thundering herd.
  • Idempotency: Properti di mana operasi yang sama dijalankan berkali-kali menghasilkan efek yang sama.
  • Visibility Timeout: Periode di mana pesan invisible untuk consumer lain setelah diambil (SQS).
  • Acknowledge (Ack): Konfirmasi bahwa pesan sudah diproses dan boleh dihapus dari queue.
  • Poison Pill: Pesan yang selalu gagal diproses dan mengulang terus.
  • At-Least-Once Delivery: Garansi pesan tidak hilang, tapi mungkin duplikat.
  • Exactly-Once Delivery: Garansi pesan diproses tepat sekali, biasanya dicapai dengan idempotency.
  • Fanout: Pola di mana satu pesan dikirim ke multiple subscriber.

Baca Cheat Sheet Lengkap

Login atau daftar akun gratis untuk membaca cheat sheet ini.

LoginDaftar Gratis
Share: