ApsaraMQ for RocketMQ menyediakan dua jenis konsumen: Push consumer dan Simple consumer. Masing-masing menangani pengambilan pesan, konkurensi, dan mekanisme retry secara berbeda. Pilih jenis yang sesuai dengan model pemrosesan dan kebutuhan keandalan Anda.
Jenis konsumen yang harus digunakan
| Skenario | Jenis yang direkomendasikan | Alasan |
|---|---|---|
| Waktu pemrosesan dapat diprediksi, tanpa threading kustom | Push consumer | SDK mengelola pengambilan pesan, konkurensi, dan retry. Cukup daftarkan listener, proses setiap pesan, lalu kembalikan hasilnya. |
| Waktu pemrosesan bervariasi, alur kerja kustom | Simple consumer | Aplikasi Anda mengontrol kapan mengambil pesan, cara mendistribusikannya ke berbagai thread, dan kapan mengonfirmasi penyelesaian. |
Mengganti jenis konsumen tidak memengaruhi sumber daya ApsaraMQ for RocketMQ yang sudah ada atau pemrosesan bisnis.
Tahapan pemrosesan pesan
Kedua jenis konsumen mengikuti siklus hidup tiga tahap:
Receive — Ambil pesan dari server.
Process — Jalankan logika bisnis pada setiap pesan.
Commit — Laporkan hasil (sukses atau gagal) ke server.

Kedua jenis ini berbeda dalam cara menangani setiap tahap:
| Fitur | Push consumer | Simple consumer |
|---|---|---|
| Antarmuka | Callback listener — implementasikan logika di dalam listener dan kembalikan hasil | Aplikasi memanggil operasi API untuk menerima, memproses, dan mengakui pesan |
| Konkurensi | Dikelola oleh SDK | Dikelola oleh aplikasi |
| Fleksibilitas | Sangat terenkapsulasi, kurang fleksibel | Operasi atomik, sangat dapat dikustomisasi |
| Paling cocok untuk | Konsumsi standar dengan waktu pemrosesan yang dapat diprediksi | Alur kerja kustom, distribusi asinkron, atau konsumsi batch |
| Kelas SDK | PushConsumer, LitePushConsumer | SimpleConsumer |
Push consumers
Push consumer mengenkapsulasi logika pengambilan pesan, threading, dan retry. Daftarkan listener pesan saat inisialisasi, dan SDK akan menangani sisanya.
Cara kerja
SDK menggunakan model thread Reactor secara internal:
Thread long-polling bawaan menarik pesan dari server secara asinkron.
Pesan ditempatkan ke dalam antrian cache internal.
SDK mendistribusikan pesan ke thread konsumen, yang kemudian memanggil listener Anda.

