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.
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.
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 | Fungsi |
|---|---|
| Producer / Publisher | Mengirim pesan ke queue atau exchange |
| Consumer / Subscriber | Menerima dan memproses pesan |
| Queue | Penampung pesan sementara |
| Exchange / Topic | Router yang mendistribusikan pesan (RabbitMQ) |
| Broker | Server yang menyimpan dan mengelola pesan |
Producer ----> [Broker: Queue/Exchange] ----> Consumer
|
+---> Dead Letter Queue (jika gagal)Satu producer, satu consumer. Setiap pesan diproses oleh tepat satu consumer.
Producer ----> [Queue] ----> ConsumerContoh RabbitMQ:
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()Satu producer, banyak consumer. Setiap consumer menerima copy pesan yang sama.
+---> Consumer A
Producer ----> [Exchange]
+---> Consumer B
+---> Consumer CContoh RabbitMQ dengan fanout exchange:
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()Satu queue, banyak consumer. Setiap pesan diproses oleh tepat satu consumer, tapi beban tersebar.
+---> Worker 1
Producer ----> [Queue]
+---> Worker 2
+---> Worker 3Contoh dengan fair dispatch di RabbitMQ:
# 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)Producer mengirim pesan dan menunggu balasan dari consumer. Dipakai untuk pola RPC async.
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}")| Tipe | Routing | Contoh Kasus |
|---|---|---|
direct | Routing key harus sama persis | Filter by severity |
fanout | Broadcast ke semua queue | Pub/sub, logs |
topic | Routing key dengan pattern wildcard | Hierarchical routing |
headers | Berdasarkan header message, bukan routing key | Complex routing |
Topic exchange memakai routing key dengan pattern yang dipisahkan titik. Wildcard: * (satu kata), # (nol atau lebih kata).
# 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='#')# 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 adalah struktur data log append-only yang mendukung consumer group. Cocok untuk message queue ringan tanpa broker terpisah.
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 memungkinkan multiple consumer membagi stream secara paralel, mirip work queue.
# 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)Pesan yang belum di ack tetap berada di pending entry list (PEL). Kalau consumer mati, consumer lain bisa claim pesannya.
# 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
)| Aspek | Redis Streams | RabbitMQ |
|---|---|---|
| Kompleksitas | Rendah (in-memory) | Sedang (broker terpisah) |
| Persistence | Optional (RDB/AOF) | Default persistent |
| Routing | Linear, consumer group | Exchange (direct, fanout, topic, header) |
| Throughput | Sangat tinggi | Tinggi |
| Fitur lanjutan | Terbatas | Rich (DLX, priority, TTL) |
| Cocok untuk | Streaming, log | Kompleks routing, enterprise |
SQS adalah managed queue service dari AWS. Tipe: Standard (at-least-once, tidak terurut) dan FIFO (exactly-once, terurut).
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']
)Saat consumer menerima pesan, pesan jadi invisible untuk consumer lain selama visibility timeout. Kalau tidak di ack (delete), pesan akan muncul lagi.
# 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 adalah pub/sub messaging. Topic mengirim ke multiple subscriber (SQS, Lambda, HTTP, email).
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 adalah pola di mana satu SNS topic mengirim ke multiple SQS queue, setiap queue mendapat copy pesan.
+---> SQS Queue A (email service)
SNS Topic ----> +---> SQS Queue B (inventory service)
+---> SQS Queue C (analytics)DLQ adalah queue tempat pesan yang gagal diproses dipindahkan setelah mencapai batas retry. Mencegah pesan beracun (poison pill) mengulang terus.
# 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# 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'
})
}
)| Strategi | Cara Kerja | Cocok Untuk |
|---|---|---|
| Fixed delay | Tunggu waktu tetap (mis. 5 detik) antar retry | Sederhana, error sementara |
| Linear backoff | Tunggu makin lama secara linear | Beban moderat |
| Exponential backoff | Tunggu 2x lebih lama setiap retry | Sistem sibuk, menghindari overload |
| Exponential + jitter | Exponential backoff dengan random noise | Distribusi beban, mencegah thundering herd |
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))# 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 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).
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")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| Garansi | Arti | Trade-off |
|---|---|---|
| At-most-once | Pesan mungkin hilang, tidak duplikat | Cepat tapi bisa kehilangan data |
| At-least-once | Pesan tidak hilang, mungkin duplikat | Butuh idempotency |
| Exactly-once | Tidak hilang, tidak duplikat | Sulit, biasanya simulasi dengan at-least-once + idempotency |
Login atau daftar akun gratis untuk membaca cheat sheet ini.