Catatan ini menggabungkan lima slide kuliah AAT tentang Message Orientation/Publish-Subscribe: materi dasar messaging & JMS (2024), materi AMQP/RabbitMQ (2023 & 2021), materi Apache Kafka (2024), dan materi perbandingan RabbitMQ vs Kafka plus MQTT (2024). Karena kelima slide adalah variasi/pelengkap topik yang sama, konsep yang tumpang tindih (terutama AMQP model dan contoh kode RabbitMQ yang muncul identik di dua slide) digabung menjadi satu penjelasan agar tidak diulang berkali-kali, namun tetap mempertahankan detail teknis dari masing-masing sumber. Topik ini relevan terutama untuk capaian belajar “mendesain layanan terdistribusi & microservices” karena messaging (message queue & pub/sub) adalah salah satu mekanisme komunikasi antar service yang paling umum dipakai selain REST API dan RPC.

Konsep Dasar Messaging

Messaging system adalah fasilitas peer-to-peer (proses ke proses) yang memungkinkan client saling mengirim pesan. Pesan dikirimkan ke sebuah agen perantara komunikasi (broker), sehingga messaging system pada dasarnya menyediakan loosely coupled communication antara sender dan receiver — keduanya tidak perlu tahu detail satu sama lain, cukup tahu cara berkomunikasi dengan perantaranya.

Sistem berbasis messaging dapat berupa:

  • Direct message passing: antar proses bertukar pesan secara langsung
  • Message queue
  • Pub/sub

Ketiganya berbeda dalam hal decoupling yang diberikan. Tabel berikut (dari Zimmermann/EIP) membandingkan berbagai abstraksi komunikasi dari sisi space/time decoupling (apakah sender & receiver perlu aktif bersamaan) dan synchronization decoupling (apakah pengirim harus menunggu):

AbstractionSpace/time decouplingSynchronization decoupling
Message passingNoVaries
RPC/RMINoInvoker is blocked
Async RPC/RMINoYes
Tuple spacesYesReader is blocked
Message queuingYesVaries
Pub/subYesYes

Terlihat bahwa message queuing dan pub/sub adalah dua mekanisme yang memberikan decoupling paling lengkap (baik ruang/waktu maupun sinkronisasi) dibanding RPC biasa.

Komponen Messaging System

  • Message: struktur data yang dikirim, terdiri atas Header, Properties, dan Body.
  • Message Queue: tempat penyimpanan pesan sementara.
  • Producers: client/aplikasi yang mengirimkan pesan ke queue.
  • Consumers: client/aplikasi yang mengambil pesan dari queue. Mode interaksi bisa pull (consumer aktif mengambil) atau push (broker mendorong pesan ke consumer).
  • Message Broker: aplikasi/server yang mengelola message queue — menerima pesan dari producer (send) dan meneruskannya ke consumer (recv).

Jenis Message (dari sisi desain)

Dari sisi desain aplikasi, message dapat dikategorikan menjadi beberapa jenis (mengacu pada pola dari Enterprise Integration Patterns dan Patterns for API Design, Zimmermann 2022):

  • Command Message: berisi perintah/aksi yang harus dijalankan penerima.
  • Event Message: berisi notifikasi bahwa sebuah event telah terjadi.
  • Document Message: berisi notifikasi event lengkap dengan atribut entity, bukan hanya perubahannya saja (berbeda dengan Event Message yang biasanya hanya membawa info perubahan).
  • Query Message: berisi query informasi (permintaan data).

Contoh interaksi antara command, event, query, dan document dalam arsitektur event-driven microservice (Rocha, 2021): User membuat order melalui UI → UI melakukan query ke Inventory Service (Query Message) → UI mem-publish CreateOrderCommand (Command Message) ke Message Queue → Order Service mengonsumsi command tersebut, lalu mem-publish OrderCreatedEvent (Event Message) ke Inventory Service dan OrderDocument (Document Message, berisi seluruh atribut order) ke Notification Service.

Messaging Domain: Point-to-Point vs Publish-Subscribe

Ada dua model utama pertukaran pesan (messaging domain):

Point to Point (PTP)

Pada PTP, destinasi pesan dinyatakan eksplisit menggunakan konsep message queue:

  • Sender mengirim pesan ke sebuah message queue tertentu, sehingga secara implisit menyatakan tujuan receiver.
  • Receiver memonitor message queue tertentu dan mengambil pesan darinya.

Karakteristik PTP:

  • Sebuah pesan hanya dikonsumsi oleh 1 receiver (walau ada banyak receiver yang memonitor queue yang sama — pesan akan dibagi/di-load-balance di antara mereka, bukan diduplikasi).
  • Sender tidak menunggu receiver (asynchronous by nature).
  • Receiver dapat meng-acknowledge pemrosesan pesan.
  • Pesan dapat dikonsumsi baik secara synchronous maupun asynchronous.

Publish-Subscribe (Pub/Sub)

Pada pub/sub, pesan dikirim ke sebuah resource bersama oleh publisher, dan multiple subscriber dapat mengonsumsi pesan yang sama (broadcast/multicast, bukan dibagi seperti PTP). Subscriber menyatakan pesan yang diminatinya dengan menspesifikasikan filter.