Hasil listener
Listener pesan harus mengembalikan salah satu hasil berikut:
| Hasil | Konstanta Java SDK | Perilaku |
|---|---|---|
| Sukses | ConsumeResult.SUCCESS | Server memperbarui progres konsumsi. |
| Gagal | ConsumeResult.FAILURE | Sistem melakukan retry berdasarkan kebijakan retry PushConsumer. |
| Exception dilempar | (dianggap sebagai kegagalan) | Perilaku retry sama seperti kegagalan eksplisit. |
Perilaku timeout
Jika logika pemrosesan hang dan mencegah pesan selesai dalam batas waktu yang diizinkan, SDK secara paksa mengirimkan hasil kegagalan dan menangani pesan sesuai kebijakan retry.
Timeout menyebabkan SDK mengirimkan hasil kegagalan, tetapi thread pemrosesan saat ini mungkin tidak merespons interupsi dan dapat terus berjalan.
Batasan keandalan
Push consumer menentukan sukses atau gagal secara ketat berdasarkan nilai kembali listener. Untuk menjaga jaminan ini:
Proses secara sinkron. Selesaikan seluruh pemrosesan sebelum mengembalikan hasil.
Jangan mendistribusikan ulang pesan. Jangan meneruskan pesan ke thread lain dan mengembalikan hasil sebelum thread tersebut selesai.
Jika listener mengembalikan sukses sebelum pemrosesan selesai dan pemrosesan kemudian gagal, server menganggap pesan telah dikonsumsi dan tidak melakukan retry.
Pengiriman pesan terurut
Ketika kelompok konsumen menggunakan mode konsumsi terurut, Push consumer memanggil listener dalam urutan pesan yang ketat tanpa konfigurasi tambahan. Untuk informasi lebih lanjut, lihat Ordered messages.
Pengiriman terurut memerlukan pemrosesan sinkron. Distribusi asinkron kustom di dalam listener akan membatalkan jaminan pengurutan.
Kapan menggunakan Push consumers
Push consumers paling efektif ketika:
Waktu pemrosesan dapat diprediksi. Durasi yang tidak dapat diprediksi sering memicu timeout, yang menyebabkan pesan duplikat melalui mekanisme retry.
Konsumsi standar sudah cukup. SDK mengontrol model thread dan mengirimkan pesan dengan throughput maksimum. Hal ini menyederhanakan pengembangan tetapi tidak mendukung pemrosesan asinkron atau kontrol laju kustom.
Kelas SDK
ApsaraMQ for RocketMQ menyediakan dua kelas SDK untuk Push consumers:
PushConsumer— Mengonsumsi pesan dari topik standar (non-Lite).LitePushConsumer— Mengonsumsi pesan dari topik Lite, dengan kontrol konsumsi pada tingkat granularitas topik Lite.
Contoh PushConsumer
// Mengonsumsi pesan normal dengan PushConsumer
ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "<your-topic>";
FilterExpression filterExpression = new FilterExpression("<your-filter-tag>", FilterExpressionType.TAG);
PushConsumer pushConsumer = provider.newPushConsumerBuilder()
// Setel kelompok konsumen
.setConsumerGroup("<your-consumer-group>")
// Setel titik akses
.setClientConfiguration(ClientConfiguration.newBuilder().setEndpoints("<your-endpoint>").build())
// Bind langganan
.setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
// Daftarkan listener pesan
.setMessageListener(new MessageListener() {
@Override
public ConsumeResult consume(MessageView messageView) {
// Proses pesan dan kembalikan hasil
return ConsumeResult.SUCCESS;
}
})
.build();Ganti placeholder berikut dengan nilai aktual Anda:
| Placeholder | Deskripsi | Contoh |
|---|---|---|
<your-topic> | Nama topik | order-events |
<your-filter-tag> | Tag pesan untuk filtering | payment |
<your-consumer-group> | Nama kelompok konsumen | order-service-group |
<your-endpoint> | Titik akses server | -- |
Contoh LitePushConsumer
// Mengonsumsi pesan normal dengan LitePushConsumer
ClientServiceProvider provider = ClientServiceProvider.loadService();
LitePushConsumer litePushConsumer = provider.newLitePushConsumerBuilder()
// Setel titik akses
.setClientConfiguration(ClientConfiguration.newBuilder().setEndpoints("<your-endpoint>").build())
// Setel topik
.bindTopic("<your-topic>")
// Setel kelompok konsumen
.setConsumerGroup("<your-consumer-group>")
// Daftarkan listener pesan
.setMessageListener(messageView -> {
// Proses pesan dan kembalikan hasil
return ConsumeResult.SUCCESS;
})
.build();
// Berlangganan ke topik Lite
litePushConsumer.subscribeLite("<your-lite-topic-1>");
litePushConsumer.subscribeLite("<your-lite-topic-2>");Simple consumers
Simple consumer menyediakan operasi API atomik untuk pemrosesan pesan. Aplikasi Anda mengontrol langsung pengambilan pesan, manajemen thread, dan acknowledgment.
Cara kerja
Panggil
ReceiveMessageuntuk menarik batch pesan dari server.Distribusikan pesan ke thread bisnis Anda untuk diproses.
Panggil
AckMessageuntuk setiap pesan yang berhasil diproses.
Jika pemrosesan gagal, jangan kirim acknowledgment. Pesan akan tersedia kembali setelah durasi invisibility pesan berakhir, memicu retry. Untuk informasi lebih lanjut, lihat kebijakan retry SimpleConsumer.
Operasi API
| Operasi | Tujuan | Parameter utama |
|---|---|---|
ReceiveMessage | Menarik pesan dari server | Ukuran batch: jumlah pesan per permintaan. Durasi invisibility pesan: waktu pemrosesan maksimum sebelum pesan dikirim ulang. |
AckMessage | Mengonfirmasi konsumsi yang berhasil | Tidak ada |
ChangeInvisibleDuration | Memperpanjang waktu pemrosesan untuk pesan yang sudah diterima | Durasi invisibility pesan: nilai baru, biasanya digunakan ketika pemrosesan memakan waktu lebih lama dari perkiraan awal. |
Server menggunakan penyimpanan terdistribusi, sehingga ReceiveMessage mungkin mengembalikan hasil kosong meskipun pesan tersedia. Untuk menanganinya, panggil kembali ReceiveMessage atau tingkatkan konkurensi pemanggilan.
Penanganan kegagalan
Tabel berikut menjelaskan bagaimana skenario kegagalan berbeda memengaruhi pengiriman pesan:
| Skenario kegagalan | Perilaku |
|---|---|
| Pemrosesan gagal (ACK tidak dikirim) | Pesan menjadi terlihat kembali setelah durasi invisibility berakhir. Server mengirim ulang pesan tersebut untuk retry. |
| Pemrosesan melebihi durasi periode ketidakterlihatan | Sama seperti tidak mengirim ACK — pesan menjadi terlihat dan dikirim ulang. Gunakan ChangeInvisibleDuration untuk memperpanjang jendela pemrosesan sebelum durasi berakhir. |
| Konsumen crash sebelum mengirim ACK | Pesan dikirim ulang setelah durasi invisibility berakhir. |
Pengiriman pesan terurut
Simple consumer memproses pesan terurut sesuai urutan penyimpanan. Untuk sekelompok pesan yang harus tetap berurutan, pesan berikutnya tidak dapat diambil hingga pesan sebelumnya selesai diproses.
Kapan menggunakan Simple consumers
Simple consumers paling efektif ketika:
Waktu pemrosesan tidak dapat diprediksi. Tentukan durasi invisibility pesan awal saat memanggil
ReceiveMessage, lalu perpanjang denganChangeInvisibleDurationjika diperlukan.Diperlukan alur kerja kustom. SDK tidak memaksakan model threading — implementasikan distribusi asinkron, konsumsi batch, atau pola kustom apa pun.
Kontrol laju penting. Kode Anda menentukan kapan dan seberapa sering memanggil
ReceiveMessage, memberikan kendali langsung atas throughput.
Kelas SDK
ApsaraMQ for RocketMQ menyediakan satu kelas SDK untuk Simple consumers: SimpleConsumer. Kelas ini tidak dapat mengonsumsi pesan dari topik Lite.
Contoh SimpleConsumer
// Mengonsumsi pesan normal dengan SimpleConsumer
ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "<your-topic>";
FilterExpression filterExpression = new FilterExpression("<your-filter-tag>", FilterExpressionType.TAG);
SimpleConsumer simpleConsumer = provider.newSimpleConsumerBuilder()
// Setel kelompok konsumen
.setConsumerGroup("<your-consumer-group>")
// Setel titik akses
.setClientConfiguration(ClientConfiguration.newBuilder().setEndpoints("<your-endpoint>").build())
// Bind langganan
.setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
.build();
try {
// Tarik hingga 10 pesan, tunggu maksimal 30 detik
List<MessageView> messageViewList = simpleConsumer.receive(10, Duration.ofSeconds(30));
messageViewList.forEach(messageView -> {
System.out.println(messageView);
// Akui setiap pesan setelah pemrosesan berhasil
try {
simpleConsumer.ack(messageView);
} catch (ClientException e) {
e.printStackTrace();
}
});
} catch (ClientException e) {
// Tangani kegagalan seperti pembatasan kecepatan, lalu coba lagi pemanggilan receive
e.printStackTrace();
}Praktik terbaik
Kontrol waktu pemrosesan untuk Push consumers
Pastikan pemrosesan pesan tetap dalam ambang batas timeout. Timeout yang sering menyebabkan retry tidak perlu dan pesan duplikat. Jika aplikasi Anda secara rutin menangani tugas berdurasi panjang, beralihlah ke Simple consumer dan atur durasi invisibility pesan yang sesuai.
Migrasi dari LitePullConsumer ke SimpleConsumer
LitePullConsumer adalah jenis konsumen dalam Apache RocketMQ 4.x Remoting SDK. Jenis ini telah dihapus pada SDK gRPC 5.x, dan SDK Go, Java, serta Python versi 5.x tidak menyediakan jenis ini.
Jika sebelumnya Anda menggunakan LitePullConsumer untuk konsumsi berbasis pull, migrasikan ke SimpleConsumer setelah melakukan upgrade ke SDK 5.x. SimpleConsumer menyediakan operasi berikut:
ReceiveMessage: menarik pesan dari server sesuai permintaan.AckMessage: mengonfirmasi konsumsi yang berhasil.ChangeInvisibleDuration: mengubah durasi invisibility pesan untuk mengontrol interval retry.
Pada SDK Go, metode yang sesuai adalah Receive dan Ack. Pada SDK Java, metode yang sesuai adalah receive dan ack.