Komponen sebuah sistem pub/sub secara umum terdiri atas:

  • Publisher: membuat notifikasi (event instance) berdasarkan sebuah situation dan mengirimkannya ke Pub/Sub Service.
  • Subscriber: mengeluarkan (issue) subscription request/response ke Pub/Sub Service, lalu menerima notifikasi yang relevan.
  • Subscription Manager: mengelola subscription atas nama publisher, disimpan di subscription storage.
  • Notification Engine: mencocokkan (match) event dengan subscription yang ada, dan mengirim notifikasi ke Notification Consumer yang sesuai; event yang masuk dapat disimpan di event storage.
  • Notification Consumer: menerima notifikasi dari engine lalu meneruskan (notify) ke subscriber.

Variasi/Karakteristik Publish-Subscribe System

Implementasi pub/sub bisa sangat bervariasi tergantung pilihan desain berikut:

  • Persistency: apakah event dan subscription disimpan secara persisten atau tidak (memengaruhi apakah subscriber yang sedang offline tetap bisa menerima pesan yang terlewat).
  • Filtering mechanisms: bagaimana subscriber menyaring pesan yang relevan — berdasarkan source (sumber pesan), header, atau content (isi pesan, content-based filtering/routing).
  • Broker hierarchy: topologi broker yang digunakan —
    • Single broker: satu broker saja.
    • Clustered broker: beberapa broker yang setara, saling berbagi beban.
    • Hierarchical/partitioned broker: broker disusun berjenjang atau dipartisi.
  • Event delivery guarantee: jaminan pengiriman pesan (dibahas detail di bawah).
  • Acknowledgment: auto acknowledgment (pesan dianggap terkirim otomatis begitu diterima consumer) vs manual acknowledgment (consumer secara eksplisit mengonfirmasi setelah selesai memproses pesan).

Tabel Delivery Guarantee

Jenis GuaranteePenjelasan
Exactly onceSetiap pesan dijamin sampai ke consumer tepat satu kali — tidak hilang dan tidak terduplikasi. Paling mahal untuk diimplementasikan (butuh koordinasi/transaksi ekstra).
At least oncePesan dijamin sampai, tapi bisa terduplikasi (misal karena retry setelah ack hilang/timeout). Consumer harus siap menangani duplikasi (idempotent).
At most oncePesan dikirim maksimal satu kali — bisa saja pesan hilang (tidak ada retry), tapi tidak akan pernah terduplikasi.
No guarantee / best effortTidak ada jaminan sama sekali; pesan bisa hilang atau terduplikasi tanpa mekanisme kompensasi.

Aspek Reliabilitas Lain yang Perlu Diperhatikan

Ada trade-off antara data safety dan performance, terutama terkait penanganan kegagalan pada tiga titik: pengiriman pesan dari producer, pengambilan pesan oleh consumer, dan kegagalan broker itu sendiri. Mekanisme yang umum digunakan:

  • Publisher confirms: broker memberikan acknowledgment (ack) ke producer untuk setiap pesan yang diterimanya, sehingga producer tahu pesannya sudah aman di broker.
  • Persistent messages & message queues: pesan dan queue disimpan ke disk (bukan hanya di memori) agar bertahan saat broker restart/crash.
  • Consumer manual ack: consumer baru mengonfirmasi (ack) setelah benar-benar selesai memproses pesan, sehingga jika consumer crash di tengah pemrosesan, pesan akan di-redeliver.

Ada juga trade-off antara availability dan performance, terkait konfigurasi cluster (berapa banyak node, bagaimana replikasi dilakukan) dan cara penanganan failure pada cluster tersebut.

Ragam Platform Message Queue

Beberapa platform message queue/pub-sub yang umum digunakan di industri:

  • RabbitMQ — AMQP broker
  • ActiveMQ — JMS broker
  • ZeroMQ — messaging library (bukan broker penuh)
  • RocketMQ — mendukung OpenMessaging, JMS
  • Kafka — high throughput messaging
  • NATS — simplified messaging platform
  • Mosquitto — implementasi MQTT
  • Layanan cloud proprietary: Google Cloud Pub/Sub, Amazon SQS/SNS, Azure Service Bus, Alicloud SMQ

Java Messaging Service (JMS)

JMS adalah spesifikasi (bukan implementasi) yang memungkinkan pengembangan message service di lingkungan Java secara portable antar vendor. Implementasi konkret dari spesifikasi JMS disebut JMS Provider (contoh: ActiveMQ).

JMS menyediakan 2 jenis messaging domain/resource:

  • Message Queues → untuk model Point to Point
  • Topics → untuk model Publish Subscribe

Administered Objects

Connection Factories dan Destinations (Queue/Topic) adalah administered objects — dikelola secara administratif (misal via tool j2eeadmin), bukan dibuat dari kode program. Cara pengelolaannya berbeda-beda antar vendor, tetapi diakses dari kode aplikasi melalui interface portable yang sama (via lookup JNDI).

$ j2eeadmin -addJmsFactory jndi_name queue
$ j2eeadmin -addJmsFactory jndi_name topic
$ j2eeadmin -addJmsDestination queue_name queue
$ j2eeadmin -addJmsDestination topic_name topic

Ada 2 jenis Connection Factory: QueueConnectionFactory dan TopicConnectionFactory. Client mengaksesnya via JNDI lookup:

Context ctx = new InitialContext(); // get the JNDI context
 
QueueConnectionFactory queueConnectionFactory =
    (QueueConnectionFactory) ctx.lookup("QueueConnectionFactory");
 
TopicConnectionFactory topicConnectionFactory =
    (TopicConnectionFactory) ctx.lookup("TopicConnectionFactory");
 
Queue myQueue = (Queue) ctx.lookup("MyQueue");
Topic myTopic = (Topic) ctx.lookup("MyTopic");

JMS API Programming Model

Alur objek dalam JMS API: Connection Factory membuat Connection → Connection membuat Session → Session membuat Message Producer (mengirim pesan ke Destination) dan Message Consumer (menerima pesan dari Destination).

  • Connections: merepresentasikan koneksi ke JMS Provider. Ada 2 jenis: QueueConnection dan TopicConnection.
  • Sessions: merepresentasikan single-thread context yang memproduksi/mengonsumsi message; menyediakan dukungan transaksional dan menserialisasi eksekusi message listener. Ada QueueSession dan TopicSession.
TopicSession topicSession = topicConnection.createTopicSession(
    false, Session.AUTO_ACKNOWLEDGE); // non-transacted, automatic ack
  • Message Producer: memproduksi pesan yang dikirim ke destination. Ada QueueSender (method send) dan TopicPublisher (method publish):
QueueSender queueSender = queueSession.createSender(myQueue);
TopicPublisher topicPublisher = topicSession.createPublisher(myTopic);
 
queueSender.send(message);
topicPublisher.publish(message);
  • Message Consumer: objek yang menerima pesan. Ada QueueReceiver dan TopicSubscriber. TopicSubscriber dapat dibuat durable, artinya dapat menerima pesan yang dipublish selagi subscriber sedang tidak aktif (di-buffer oleh broker sampai subscriber kembali online).

Synchronous consumption — consumer secara aktif memanggil receive() dan menunggu (blocking) hingga pesan datang:

QueueReceiver queueReceiver = queueSession.createReceiver(myQueue);
TopicSubscriber topicSubscriber = topicSession.createSubscriber(myTopic);
 
queueConnection.start();
Message m = queueReceiver.receive();
 
topicConnection.start();
Message m2 = topicSubscriber.receive(1000); // time out after a second

Catatan: message tidak akan di-deliver hingga connection di-start().

Asynchronous consumption — aplikasi mendefinisikan sebuah message listener yang meng-implement interface MessageListener, lalu diasosiasikan ke consumer:

public interface MessageListener {
    public void onMessage(Message message);
}
 
TopicListener topicListener = new TopicListener();
topicSubscriber.setMessageListener(topicListener);

Message Selector

Message selector digunakan untuk memfilter pesan mana yang boleh sampai ke consumer tertentu. Yang penting: filtering dilakukan oleh JMS provider (di sisi broker), bukan oleh aplikasi consumer — sehingga lebih efisien karena pesan yang tidak relevan tidak perlu dikirim ke consumer sama sekali. Selector dinyatakan sebagai statement yang merupakan subset dari sintaks SQL92 conditional expression, dan dapat diberikan sebagai parameter pada method createReceiver, createSubscriber, dan createDurableSubscriber.

Tipe Message JMS

Pesan JMS terdiri atas header, properties, dan body. Ada 5 jenis message yang didefinisikan API JMS:

  • TextMessage: body berupa teks (misalnya XML)
  • MapMessage: sekumpulan pasangan name/value
  • BytesMessage: sebuah stream of bytes
  • StreamMessage: sebuah stream of primitive Java values, diisi dan dibaca secara sekuensial
  • ObjectMessage: sebuah objek Java yang Serializable

AMQP (Advanced Message Queuing Protocol)

AMQP adalah sebuah open standard pada application layer untuk messaging & pub/sub, distandardisasi oleh OASIS (amqp.org). Key features-nya: message orientation, queuing, routing, reliability, dan security. AMQP mendefinisikan behaviour dari message provider, client, broker, dan wire-level protocol-nya secara detail — sehingga (berbeda dengan JMS yang hanya spesifikasi API) client dan broker dari vendor yang berbeda tetap bisa saling berkomunikasi karena protokolnya sudah distandardisasi hingga level wire-protocol. Standard port AMQP adalah 5672/tcp.

Basic pattern yang didukung AMQP: Request-response dan Pub/sub.

AMQP Model: Broker, Node, Container, Connection, Channel, Session

AMQP model mendefinisikan istilah-istilah berikut:

  • Message Broker: server message tempat client terhubung.
  • Node: entity yang bertugas menyimpan atau mengirim message, misalnya: Producer, Consumer, Queue, Exchange.
  • Container: entity yang meng-host sebuah node, misalnya Broker application atau Client application.
  • Connection: koneksi fisik yang berasosiasi dengan sebuah User.
  • Channel: koneksi logik yang berasosiasi dengan sebuah Connection.
  • Session: dua channel unidirectional antara 2 node.

Message ditransfer antar node melalui sebuah Link yang menghubungkan sisi source (Src) dan target (Tgt), dan dapat melewati sebuah filter:

(Figure 2.1 – Message Transfer between Nodes: node A mengirim MSG_1..4 melalui Link(Src,Tgt) menuju node B, dengan opsi filter di tengah)

Relasi antara Container dan Node digambarkan sebagai class diagram berikut — sebuah Container (Broker atau Client) meng-host 0..n Node (Producer, Consumer, atau Queue):

(Figure 2.2 – Class Diagram of Concrete Containers and Nodes)

Connection adalah koneksi fisik antara Client App dan Broker (menghubungkan node Consumer di sisi client dengan node Queue di sisi broker):

Di dalam sebuah Connection, dapat dibuat satu atau lebih Session — koneksi logik yang berjalan di atas Connection fisik yang sama:

Contoh alur AMQP end-to-end: Publisher mengirim pesan ke Exchange di sisi Server/Broker, Exchange meneruskan pesan ke Queue yang sesuai, dan Consumer mengambil pesan dari Queue tersebut:

Message, Exchange, Queue

Producer dan Consumer menggunakan Client API untuk mengirim dan menerima pesan dari broker.

  • Exchange: abstraksi yang bertugas menerima pesan dari producer dan meneruskannya ke queue yang sesuai.
  • Queue: menampung pesan untuk diambil consumer secara FIFO.

Message Queue dapat menyimpan message di disk atau di memori; setiap message queue bersifat independen satu sama lain. Properti sebuah message queue:

  • private atau shared
  • durable atau temporary
  • client named atau server named

Contoh penggunaan: shared store-and-forward queue, private reply queue, private subscription queue.

Exchange: Routing Message

Exchange bertugas menentukan routing message ke queue dengan cara memeriksa header dan body pesan. Routing key adalah nilai yang digunakan untuk menentukan queue tujuan. Sebuah exchange dapat bersifat durable, temporary, atau autodelete. Ada 3 (empat, termasuk header) jenis exchange: direct, fanout, topic (dan header).

Direct Exchange — cocok untuk simple/load-balanced/abstracted point-to-point queue delivery. Producer mempublish dengan sebuah routing key, dan queue yang di-bind dengan routing key yang sama akan menerima pesan tersebut:

Contoh: producer mempublish dengan routing key "France" ke sebuah direct exchange, dan hanya consumer yang men-declare queue lalu melakukan queueBind dengan routing key "France" yang akan menerima pesan tersebut (mirip mekanisme publish/subscribe dengan filter berbasis routing key, lihat contoh kode pada bagian RabbitMQ di bawah).

Fanout Exchange — model publish/subscribe murni: pesan yang sama dikirim ke semua queue yang ter-bind ke exchange tersebut (broadcast), tanpa memandang routing key:

Topic Exchange — content-based routing: routing pesan ditentukan oleh kecocokan binding pattern (misalnya pola seperti a.b.c, a.b.*, a.#) terhadap routing key pesan, sehingga satu pesan bisa masuk ke beberapa queue sekaligus berdasarkan pola yang cocok:

Implementasi AMQP

  • OpenAMQ — opensource, berbasis C (openamq.org)
  • RabbitMQ — opensource (rabbitmq.com)
  • Apache Qpid — opensource (qpid.apache.org)

Contoh Kode: RabbitMQ (Point-to-Point)

Contoh dasar mengirim (Send) dan menerima (Recv) pesan langsung ke/dari sebuah queue bernama "hello":

import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.Channel;
 
public class Send {
  private final static String QUEUE_NAME = "hello";
  public static void main(String[] argv) throws Exception {
    ConnectionFactory factory = new ConnectionFactory();
    factory.setHost("localhost");
    Connection connection = factory.newConnection();
    Channel channel = connection.createChannel();
    channel.queueDeclare(QUEUE_NAME, false, false, false, null);
    String message = "Hello World!";
    channel.basicPublish("", QUEUE_NAME, null, message.getBytes());
    System.out.println(" [x] Sent '" + message + "'");
    channel.close();
    connection.close();
  }
}
public class Recv {
  private final static String QUEUE_NAME = "hello";
  public static void main(String[] argv) throws Exception {
    ConnectionFactory factory = new ConnectionFactory();
    factory.setHost("localhost");
    Connection connection = factory.newConnection();
    Channel channel = connection.createChannel();
    channel.queueDeclare(QUEUE_NAME, false, false, false, null);
    System.out.println(" [*] Waiting for messages. To exit press CTRL+C");
    QueueingConsumer consumer = new QueueingConsumer(channel);
    channel.basicConsume(QUEUE_NAME, true, consumer);
    while (true) {
      QueueingConsumer.Delivery delivery = consumer.nextDelivery();
      String message = new String(delivery.getBody());
      System.out.println(" [x] Received '" + message + "'");
    }
  }
}

Contoh Kode: RabbitMQ Publish/Subscribe (Fanout Exchange)

Berbeda dengan contoh di atas, pada pola publish/subscribe producer tidak mengirim pesan langsung ke queue, melainkan ke sebuah exchange (di sini bertipe "fanout"), dan exchange itulah yang menentukan pesan diteruskan ke queue mana. Fanout exchange akan mengirimkan pesan yang sama ke sejumlah queue sekaligus:

public class EmitLog {
  private static final String EXCHANGE_NAME = "logs";
  public static void main(String[] argv) throws Exception {
    ConnectionFactory factory = new ConnectionFactory();
    factory.setHost("localhost");
    Connection connection = factory.newConnection();
    Channel channel = connection.createChannel();
    channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
    String message = getMessage(argv);
    channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes());
    System.out.println(" [x] Sent '" + message + "'");
    channel.close();
    connection.close();
  }
}
public class ReceiveLogs {
  private static final String EXCHANGE_NAME = "logs";
  public static void main(String[] argv) throws Exception {
    ConnectionFactory factory = new ConnectionFactory();
    factory.setHost("localhost");
    Connection connection = factory.newConnection();
    Channel channel = connection.createChannel();
    channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
    String queueName = channel.queueDeclare().getQueue();
    channel.queueBind(queueName, EXCHANGE_NAME, "");
    System.out.println(" [*] Waiting for messages. To exit press CTRL+C");
    QueueingConsumer consumer = new QueueingConsumer(channel);
    channel.basicConsume(queueName, true, consumer);
    while (true) {
      QueueingConsumer.Delivery delivery = consumer.nextDelivery();
      String message = new String(delivery.getBody());
      System.out.println(" [x] Received '" + message + "'");
    }
  }
}

Perhatikan bahwa setiap instance ReceiveLogs membuat queue anonim/sementara sendiri (channel.queueDeclare().getQueue()) lalu mem-bind-nya ke exchange "logs" — inilah yang membuat setiap subscriber menerima salinan pesan yang sama (ciri khas pub/sub), berbeda dengan contoh Send/Recv di atas yang berbagi satu queue bernama tetap (ciri khas PTP, pesan dibagi habis, bukan diduplikasi).

Contoh kode di atas untuk direct exchange (routing key "France") juga menunjukkan pola serupa: producer publish ke exchange dengan routing key tertentu, consumer membuat queue lalu queueBind dengan routing key yang sama:

// Producer
channel.exchangeDeclare(EXCHANGE_NAME, "direct");
channel.basicPublish(EXCHANGE_NAME, "France", null, message.getBytes());
 
// Consumer
String queueName = channel.queueDeclare().getQueue();
channel.queueBind(queueName, EXCHANGE_NAME, "France");

Apache Kafka

Apache Kafka dikembangkan di LinkedIn mulai 2011, diimplementasikan menggunakan Scala dan Java. Kafka dirancang dengan tujuan (goal):

  • high throughput
  • mendukung realtime processing
  • fault tolerant
  • fast
  • mampu menangani large data

Konsep Dasar

Komponen dasar Kafka sama seperti messaging system pada umumnya: Producer, Consumer, dan Broker (dikumpulkan dalam sebuah Kafka cluster).

Zookeeper berperan mengelola state dan koordinasi antar broker dan consumer dalam cluster Kafka (pada versi Kafka yang lebih baru, sebagian peran ini bisa digantikan oleh KRaft, namun materi kuliah ini mengacu pada arsitektur berbasis Zookeeper).

Topic, Partition, dan Offset

Data pada Kafka disimpan pada topics. Sebuah topic dibagi ke dalam beberapa partition, yang masing-masing direplikasi. Setiap message pada sebuah partition diberikan id unik yang disebut offset.

Poin penting:

  • Producer saat mengirim message ke sebuah topik dapat menuliskannya ke salah satu partisi yang ada di topik tersebut.
  • Message dalam satu partisi yang sama terjamin urutannya (order) berdasarkan offset id — namun urutan antar partisi berbeda tidak dijamin. Ini adalah trade-off khas Kafka: pararelisme tinggi (banyak partisi) mengorbankan jaminan urutan global.

Jumlah partisi dapat dikonfigurasi, dan jumlah partisi menentukan tingkat paralelisasi yang bisa dicapai per consumer group — semakin banyak partisi, semakin banyak consumer dalam satu consumer group yang bisa memproses topik tersebut secara paralel.

Consumer Group

Consumer Group (CG): pembacaan message oleh consumer dijamin unik per consumer group — artinya sebuah pesan hanya akan diproses oleh satu consumer dalam satu consumer group yang sama, mirip semantik PTP di dalam grup tersebut. Namun, consumer group yang berbeda akan menerima salinan pesan yang sama dari topik yang sama — inilah yang membuat Kafka bisa mendukung dua pola sekaligus:

  • Publish-subscribe: setiap consumer memiliki consumer group sendiri-sendiri → semua consumer menerima semua pesan (broadcast).
  • Load balancing queue: semua consumer terhubung ke consumer group yang sama → pesan dibagi (load balanced) di antara consumer-consumer tersebut.

Consumer menggunakan offset pointer untuk mencatat sejauh mana pesan sudah dibaca pada tiap partisi — pointer ini yang membedakan Kafka dari message queue tradisional, karena posisi baca disimpan per consumer group, bukan dihapus dari broker setelah dibaca (pesan tetap ada di log Kafka hingga masa retensi habis, sehingga consumer bisa “replay” dari offset awal jika diperlukan).

Jumlah partisi juga memengaruhi bagaimana beban dibagi antar consumer dalam consumer group yang berbeda — pada contoh konfigurasi 4 partisi (P0-P3) tersebar di 2 server, Consumer Group A (2 consumer) dan Consumer Group B (4 consumer) masing-masing membaca dari keempat partisi tersebut secara independen.

Replikasi: Leader dan Follower

Setiap partisi dapat memiliki replika, sesuai konfigurasi replication factor. Dalam sebuah partisi yang memiliki beberapa replika:

  • Satu replika akan menjadi leader, sisanya menjadi follower/backup.
  • Semua write dan read hanya dilakukan pada leader (follower hanya mereplikasi data dari leader).
  • Jika leader gagal (fail), salah satu follower akan dipromosikan menjadi leader baru — mekanisme inilah yang membuat Kafka fault tolerant.

Menjalankan Kafka & Konfigurasi Dasar

Distribusi Apache Kafka menyediakan berbagai script untuk menjalankan Zookeeper, server (broker), topic management, console publisher, dan console consumer. Contoh membuat topic via CLI:

> bin/kafka-topics.sh --create --zookeeper localhost:2181 \
    --replication-factor 1 --partitions 1 --topic test

Auto-create topic dapat diaktifkan dengan properti auto.create.topics.enable=true pada config/server.properties.

Contoh kode Producer (API lama, kafka.javaapi.producer.Producer<K,V>):

Properties props = new Properties();
props.put("metadata.broker.list", "...");
ProducerConfig config = new ProducerConfig(props);
 
Producer p = new Producer(ProducerConfig config);
KeyedMessage<K, V> msg = ...; // cf. later slides
p.send(KeyedMessage<K,V> message);
class kafka.javaapi.producer.Producer<K,V> {
    public Producer(ProducerConfig config);
 
    /** Sends the data to a single topic, partitioned by key, using either the
      * synchronous or the asynchronous producer. */
    public void send(KeyedMessage<K,V> message);
 
    /** Use this API to send data to multiple topics. */
    public void send(List<KeyedMessage<K,V>> messages);
 
    /** Close API to close the producer pool connections to all Kafka brokers. */
    public void close();
}

Konfigurasi Producer yang penting:

KonfigurasiKeterangan
client.ididentifies producer app, e.g. in system logs
producer.typeasync atau sync
request.required.acksacking semantics
serializer.classconfigure encoder
metadata.broker.listbootstrapping list of brokers

Konfigurasi Consumer yang penting:

KonfigurasiKeterangan
group.idassigns an individual consumer to a “group”
zookeeper.connectto discover brokers/topics/etc., dan menyimpan consumer state (high-level consumer API)
fetch.message.max.bytesjumlah bytes pesan yang diambil per partisi; harus >= message.max.bytes broker

Kenapa Kafka Cocok untuk High-Throughput, Real-Time, Large-Scale, Fault-Tolerant

  • High throughput: penulisan bersifat append-only ke log per partisi (sequential disk write yang sangat cepat), dan partisi memungkinkan paralelisme baca/tulis.
  • Realtime processing: consumer dapat membaca pesan segera setelah ditulis (near real-time), cocok untuk stream processing.
  • Large scale: topic dapat dipartisi lintas banyak broker, sehingga throughput dan kapasitas penyimpanan dapat diskalakan secara horizontal.
  • Fault tolerant: replikasi partisi dengan mekanisme leader-follower memastikan data tidak hilang meski satu broker mati.

MQTT

MQTT (MQ Telemetry Transport) dikembangkan oleh IBM sejak 1999, menjadi open standard OASIS pada 2014 (MQTT 3.1.1). MQTT adalah Client/Server publish/subscribe messaging transport protocol — penting dicatat bahwa MQTT bukan protokol tentang message queue, melainkan murni protokol pub/sub. Target utamanya adalah embedded device dan IoT.

Fitur MQTT

  • Lightweight & binary protocol — dirancang untuk:
    • jaringan dengan low bandwidth, high latency
    • low powered device (baterai terbatas)
    • payload berupa UTF-8 encoded string
  • Berjalan over TCP
  • Mengikuti publish/subscribe pattern — client melakukan subscribe ke topics
  • Memiliki Keep-Alive timer (mekanisme ping berkala agar broker tahu client masih terhubung)
  • Mendukung QoS (Quality of Service) bertingkat

Model Publish/Subscribe MQTT

Alur kerja MQTT:

  1. Client subscribe ke sebuah topik.
  2. Client lain mempublish message ke topik spesifik tersebut.
  3. Client terhubung ke server/broker melalui CONNECT message.
  4. Broker mengautentikasi client.
  5. Saat sebuah client mempublish message ke topik tertentu, broker meneruskan message tersebut ke semua client yang subscribe ke topik itu.
  6. Broker dapat menyimpan message dengan flag khusus: RETAIN dan WILL message.

Jenis Message (Control Packet) MQTT

MQTT mendefinisikan jenis-jenis message berikut: CONNECT, CONNACK, PUBLISH, PUBACK, PUBREC, PUBREL, PUBCOMP, SUBSCRIBE, SUBACK, UNSUBSCRIBE, UNSUBACK, PINGREQ, PINGRESP, DISCONNECT.

Format fixed header MQTT terdiri atas 2 byte: byte pertama berisi Message Type, DUP flag, QoS level, dan RETAIN flag; byte kedua berisi Remaining Length:

QoS (Quality of Service) pada MQTT

Level QoSNamaMekanisme
0Best Efforthanya PUBLISH (kirim sekali, tanpa konfirmasi — setara at most once)
1At Least OncePUBLISH + PUBACK (penerima mengonfirmasi, pengirim retry jika ack tidak diterima — bisa terjadi duplikasi)
2Exactly OncePUBLISH + PUBREC + PUBREL + PUBCOMP (4-way handshake untuk menjamin pesan sampai tepat satu kali)

Terlihat bahwa 3 level QoS MQTT ini persis memetakan ke tiga dari empat kategori delivery guarantee yang dibahas pada bagian variasi pub/sub di atas (best effort, at-least-once, exactly-once).

RETAIN Message

Sebuah message saat dipublish dapat diberi flag RETAIN. Jika flag ini diset, server/broker akan menyimpan message tersebut dan akan meneruskannya ke client baru yang melakukan subscribe ke topik yang sesuai — walaupun subscriber tersebut baru bergabung setelah pesan dipublish. Ini berguna agar subscriber baru langsung mendapat “state terakhir” dari sebuah topik tanpa harus menunggu publish berikutnya.

Last Will and Testament (LWT)

Saat melakukan CONNECT, sebuah client dapat menyertakan sebuah message LAST WILL AND TESTAMENT. Pesan “wasiat” ini akan otomatis dikirimkan oleh broker ke semua client lain yang relevan apabila client tersebut terputus secara tiba-tiba/tidak normal (misalnya koneksi jaringan putus tanpa sempat mengirim DISCONNECT). Mekanisme ini memungkinkan sistem lain mendeteksi bahwa sebuah device IoT “mati” secara mendadak.

Perbandingan: RabbitMQ vs Apache Kafka

RabbitMQ

  • General purpose messaging broker, memenuhi standar AMQP compliance.
  • Mendukung beragam pola interaksi producer-consumer.
  • Memiliki exchange dan queue: producer publish ke exchange (dengan konfigurasi tertentu), consumer mengambil pesan dari queue.
  • Mendukung baik point-to-point maupun publish-subscribe.
  • Broker menyediakan fasilitas routing — RabbitMQ adalah smart broker, dumb client: logika routing kompleks ada di sisi broker, client relatif sederhana.
  • Secara default message tidak persistent, namun producer dapat men-set persistency per pesan.
  • Mendukung transactional support.

Ilustrasi topologi RabbitMQ: producer mengirim ke exchange via channel, exchange dapat melakukan load-balancing (satu queue dengan banyak consumer berbagi beban) maupun multicast ke banyak queue berbeda berdasarkan binding pattern (topic), dengan mekanisme (N)ack untuk konfirmasi pengiriman:

Apache Kafka

  • Dirancang untuk high volume, durable data.
  • Message akan di-keep hingga durasi tertentu (masa retensi) — Kafka dapat berperan sebagai durable message store, bukan sekadar antrian sementara.
  • Kafka adalah smart client, dumb broker (kebalikan dari RabbitMQ): client-lah yang mencatat sendiri posisi pembacaan pesan (offset, dahulu disimpan di Zookeeper), sehingga client dapat me-replay message dari awal kapan saja.

Tabel Perbandingan RabbitMQ vs Apache Kafka (dan MQTT)

AspekRabbitMQApache KafkaMQTT
Standar/protokolAMQP compliantProtokol proprietary KafkaStandar OASIS, protokol sendiri (bukan message queue)
Filosofi arsitekturSmart broker, dumb clientSmart client, dumb brokerClient/server pub-sub sederhana, ringan
Model utamaPoint-to-point & publish-subscribe (via exchange: direct/fanout/topic/header)Publish-subscribe berbasis topic-partition & consumer groupPublish-subscribe murni berbasis topic
Penyimpanan pesanDefault tidak persistent (bisa diset persistent), pesan dihapus setelah di-ackDurable, disimpan hingga masa retensi tertentu (dapat di-replay)Umumnya tidak persistent, kecuali dengan flag RETAIN
Skala/use caseMessaging umum, decouple web request dari long-running backend process, request-reply, load balancingHigh volume, stream processing skala besar (log processing, big data analytics, event sourcing)Embedded device & IoT, jaringan low-bandwidth/low-power
Replikasi/HAClustering — queue dapat di-mirror antar node (ha-mode), exchange & binding ada di semua nodeReplikasi per-partisi dengan skema leader-followerTidak spesifik terkait replikasi broker (fokus pada protokol client-broker)
Contoh use caseTicketing (pencarian harga tiket dilempar dari frontend ke backend via RabbitMQ)Log processing, Telecom call data record processing, input Hadoop/Spark, event sourcingSensor IoT, telemetry, sistem embedded

Replikasi dan Fault Tolerant: RabbitMQ Clustering

Pada RabbitMQ Clustering:

  • Semua node yang terhubung ke sebuah cluster mereplikasi semua data/state yang diperlukan untuk operasi broker.
  • Client dapat terhubung ke node mana pun (setiap node adalah equal peer).
  • Exchange dan binding ada pada setiap node, tetapi queue hanya ada pada node tempat ia di-declare.
  • Queue dapat di-mirror ke beberapa node menggunakan policy dengan atribut ha-mode.
  • Setiap queue memiliki master replika pada satu node tertentu; semua operasi ke queue tersebut dikirim ke master replika ini terlebih dahulu. Lokasi master replika ditentukan oleh: node dengan bound master terkecil, node tempat client men-declare queue, atau secara random.

Bandingkan dengan replikasi Kafka yang berbasis partisi (leader-follower per partisi, dibahas di bagian Kafka) — RabbitMQ mereplikasi di level queue (via mirroring), sedangkan Kafka mereplikasi di level partisi topic.

Flashcard

flashcards Apa perbedaan utama antara messaging Point-to-Point dan Publish-Subscribe? :: Pada Point-to-Point, sebuah pesan hanya dikonsumsi oleh 1 receiver (dibagi/load-balanced antar receiver pada queue yang sama); pada Publish-Subscribe, pesan dikirim ke resource bersama dan dapat dikonsumsi oleh banyak subscriber sekaligus (broadcast/multicast). Sebutkan 4 kategori event delivery guarantee pada sistem pub/sub :: Exactly once, at least once, at most once, dan no guarantee/best effort. Apa itu message selector pada JMS, dan di mana filtering dilakukan? :: Message selector adalah mekanisme filter pesan (subset sintaks SQL92) yang dapat diberikan ke createReceiver/createSubscriber; filtering dilakukan oleh JMS provider (broker), bukan oleh aplikasi consumer. Sebutkan 5 jenis message body yang didefinisikan JMS API :: TextMessage, MapMessage, BytesMessage, StreamMessage, dan ObjectMessage. Apa fungsi Exchange pada AMQP, dan sebutkan 3 jenis exchange utamanya :: Exchange bertugas menentukan routing message ke queue yang sesuai berdasarkan routing key/header/body. Tiga jenis utamanya: Direct (point-to-point berbasis routing key persis), Fanout (broadcast ke semua queue ter-bind, model pub/sub), dan Topic (content-based routing berdasarkan pattern matching). Pada Apache Kafka, apa fungsi offset dan mengapa pesan dalam satu partisi terjamin urutannya sedangkan antar partisi tidak? :: Offset adalah id unik berurutan untuk tiap pesan dalam sebuah partisi, digunakan consumer sebagai pointer pembacaan. Urutan hanya dijamin dalam satu partisi karena pesan ditulis sekuensial di situ; antar partisi berbeda tidak dijamin karena penulisan tiap partisi independen satu sama lain. Jelaskan perbedaan perilaku Consumer Group pada Kafka untuk mode publish-subscribe vs load balancing queue :: Jika tiap consumer punya consumer group sendiri-sendiri, semua consumer menerima semua pesan (mode publish-subscribe/broadcast). Jika semua consumer bergabung ke consumer group yang sama, pesan dibagi/di-load-balance di antara mereka (mode queue). Apa peran leader dan follower dalam replikasi partisi Kafka? :: Setiap partisi yang direplikasi memiliki satu leader (menangani semua read/write) dan sisanya follower/backup (mereplikasi data dari leader); jika leader gagal, salah satu follower dipromosikan menjadi leader baru. Sebutkan perbedaan filosofi arsitektur RabbitMQ vs Kafka dalam satu istilah masing-masing :: RabbitMQ adalah “smart broker, dumb client” (logika routing kompleks di broker); Kafka adalah “smart client, dumb broker” (client mencatat sendiri offset pembacaan dan bisa replay pesan). Jelaskan tiga level QoS pada MQTT beserta mekanismenya :: QoS 0 (Best Effort): hanya PUBLISH tanpa konfirmasi. QoS 1 (At Least Once): PUBLISH+PUBACK, bisa terjadi duplikasi. QoS 2 (Exactly Once): PUBLISH+PUBREC+PUBREL+PUBCOMP, handshake 4 langkah menjamin pesan sampai tepat sekali. Apa fungsi flag RETAIN dan Last Will and Testament (LWT) pada MQTT? :: RETAIN membuat broker menyimpan pesan terakhir pada sebuah topik dan mengirimkannya ke subscriber baru yang subscribe setelahnya. LWT adalah pesan yang didaftarkan client saat CONNECT, yang akan otomatis dikirim broker ke client lain jika client tersebut terputus secara tiba-tiba. Mengapa MQTT dikatakan “bukan protokol message queue”? :: Karena MQTT murni protokol publish/subscribe berbasis topic tanpa konsep antrian pesan (queue) seperti pada AMQP/JMS — pesan diteruskan langsung ke subscriber yang aktif (dengan opsi RETAIN untuk subscriber baru), bukan disimpan dalam struktur queue untuk diambil satu per satu.