Pitanja za intervju o redov poruka — RocketMQ deo, 23 pitanja o RocketMQ (11.000 reči, 45 crteža), mora pročitati ko se priprema za intervju

Uvod
11.000 reči, 45 crteža, detaljno objašnjeno 23 čestih pitanja za intervju o RocketMQ. Kandidati koji se pripremaju moraju da ovo nauče, siguran sam da ćeš ovog puta "prebiti" intervjutera (manualni dog). Priredio: Chenmo Wang Er, link za preuzimanje, autor: Sanfen E, link do originalnog teksta.
Svetla verzija je bolja za štampanje, što je i način koji mnogi studenti vole — u štampanom obliku učenje je efikasnije.

Dana 02. 11. 2025. godine počeo sam sa radom na drugom izdanju.
Za česta pitanja, označiću gde se pojavljuju u "Vodiču za intervju iz Jave" — koja kompanija, koji je originalni zadatak, i dodati 🌟, tako da je sadržaj odmah jasan; ako želiš da uštediš vreme, možeš prvo da učiš ova pitanja, brzo saznavši protivnika i mogući ishod.
Razlikujem izvorne odgovore na učna pitanja i objašnjenja principa i temeljnih mehanizama, tako da kandidati razumeju kako i zašto, a istovremeno mogu efikasno da odgovaraju na intervjuu.
Uključujem projekte (Spring Boot+React web projekat podeljen na frontend i backend — Tehnička škola, mikroservisi pmhub, RAG projekat Pametna AI baza znanja — Pai Smart) kako bi se obrazloženje artikulisalo, a intervjuter maksimalno osetio tvoju iskrenost, a ne mehanično učenje.
Popravio sam probleme iz prvog izdanja, uključujući povratne informacije članova sajta, komentare sa sekcije za komentare na sajtu, kao i issue iz GitHub repozitorijuma, kako bi ovaj vodič za intervju bio potpuniji.
Dodao sam neke ponude koje su dobili članovi Ergeove programerske planete, zahvalnice na "Pripremi za intervju" i priznanja za izmenu životopisa, kako bih inspirisao sve i dao više samopouzdanja.
Unapredio sam raspored, dodao crteže, reorganizovao odgovore kako bi bili govornijim i bliži onome što intervjuteri očekuju.

Naravno, dozvoli da imam malu sebičnost — PDF verzija sa planete će biti objavljena mesec dana ranije nego na javnom nalogu, jer su članovi planete platili, pa ih moram prvo da nešto uživaju. Verujem da svi razumiju, jer je online verzija besplatna, a CDN, serveri, domeni, OSS itd. koštaju.
Nemoj ni da pominjem moje vreme i energiju — ako vam se ovo pomaže, molim vas dajte reč, neka i vaše kolege i saučesnici profitiraju.
Uključio sam Ergeov "Napredni put kroz Javu", "Napredni put kroz JVM", "Napredni put kroz konkurentno programiranje", kao i sve verzije "Pripreme za intervju" — ukupno 16 velikih tema: Osnove Jave, Java kolekcije, Java konkurentnost, JVM, Spring, MyBatis, računarske mreže, operativni sistemi, MySQL, Redis, RocketMQ, distribuirani sistemi, mikroservisi, dizajn obrasci, Linux itd. — ukupno preko 400.000 reči, više od 2000 crteža, zaista sam dao sve od sebe.
Pokažimo PDF u tamnoj verziji — raspored je jasan, fontovi su elegantni, pogodnije je za noćno čitanje, noću je ugodnije za oči.

Osnove
1.Zašto koristiti redove poruka?
Mislim da se osnovna vrednost redova poruka ogleda u četiri aspekta. Prva je dekupljaža, i to je najvažije.

Na primer, u Pametni RAG projektu, nakon otpremanja fajla sledi mnogo posla — ekstrakcija metapodataka, pravljenje punoeksternog indeksa, AI vektorizacija.
Ovi procesi obiluju velikim količinama podataka i sami su veoma resursno intenzivni. Bez reda poruka, servis za otpremanje fajlova morao bi da čeka da se svi zadaci završe pre nego što vrati rezultat korisniku — iskustvo bi bilo loše.

Zato uveđemo Kafaku kao red poruka; servis za otpremanje fajlova samo šalje zadatak u red i završava. Ostali servisi sami konzumiraju ovu poruku i nezavisno obrađuju svoju poslovnu logiku. Čak i ako neki servis padne, to ne utiče na osnovni tok otpremanja fajlova, a tolerancija sistema na greške se znatno povećava.
Još jedan primer: u PmHub se proces odobrenja zadataka oslanja na RocketMQ radi dekupljaže.

Druga je asinhrona obrada — sistem može da stavi dugme, resursno intenzivne zadatke u red poruka za asinhronu obradu, čime brzo odgovara na korisničke zahteve. Na primer, nakon što korisnik naruči, sistem prvo može da vrati poruku o uspešnoj narudžbini, a potom narudžbinu stavi u red; pozadinski sistem zatim obrađuje narudžbinu.

Treća je "rezanje vrhova i punjenje dolina" — u scenama velike konkurentnosti ovo je izuzetno važno. Na primer, akcija "sekil" može momentalno da donese stotine hiljada zahteva. Ako oimo izravno na bazu, sistem sigurno će pasti. Ali preko reda poruka svi prvo ulaze u red, a pozadinski konzumeri, u skladu sa svojim kapacitetom, jedan po jedan konzumiraju. Čak i privremeno ne mogu da obrade, poruka se sigurno čuva u redu. Tako sistem neće pasti od iznenadnog saobraćaja.

Pored toga, redovi poruka podržavaju perzistentno skladištenje, pokušaje ponovnog dostavljanja i transakcione mehanizme. Čak i ako konzumer pri obradi poruke dođe do izuzetka, poruka se neće izgubiti — može se ponovo isporučiti i na kraju garantovati da će poslovna logika biti izvršena.
Kako RocketMQ koristiti za "rezanje vrhova i punjenje dolina"?
Moje razumevanje je: korisnikovi zahtevi ne odlaze direktno na pozadinski servis, već prvo u RocketMQ red. RocketMQ, kao srednji sloj velike propusnosti, može brzo da primi te zahteve. Zatim konzumeri, u skladu sa svojim kapacitetom, iz reda izvlače poruke određenom brzinom i obrađuju ih. Tako se formira bafer koji upija iznenadni saobraćaj.
Uzmimo scenario "sekil"-a. Prvo, korisnikov zahtev za "sekil" ne odlaze direktno na umanjenje zalihe, već se šalje poruka u RocketMQ. Ova operacija je brza jer se poruka samo baca u red, ne uključuje nikakvu poslovnu logiku. Zatim na konzumerskoj strani pokrećemo niti konzumera; ovi konzumeri relativno stabilnom brzinom konzumiraju poruku po poruku i obrađuju "sekil" logiku — provera zalihe, umanjenje, kreiranje narudžbine itd.
// proizvođač — prijem sekil zahteva
@PostMapping("/seckill")
public Result seckill(Long productId, Long userId) {
// direktno šaljemo poruku u RocketMQ, brzo vraćamo
Message message = new Message("seckill_topic",
JSON.toJSONString(new SeckillRequest(productId, userId)).getBytes());
try {
SendResult sendResult = rocketMQTemplate.syncSend("seckill_topic", message);
return Result.success("Sekil zahtev je podnet, malo strpljenja");
} catch (Exception e) {
return Result.fail("Sistem je zauzet, pokušajte kasnije");
}
}
// konzumer — konzumira po svom kapacitetu
@RocketMQMessageListener(topic = "seckill_topic",
consumerGroup = "seckill_consumer_group")
public class SeckillConsumer implements RocketMQListener<SeckillRequest> {
@Override
public void onMessage(SeckillRequest request) {
// ovde relativno stabilnom brzinom obrađujemo sekil logiku
// konzumer može koliko može, neće ga sresti iznenadni saobraćaj
seckillService.processSeckill(request.getProductId(), request.getUserId());
}
}Ovde postoji jedna stvar na koju treba obratiti pažnju — problem nakupljanja poruka. Ako konzumeri stalno zaostaju za proizvođačem, poruke će se beskonačno nakupljati, što na kraju može da dovede do punog diska ili isteka i brisanja poruka. Zato u realnim projektima moramo pratiti nakupljanje u redovima i po potrebi povećavati broj konzumera ili optimisati konzumersku logiku kako bismo ubrzali obradu.
- Vodič za intervju iz Jave sadrži originalno pitanje kandidata broj 7 iz ByteDancea, prvo razgovor za praksu u Javi: Da li ste čuli za MQ?
- Vodič za intervju iz Jave sadrži originalno pitanje kandidata broj 24 iz Tencent-a: Kako red poruka koristiti za "rezanje vrhova i punjenje dolina"?
- Vodič za intervju iz Jave sadrži originalno pitanje kandidata broj 4 iz Meituana, prvo razgovor: U projektu RocketMQ za "rezanje vrhova", koje su još scene pogodne za redove poruka?
- Vodič za intervju iz Jave sadrži originalno pitanje kandidata broj 20 iz ByteDancea, test razvoj, prvo razgovor: Šta RocketMQ radi, šta ti obično radiš njime?
memo: 03. 11. 2025. izmenjeno do ovde. Danas član planete je javio da je Sangfor objavio ponude, AI softver može da dođe do 30k+, zaista je veoma visoko, nadmašilo je SSP velikih internet kompanija. AI softverski pozivi će u narednih nekoliko godina biti veoma traženi, svi mogu da obrate pažnju, član planete koristi PmHub+RAG projekat.

2.Zašto izabrati RocketMQ?
Prva, podrška za transakcije — RocketMQ je odlično odradio. Na primer, u scenariju prenosa, moramo garantovati da "naplata" (lokalna transakcija) i "slanje poruke o prenosu" budu ili obe uspešne ili obe neuspešne, ne sme doći do situacije da je novac skinut, a poruka nije poslata. Transakcioni mehanizam RocketMQ može lepo da reši ovaj problem.
// RocketMQ transakcioni primer
TransactionSendResult sendResult = rocketMQTemplate.executeAndReplyTransaction(
"transfer_topic",
new Message("transfer_topic",
JSON.toJSONString(new TransferRequest(fromId, toId, amount)).getBytes()),
new RocketMQLocalTransactionListener() {
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
// izvršavamo lokalnu transakciju — naplata
accountService.deductAccount(fromId, amount);
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
return RocketMQLocalTransactionState.ROLLBACK;
}
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
// logika provere transakcije
return accountService.isDeducted(fromId, amount) ?
RocketMQLocalTransactionState.COMMIT :
RocketMQLocalTransactionState.ROLLBACK;
}
}
);Druga, podrška za sekvencijalne poruke. U mnogim scenarijima je redosled poruka važan. Na primer, životni ciklus narudžbine treba da bude: narudžbina → plaćanje → slanje → potvrda prijema. Ako se poremeća redosled, može doći do logičke konfuzije — još nije plaćeno, a već je poslato.
RocketMQ sekvencijalne poruke mogu da garantuju da poruke istog OrderID-a odlaze u isti red, a zatim ih isti konzumer konzumira uzastopno.
PmHub proces odobrenja zadataka koristi RocketMQ kako bi garantovao ispravan redosled koraka odobrenja.

Treća, RocketMQ podržava Master-Slave režim visoke dostupnosti. Kada Master čvor padne, Slave se automatski može prebaciti u Master, čime se poboljšava tolerancija na greške.

Ako je reč o prikupljanju logova i strujnoj obradi, Kafka je pogodnija jer je rođena za velike podatke. Pametni RAG projekat nakon otpremanja fajlova za vektorizaciju i izgradnju indeksa koristi Kafaku kao red poruka.

Ako treba lagan transfer poruka, RabbitMQ je bolji jer implementira AMQP protokol, podržava bogate rute i tipove razmenjivača.
Tehnički školski projekat za asinhronu obradu lajkova, omiljenih, komentara itd. koristi RabbitMQ.

3.Koje su prednosti i mane RocketMQ-a?
Prva, podrška za transakcione poruke — ovo je najveće zlatni RocketMQ-a. Moje razumevanje je da RocketMQ dvofaznom potvrdom garantuje atomičnost slanja poruka i lokalnih transakcija. Ovo je izuzetno važno za scene gde treba garantovati konzistentnost podataka. Na primer, umanjenje zaliha i kreiranje narudžbine — ili obe uspešne ili obe vraćene, ne sme biti nekonzistentno stanje.
Druga, podrška za sekvencijalne poruke. Za isti poslovni subjekat (npr. ista narudžbina) RocketMQ može da garantuje poredak poruka. Ovaj dizajn je veoma pametan: preko ShardingKey-a poruku rutira u isti red, a zatim ista konzumer nit uzastopno konzumira. To rešava mnoge realne poslovne potrebe.
Treća, velika propusnost i mali kašnjenje. Pojedinačni RocketMQ može da doseže stotine hiljada TPS.
Mana: RocketMQ ne eliminiše duplikate automatski, konzumer mora sam da implementira idempotentnost, inače može doći do ponovljenog konzumiranja.
// konzumer mora sam da implementira idempotentnost
@RocketMQMessageListener(topic = "order_topic",
consumerGroup = "order_consumer_group")
public class OrderConsumer implements RocketMQListener<OrderMessage> {
@Override
public void onMessage(OrderMessage message) {
// treba proveriti da li je jedinstveni identifikator poruke već obrađen
if (orderService.isProcessed(message.getOrderId())) {
// već obrađeno, samo se vraćamo
return;
}
// obrada narudžbine
orderService.processOrder(message);
}
}Šta mislite o RocketMQ-u?

RocketMQ je Alibaba-ov open-source distribuirani srednji sloj za poruke, karakterisan velikom propusnošću, malim kašnjenjem i visokom dostupnošću. Glavne komponente uključuju proizvođača, konzumera, Broker, Topic i redove. Proizvođač šalje poruke na Broker, koje se zatim, na osnovu rutiranih pravila, skladište u redovima; konzumeri izvlače poruke iz redova i obrađuju ih. Pogodno je za asinhronu dekupljažu i "rezanje vrhova".
- Vodič za intervju iz Jave sadrži originalno pitanje kandidata broj 4 iz JD.com, obla praksi, intervju: Recite Vaše mišljenje o RocketMQ-u
memo: 04. 11. 2025. izmenjeno do ovde. Danas član planete baš intervjuisao u kompaniji gde su gosti planete, tri projekta, mydb+Tehnička škola+Pametni, sada je i ovo sigurno.

4.Koji modeli poruka postoje u redovima poruka?
Mislim da se modeli poruka redova poruka mogu podeliti u dve velike kategorije: model tačka-u-tačku i model izdavač-pretplatnik.
Model tačka-u-tačku karakteriše to što jedna poruka može da konzumira samo jedan konzumer. Proizvođač šalje poruku u jedan red, a konzumer izvlači iz tog reda. Kada poruku konzumuje jedan konzumer, ona se briše, ostali konzumeri je ne vide.

Model izdavač-pretplatnik karakteriše to što jedna poruka može da konzumira više pretplatnika. Proizvođač objavljuje poruku na temu (Topic), svi konzumeri koji su pretplaćeni na tu temu dobijaju tu poruku.

Ovaj model je naročito pogodan za obaveštenja o događajima. Na primer, u Tehničkoj školi, kada autor objavi sadržaj, sistem istovremeno može da obavesti sve korisnike koji prate tog autora, dobijajući obaveštajnu notifikaciju. Sistemska obaveštenja su slična.

5.A šta model poruka kod RocketMQ-a?
RocketMQ koristi objedinjeni model baziran na Topic i Group. Unutar iste konzumer grupe može se smatrati kao tačka-u-tačku, između različitih konzumer grupa kao izdavač-pretplatnik.

U RocketMQ-u, tema (Topic) je logička klasifikacija poruka. Proizvođač šalje poruku na Topic, a konzumer izvlači sa Topic-a. Jedan Topic može imati više proizvođača koji mu šalju poruke, i više konzumera koji iz njega konzumiraju.
Jedan Topic je fizički podeljen u više redova (Queue). Prilikom slanja poruka, na osnovu nekog ključa poruka se rutira u različite Queue-ove. Ovaj dizajn je pametan jer i garantuje poredak unutar pojedinog Queue-a, i kroz više Queue-ova omogućuje paralelnu obradu.

Kod konzumera, konzumeri pripadaju određenoj konzumer grupi (Consumer Group). Više konzumera unutar konzumer grupe zajednički konzumira poruke iste teme. RocketMQ dodeljuje queue-e iz temene konzumerima u toj grupi.
Scena 1: jedan konzumer, jedan Topic sa 4 Queue
ConsumerGroup: order_consumer_group
└─ Consumer 1
├─ Queue 0
├─ Queue 1
├─ Queue 2
└─ Queue 3
Jedan konzumer konzumira svih 4 reda.Scena 2: dva konzumera, jedan Topic sa 4 Queue
ConsumerGroup: order_consumer_group
├─ Consumer 1
│ ├─ Queue 0
│ └─ Queue 2
│
└─ Consumer 2
├─ Queue 1
└─ Queue 3
Dva konzumera konzumiraju po 2 reda, ostvaruje se balansiranje opterećenja.Scena 3: četiri konzumera, jedan Topic sa 4 Queue
ConsumerGroup: order_consumer_group
├─ Consumer 1 → Queue 0
├─ Consumer 2 → Queue 1
├─ Consumer 3 → Queue 2
└─ Consumer 4 → Queue 3
Četiri konzumera konzumiraju po jedan red, puna paralelnost.memo: 15. 11. 2025. izmenjeno do ovde. Danas član planete javio je dobru vest — Trip.com je otvorio SP, vrlo zadovoljan, zahvaljujući projektima sa planete, koristio je Pametni RAG+mydb točak.

6.Konsumni modeli poruka su poznati?
Mislim da se konzumni modeli mogu kategorisati u dve dimenzije: smer konzumcije i opseg konzumcije.
Po smeru, postoje dva modela: pull (povuci) i push (guraj).

Pull model zahteva da konzumer aktivno izvlači poruke iz reda, može da kontroliše brzinu i količinu povlačenja, ali mora stalno da anketa, što je više troši resurse.
Push model je kada red sam gura poruke konzumeru; konzumer samo registrujeListenera, čim poruka stigne, okida se callback i obrada, brz odgovor, ali može doći do nakupljanja poruka.
Po opsegu, takođe postoje dva modela: klaster konzumpcija i difuzna konzumpcija.

Klaster konzumpcija znači da više konzumera u istoj konzumer grupi zajednički konzumira poruke jedne teme. Poruke se distribuiraju konzumerima u toj grupi, svaka poruka se konzumira samo jednom od strane jednog konzumenta iz te grupe.
Drugim rečima, RocketMQ ravnomerno dodeljuje sve queue-e iz temene konzumerima u grupi, ostvarujući balansiranje opterećenja. Takođe se garantuje da istu poruku konzumira samo jedan konzumer iz grupe, izbegava se dupla konzumpcija.
Topic: order_topic
├─ Queue 0 → [poruka1] [poruka3] [poruka5]
├─ Queue 1 → [poruka2] [poruka4] [poruka6]
├─ Queue 2 → [poruka7] [poruka9] [poruka11]
└─ Queue 3 → [poruka8] [poruka10] [poruka12]
ConsumerGroup: order_consumer_group
├─ Consumer 1 konzumira Queue 0, 1
├─ Consumer 2 konzumira Queue 2, 3
Ista poruku konzumira samo Consumer 1 ili Consumer 2, ne će se duplo konzumirati.Difuzna konzumpcija znači da svaku poruku iz temene svaki konzumer u grupi konzumira jednom. Drugim rečima, svaki konzumer u grupi dobija sve poruke iz temene, ostvarujući efekat difuzije.
Topic: config_update_topic
└─ [konfiguraciona poruka1] [konfiguraciona poruka2] [konfiguraciona poruka3]
ConsumerGroup: config_consumer_group
├─ Consumer 1 (aplikacija na serveru 1)
│ └─ Primio: poruku1, poruku2, poruku3 (potpuno)
│
├─ Consumer 2 (aplikacija na serveru 2)
│ └─ Primio: poruku1, poruku2, poruku3 (potpuno)
│
└─ Consumer 3 (aplikacija na serveru 3)
└─ Primio: poruku1, poruku2, poruku3 (potpuno)
Sva trojica konzumera su primila sve poruke, svaki nezavisno obrađuje.7.Osnovna arhitektura RocketMQ-a?
Arhitekturu RocketMQ-a četiri osnovna dela čine: NameServer, Broker, proizvođač i konzumer.

Moje razumevanje je da je NameServer centar za rutiranje RocketMQ-a, odgovaran za održavanje informacija o rutiranju između Topic-a i Broker-a. Pre slanja poruka, proizvođač i pre konzumiranja konzumer prvo dobijaju najnovije rutirane informacije sa NameServer-a.
- Svaki Broker se registrova na NameServer, uključujući adresu, port, Topic i Queue koje skladišti.
- NameServer na osnovu naziva Teme proizvođaču i konzumeru kaže odgovarajuću adresu Broker-a.
- Broker periodicno šalje heartbeate na NameServer, izveštavajući o svom stanju.

Broker je centar za skladištenje poruka, odgovornosti:
- Sve poruke koje proizvođači šalju se skladište na Broker-u, u obliku fajlova perzistentno na disku.
- Konzumeri pri konzumiranju izvlače poruke sa Broker-a; Broker mora na osnovu Offset-a konzumera da pronađe odgovarajuću poruku i vrati je konzumeru.
- Ako je konfigurisana visoka dostupnost, Broker Master sinhronizuje poruke na Broker Slave, ostvarujući master-slave backup.
Proizvođač pri slanju poruka prvo dobija rutirane informacije o Topic-u sa NameServer-a, a zatim na osnovu njih šalje poruku na odgovarajući Broker.
Konzumer pri konzumiranju takođe prvo dobija rutirane informacije o Topic-u sa NameServer-a, a zatim na osnovu njih izvlači poruke sa odgovarajućeg Broker-a.
memo: 25. 11. 2025. izmenjeno do ovode. Danas član planete je javio da je dobio ponude iz iFlyTEK-a i Huawei-a, sve je objavljeno, jesenji regrutacija je zvanično završena. Posebno se zahvalila Tehničkoj školi što joj je pronašla letnju praksu, a za jesenji regrutaciju je zahvaljujući Pametnom RAG-u i praktičnim prijateljima sa planete sigurno prošla. Takođe sam joj izmenio životopis i dodao mnogo poena, na licu mesta joj je hiljadu puta hvale.

8.Možeš detaljno da predstaviš RocketMQ NameServer?
NameServer je centar za rutiranje i otkrivanje servisa. Prva dužnost mu je da skladišti i održava rutirane informacije. Kada Broker startuje, registrova se na NameServer.

NameServer te informacije čuva u memoriji, formirajući rutiranu tablicu.

Druga dužnost je da pruža servis rutiranja. Kada proizvođač ili konzumer želi da zna na kom Broker-u se tema nalazi, upita NameServer. NameServer na osnovu naziva temene vraća adresu Broker-a i informacije o redu.
Treća dužnost je monitoring stanja Broker-a. Broker periodicno šalje heartbeate na NameServer; ako neki Broker dugo ne šalje heartbeat, NameServer ga označava kao nedostupan i uklanja ga iz rutirane tablice.
Šta Broker radi?
Broker je server za skladištenje poruka, odgovoran je da primi poruke od proizvođača i da ih skladišti, a zatim ih vraća konzumerima pri konzumiranju.

Šta proizvođač?
Osnovna dužnost proizvođača je da pretvori podatke aplikacije u poruke i da ih pošalje na Broker.

RocketMQ nudi tri načina slanja poruka: sinhrono, asinhrono i jednostrano.
- Sinhrono slanje: proizvođač nakon slanja poruke blokira i čeka odgovor od Broker-a. Tek kada primi potvrdu da je Broker uspešno sačuvao poruku, vraća se aplikaciji.
- Asinhrono slanje: proizvođač se vraća odmah nakon slanja, ne blokira. Odgovor od Broker-a stiže aplikaciji kroz callback funkciju.
- Jednosmerno slanje: proizvođač se vraća odmah nakon slanja, ne čeka odgovor, niti zahteva callback. Ovaj mod se koristi u scenama gde ne mari rezultatom slanja.
Šta konzumer?
Konzumer je primalac poruka, osnovna dužnost mu je izvlačiti poruke sa Broker-a, obrađivati poslovnu logiku i potvrđivati konzumiranu poziciju.

RocketMQ istovremeno podržava Pull, Push, Pop tri konzumna modela.
Pull model je najosnovniji način konzumiranja. Konzumer aktivno inicira zahtev ka Broker-u i izvlači poruke. Push model izgleda kao da server gura poruke, ali se u temelju ipak oslanja na Pull model.
Kada je konzumera mnogo, rebalansiranje konzumcije može potrajati, pa RocketMQ nudi Pop model. Pop model potpuno prebacuje rebalansiranje na serversku stranu, smanjujući opterećenje konzumera.
memo: 27. 11. 2025. izmenjeno do ovde. Danas član planete je javio da je dobio ponudu iz ByteDance-a, SSP ponuda, još 80.000 potpisa za bonus, posebno se zahvalio projektima sa planete i temama za učenje Priprema za intervju.

Napredno
9.Kako garantovati dostupnost/pouzdanost/ne-izgubljenost poruka?
Poruke se mogu izgubiti u tri faze: proizvodnja, skladištenje, konzumpcija.
Dakle, treba razmotriti ove tri faze:

Proizvodnja
U fazi proizvodnje, kroz mehanizam potvrđivanja zahteva garantuje se pouzdan isporuka poruka.
- Pri sinhronom slanju obratiti pažnju na rezultate odgovora i izuzetke. Ako se vrati OK, poruka je uspešno stigla na Broker; ako se vrati greška ili dođe do drugog izuzetka, treba pokušati ponovo.
- Pri asinhronom slanju, treba proveriti u callback metodi; ako slanje ne uspe ili dođe do izuzetka, treba pokušati ponovo.
- Ako dođe do isteka vremena, može se preko API-ja za proveru loga utvrditi da li je uspešno sačuvano na Broker-u.
Skladištenje
U fazi skladištenja, kroz konfiguraciju parametara za pouzdanost Broker-a izbegava se gubitak poruka usled pada. Jednostavno rečeno, u scenama pouzdanosti uvek treba koristiti sinhroni način.
- Dok se poruka ne perzistira u CommitLog (log fajl), čak i ako Broker padne, nekonzumirane poruke se mogu oporaviti i ponovo konzumirati.
- Broker-ov mehanizam čišćenja diska: sinhrono i asinhrono čišćenje; obe vrste garantuju da je poruka sigurno u pagecache-u (u memoriji), ali je sinhrono čišćenje pouzdanije — proizvođač šalje poruku, čeka da se podaci perzistentiraju na disku, pa tek tada vraća odgovor.
- Broker kroz master-slav režim garantuje visoku dostupnost; Broker podržava sinhronu i asinhronu replikaciju Master-Slave, proizvođačeve poruke svih slanje ka Master-u, a konzumpcija može i sa Master-a i sa Slave-a. Sinhrona replikacija garantuje da čak i ako Master padne, poruka sigurno ima backup na Slave-u, neće se izgubiti.
Konzumpcija
S gledišta Consumer-a, kako garantovati da je poruka uspešno konzumirana?
- Ključno je vreme potvrde — ne šaljite potvrdu odmah po prijemu poruke, već nakon što se izvrši cela poslovna logika konzumcije. Jer red poruka održava poziciju konzumcije; ako se logika ne izvrši i ne potvrdi, sledeći izvlačenje iz reda donosi istu poruku.
10.Kako rešavati problem dupliranih poruka?
RocketMQ može garantovati isporuku, ne-izgubljenost, ali ne garantuje neduplikatnu konzumpciju.
Zato, na poslovnom sloju treba osigurati idempotentnost ili deduplikaciju.

Idempotentnost znači da operacija može više puta da se izvrši a neće proizvesti nuspojave, odnosno bez obzira koliko puta se izvrši, rezultat je isti. U poslovnu logiku može se dodati logika provere, osiguravajući da višestruka konzumcija iste poruke neće proizvesti nuspojave.
Na primer, u scenariju plaćanja, konzumer konzumira poruku o naplati, izvršava naplatu od 100 RMB za jednu narudžbinu.
Ako se usled nestabilnosti mreže itd. poruka o naplati ponovo isporuči, konzument ponovo konzumira tu poruku, ali konačni poslovni rezultat mora garantovati samo jednu naplatu od 100 RMB. Ako je operacija naplate ispravna, smatra se da je ceo proces konzumcije idempotentan.
Deduplikacija poruka znači da pre konzumiranja konzumer prvo proveri da li je već konzumirao tu poruku; ako jeste, ne konzumira je.
Poslovni sloj može imati specijalnu tablicu koja čuva ID-ove već konzumiranih poruka; pre konzumiranja proveri tu tablicu, ako ID već postoji, ne konzumira.
public void processMessage(String messageId, String message) {
if (!isMessageProcessed(messageId)) {
// obrada poruke
markMessageAsProcessed(messageId);
}
}
private boolean isMessageProcessed(String messageId) {
// proveri tabelu deduplikacije, da li postoji messageId
}
private void markMessageAsProcessed(String messageId) {
// ubaci messageId u tabelu deduplikacije
}Kako garantovati idempotentnost poruka?

Prva, poruka mora nositi jedinstven poslovni identifikator; može se generisati kroz Snowflake algoritam globalni jedinstveni ID.
Message msg = new Message(TOPIC /* Topic */,
TAG /* Tag */,
("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET) /* Message body */
);
message.setKey("ORDERID_100"); // šifra narudžbine
SendResult sendResult = producer.send(message);Druga, kada konzumer primi poruku, proveri u Redis-u postojanje oznake za taj poslovni ključ; ako postoji, smatra se da je konzumcija uspešna, inače izvrši poslovnu logiku, a nakon završetka doda oznaku u keš.
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
try {
for (MessageExt messageExt : msgs) {
String bizKey = messageExt.getKeys(); // jedinstveni poslovni ključ
//1. provera postojanja oznake
if(redisTemplate.hasKey(RedisKeyConstants.WAITING_SEND_LOCK + bizKey)) {
continue;
}
//2. izvršavanje poslovne logike
//TODO do business
//3. postavljanje oznake
redisTemplate.opsForValue().set(RedisKeyConstants.WAITING_SEND_LOCK + bizKey, "1", 72, TimeUnit.HOURS);
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
logger.error("consumeMessage error: ", e);
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}Treća, iskoristiti jedinstveni indeks baze podataka da bi se sprečila duplicirana unos.
CREATE TABLE `t_order` (
`id` bigint(20) NOT NULL AUTO_INCREMENT,
`order_id` varchar(64) NOT NULL COMMENT 'šifra narudžbine',
`order_name` varchar(64) NOT NULL COMMENT 'naziv narudžbine',
PRIMARY KEY (`id`),
UNIQUE KEY `order_id` (`order_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='tabela narudžbina';Konačno, koristiti verzije u tabeli baze podataka, kroz mehanizam optimističkog zaključavanja garantovati idempotentnost. Pri svakoj operaciji ažuriranja proveriti da li je broj verzije konzistentan, samo ako jeste izvršiti ažuriranje i povećati broj verzije. Ako se broj verzije ne poklapa, znači da je operacija već izvršena, odbiti dupliranu operaciju.
public void updateRecordWithOptimisticLock(int id, String newValue, int expectedVersion) {
int updatedRows = jdbcTemplate.update(
"UPDATE records SET value = ?, version = version + 1 WHERE id = ? AND version = ?",
newValue, id, expectedVersion
);
if (updatedRows == 0) {
throw new OptimisticLockingFailureException("Record has been modified by another transaction");
}
}Ili mehanizam pesimističkog zaključavanja, kroz mehanizam zaključavanja baze podataka garantovati idempotentnost.
public void updateRecordWithPessimisticLock(int id) {
jdbcTemplate.queryForObject("SELECT * FROM records WHERE id = ? FOR UPDATE", id);
jdbcTemplate.update("UPDATE records SET value = ? WHERE id = ?", "newValue", id);
}Snowflake algoritam je poznat?
Snowflake algoritam je Twitter-ov razvijeni algoritam za generisanje distribuiranih jedinstvenih ID-ova.

Snowflake algoritam čuva ID u 64 bita, sastoji se od 4 dela:
- Najviši bit zauzima 1 bit, uvek je 0, označava pozitivan broj.
- Srednji deo zauzima 41 bit, vrednost je vremenska oznaka na nivou milisekunde;
- Donji srednji deo zauzima 10 bitova, ID mašine (uključuje ID data centra i ID mašine), može podržati 1024 čvora.
- Najniži deo zauzima 12 bitova, vrednost je različita inkrementalna sekvenca generisana u okviru tekuće milisekunde, gornja granica je 4096;
Trenutno postoji mnogo implementacija Snowflake algoritma, može se direktno koristiti IdUtil.getSnowflake() metoda iz Hutool biblioteke alata za dobijanje Snowflake ID-a.
long id = IdUtil.getSnowflakeNextId();
- Vodič za intervju iz Jave sadrži originalno pitanje kandidata broj 4 iz JD.com, obla praksi, intervju: Kako rešavati problem duplirane konzumcije? Kako garantovati idempotentnost? Snowflake algoritam je poznat?
11.Kako rešiti nakupljanje poruka?
Došlo je do nakupljanja poruka, treba hitro konzumirati nakupljene poruke, poboljšati kapacitet konzumcije. Uopšteno postoje dva načina:

- Proširenje konzumera: Ako broj Message Queue u trenutnoj Topic je veći od broja konzumera, može se proširiti broj konzumera, povećati kapacitet konzumcije i hitro konzumirati nakupljene poruke.
- Migracija poruka, proširenje Queue: Ako je broj Message Queue u trenutnoj Topic manji ili jednak broju konzumera, proširenje konzumera nema efekta — treba razmotriti proširenje Message Queue. Može se kreirati privremena Topic sa više Message Queue, zatim nekoliko konzumera prebacuje podatke u privremenu Topic (jer ne obuhvata poslovnu logiku, samo preusmerava poruke, brzo je). Zatim se koristi prošireni konzumeri da konzumiraju podatke iz nove Topic; nakon konzumcije, vraća se u prvobitno stanje.

12.Kako se realizuju sekvencijalne poruke?
RocketMQ nudi dva nivoa sekvencijalnih poruka: globalne sekvencijalne i lokalne sekvencijalne.

Globalna sekvenca znači da sve poruke u Topic konzumira striktno po redosledu slanja; ovaj mod je manje efikasan, u praksi se retko koristi.

Lokalna sekvenca znači da se garantuje redosled unutar određene particije — ovo je naš česti način.

Ključ za očuvanje redosleda je poslati poruke koje trebaju ostati uređene u isti MessageQueue.
// na osnovu ID-ja narudžbine biramo red, garantujući da poruke iste narudžbine idu u isti red
producer.send(message, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
String orderId = (String) arg;
int index = orderId.hashCode() % mqs.size();
return mqs.get(index);
}
}, orderId);Svaki MessageQueue na Broker-u odgovara jednom ConsumeQueue, poruke se upisuju uzastopno po redosledu dolaska na Broker.
Kada konzumer počne da konzumira određeni MessageQueue, na strani Broker-a se zaključava taj red, drugi konzumeri ne mogu istovremeno da konzumiraju taj red. Tako se osigurava da u istom trenutku samo jedan konzumer obrađuje poruke tog reda, garantujući poredak konzumcije.
- Vodič za intervju iz Jave sadrži originalno pitanje kandidata broj 2 iz JD.com, pozadinski razvoj, intervju: Recite mi princip MQ, kako garantovati redosled prijema poruka?
- Vodič za intervju iz Jave sadrži originalno pitanje kandidata broj 1 iz Shouqianba, Java pozadinski prvi razgovor: Sekvencijalne poruke kod RocketMQ-a?
- Vodič za intervju iz Jave sadrži originalno pitanje kandidata broj 6 iz srednje fabrike, Guangzhou, intervju: Kako RocketMQ garantuje poredak poruka?
memo: 15. 08. 2025. izmenjeno do ovde. Danas, dok sam pomagao članu planete da izmeni životopis, dobila sam ovu povratnu informaciju: trenutno je na letnjoj praksi u Gaode, krajem marta sam se obratio Ergeu da izmeni životopis, smatram da je vrlo dobro izmenjen.

13.Kako realizovati filtriranje poruka?
Dva rešenja:
- Jedno je filtriranje na strani Broker-a prema deduplikacionoj logici konzumera; prednost je što se izbegava nepotrebni transfer konzumera, mana je što se povećava opterećenje Broker-a, implementacija je relativno složena.
- Drugo je filtriranje na strani Consumer-a, na osnovu tag-a postavljenom na poruci; prednost je jednostavna implementacija, mana je što veliki količine nepotrebnih poruka stižu na Consumer stranu i samo se odbacuju.
Uobičajeno se koristi filtriranje na Consumer strani; ako se želi povećati propusnost, može se koristiti Broker filtriranje.
Postoje tri načina filtriranja poruka:

- Filtriranje prema Tag: ovo je najčešći, jednostavan i efikasan
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("CID_EXAMPLE");
consumer.subscribe("TOPIC", "TAGA || TAGB || TAGC");- Filtriranje SQL izrazom: fleksibilnije
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("please_rename_unique_group_name_4");
// samo ako poruka ima atribut a, a >= 0 and a <= 3
consumer.subscribe("TopicTest", MessageSelector.bySql("a between 0 and 3");
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
consumer.start();- Filter Server način: najfleksibilniji, najkompleksniji, omogućuje korisniku da definiše funkciju za filtriranje
14.Odgodivene poruka je poznata?
U e-trgovini, automatsko otkazivanje narudžbine posle isteka vremena je tipičan primer korišćenja odgodivih poruka. Korisnik podnosi narudžbinu, sistem može poslati odgodivu poruku, posle 1h proveriti status narudžbine, ako je i dalje neplaćena, otkazati i osloboditi zalihe.
RocketMQ podržava odgodivene poruke; pri proizvodnji poruke treba postaviti nivo odgodivanja:
// instanciramo proizvođača za odgodivu poruku
DefaultMQProducer producer = new DefaultMQProducer("ExampleProducerGroup");
// pokrećemo proizvođača
producer.start();
int totalMessagesToSend = 100;
for (int i = 0; i < totalMessagesToSend; i++) {
Message message = new Message("TestTopic", ("Hello scheduled message " + i).getBytes());
// postavljamo nivo 3 odgodivanja, poruka će biti poslata posle 10s (trenutno su podržani samo fiksni vremenski intervali, detaljnije u delayTimeLevel)
message.setDelayTimeLevel(3);
// šaljemo poruku
producer.send(message);
}Ali trenutno RocketMQ podržava ograničene nivoe odgodivanja:
private String messageDelayLevel = "1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h";Kako RocketMQ implementira odgodivu poruku?
Jednostavno, osam reči: privremeno skladištenje+periodični zadatak.
Broker primi odgodivu poruku, prvo je šalje u odgovarajući MessageQueue teme (SCHEDULE_TOPIC_XXXX) u odgovarajućem vremenskom intervalu, a zatim kroz periodični zadatak proleći te redove, nakon isteka, isporučuje poruku u redove ciljne Topic-e, pa konzumeri mogu normalno da konzumiraju te poruke.

15.Kako realizovati distribuirane transakcije poruka? Polu-poruka?
Polu-poruka: poruka koju konzumer privremeno ne može da konzumira. Proizvođač uspešno šalje na Broker, ali poruka se označava kao "privremeno nedostupna" stanje; tek kada Proizvođač izvrši lokalnu transakciju i nakon toga potvrdi drugi put, konzumer može da konzumira ovu poruku.
Koristeći polu-poruku, može se ostvariti distribuirana transakcija poruka, ključ je u drugoj potvrdi i povratnoj proveri:

- Proizvođač šalje polu-poruku broker-u
- Proizvođač prima odgovor, slanje je uspešno, u ovom trenutku poruka je polu-poruka, označena kao "nedostupna za isporuku", konzumer je ne može konzumirati.
- Proizvođač izvršava lokalnu transakciju.
- U normalnom slučaju lokalna transakcija se završi, Proizvođač šalje Broker-u Commit/Rollback; ako je Commit, Broker polu-poruku označava kao normalnu poruku, konzumer može da konzumira; ako je Rollback, Broker odbacuje poruku.
- U abnormalnom slučaju, Broker dugo čeka drugu potvrdu. Nakon određenog vremena, upitaće sve polu-poruke, a zatim će na strani Proizvođača proveriti status izvršenja polu-poruke.
- Proizvođač proverava stanje lokalne transakcije
- Na osnovu stanja transakcije šalje commit/rollback broker-u (5, 6, 7 su povratna provera).
- Konzumer konzumira poruku, izvršava lokalnu transakciju.
16.Red mrtvih poruka je poznat?
Red mrtvih poruka služi za čuvanje onih poruka koje se ne mogu normalno obraditi; takve poruke se zovu "mrtva slova" (Dead Letter).

Uzrok mrtvih slova je: konzumer pri obradi poruke doživi izuzetak i dostigne maksimalan broj ponovnih pokušaja. Kada se ukaže i reši uzrok neuspeha konzumcije, te mrtve poruke se mogu ponovo isporučiti, konzumeri ponovo konzumiraju; ako se privremeno ne može rešiti, kako bi se izbeglo brisanje mrtvih poruka posle isteka roka, prvo ih treba izvesti i sačuvati.
- Vodič za intervju iz Jave sadrži originalno pitanje kandidata broj 4 iz JD.com, obla praksi, intervju: Recite mi o redovima mrtvih poruka kod RocketMQ-a
17.Kako garantovati visoku dostupnost RocketMQ-a?
NameServer je bez stanja, ne komunicira međusobno, samo klaster deployment garantuje visoku dostupnost.

Visoka dostupnost RocketMQ-a se u prvom redu ogleda u čitanju i pisanju Broker-a, to se postiže kroz klaster i master-slave.

Broker može imati dve uloge: Master i Slave; Master podržava čitanje i pisanje, Slave samo čitanje; Master sinhronizuje poruke na Slave.
To znači da Proizvođač može samo na Master Broker da piše poruke, a Konzumer može sa i Master i Slave Broker-a da čita.
U konfiguraciji Konzumera nije potrebno podešavati da li se čita sa Master-a ili sa Slave-a; kada Master nije dostupan ili je zauzet, zahtev za čitanje se automatski prebacuje na Slave. Sa ovakvim automatskim prebacivanjem, kada Master mašina padne, Konzumer i dalje može sa Slave-a da čita poruke, ne utiče na čitanje, ostvarujući visoku dostupnost čitanja.
Kako ostvariti visoku dostupnost pisanja? Prilikom kreiranja Topic-e, više Message Queue iz te Topic-e se raspoređuje na više Broker grupa (isti naziv Broker-a, različit brokerId mašine čine Broker grupu); kada Master jedne Broker grupe padne, Master-i drugih grup su i dalje dostupni, Proizvođač i dalje može da šalje poruke. RocketMQ trenutno ne podržava automatsku konverziju Slave-a u Master; ako resursi nisu dovoljni, treba ručno zaustaviti Slave Broker, izmeniti konfiguracioni fajl i pokrenuti Broker sa novom konfiguracijom.
GitHub repozitorijum sa 17000+ zvezdica "Ergeov napredni put kroz Javu" prvo PDF izdanje je konačno stiglo! Uključuje osnovnu sintaksu Jave, nizove&stringove, OOP, kolekcijski framework, Java IO, obradu izuzetaka, nove karakteristike Jave, mrežno programiranje, NIO, konkurentno programiranje, JVM itd. — ukupno preko 320.000 reči, 500+ crteža, zaista lako razumljivo, zabavno... Detalji: Vrhunski, Java tutorijal na GitHub-u sa 17000+ zvezdica
Principi
18.Recite mi ukupan tok rada RocketMQ-a?
Jednostavno rečeno, RocketMQ je distribuirani red poruka, odnosno red poruka+distribuirani sistem.
Kao red poruka, to je model šalji-čuvaj-primaj, što odgovara Proizvođač, Broker, Konzumer; kao distribuirani sistem, mora imati serversku stranu, klijentsku stranu, centar za registraciju, što odgovara Broker, Proizvođač/Konzumer, NameServer
Dakle, pogledajmo osnovni tok: RocketMQ se sastoji od klastera NameServer centra za registraciju, klastera Proizvođača, klastera Konzumera i nekoliko Broker-a (RocketMQ procesa):
- Broker pri pokretanju ide na sve NameServer-e da se registruje i održava dugu vezu, svakih 30s šalje po jedan heartbeat
- Proizvođač pri slanju poruka sa NameServer-a dobija adresu Broker servera, na osnovu algoritma za balansiranje opterećenja bira jedan server i šalje poruku
- Konzumer pri konzumiranju sa NameServer-a takođe dobija adresu Broker-a, a zatim aktivno izvlači poruke za konzumiranje

19.Zašto RocketMQ ne koristi Zookeeper kao centar za registraciju?
Kao što znamo, Kafka koristi Zookeeper kao centar za registraciju — iako polako počinje da ga eliminiše. RocketMQ ne koristi Zookeeper iz uglavnom sledećih razloga:
- S gledišta dostupnosti, prema CAP teoriji, najviše dva tačka se mogu zadovoljiti istovremeno, a Zookeeper zadovoljava CP, to znači da Zookeeper ne garantuje dostupnost servisa. Tokom izbora, samo izbor traje predugo, u tom periodu ceo klaster nije dostupan, a za centar za registraciju to definitivno nije prihvatljivo. Kao servis za otkrivanje servisa treba dizajnirati za dostupnost.
- S gledišta performansi, NameServer je veoma lagan, može se horizontalno skalirati dodavanjem mašina, povećavajući otpornost klastera. Zapis Zookeeper-a nije skalabilan; rešenje je samo podeliti domene, podeliti na više Zookeeper klastera — prvo je suviše komplikovano za rad, a drugo, to narušava A u CAP-u, servisi međusobno nisu povezani.
- Problem mehanizma perzistencije: ZAB protokol ZooKeeper-a za svaki zahtev za upis čuva transakcioni log na svakom Zookeeper čvoru, plus periodično slikanje memorije na disk kako bi se garantovala konzistentnost i perzistentnost podataka. Za jednostavnu scenu otkrivanja servisa to je previše teško. Implementacija je preteška. Podaci koji se čuvaju trebalo bi da su visoko prilagođeni.
- Slanje poruka treba slabo da zavisi od centra za registraciju, a dizajn filozofija RocketMQ-a upravo je to: kada prvi put šalje poruku, Proizvođač dobija adresu Broker-a sa NameServer-a i kešira je lokalno; ako ceo NameServer klaster postane nedostupan, u kratkom roku to neće mnogo uticati na proizvođače i konzumere.
20.Kako Broker čuva podatke?
RocketMQ-ove glavne skladištne fajlove uključuju CommitLog, ConsumeQueue i Indexfile.

Ukupni dizajn skladištenja poruka:

- CommitLog: glavni skladišni subjekat za poruke i metapodatke, čuva sadržaj poruka koji je napisao Proizvođač. Sadržaj poruka nije fiksne dužine. Pojedinačni fajl je standardno 1G, naziv fajla je 20 cifara, sleva se dopunjava nulama, ostatak je početni ofset, na primer 00000000000000000000 predstavlja prvi fajl, početni ofset 0, veličina fajla 1G=1073741824; kada se prvi fajl napuni, drugi je 00000000001073741824, početni ofset 1073741824, itd. Poruke se uglavnom upisuju u log fajl uzastopno; kada se fajl napuni, prelazi se na sledeći.
CommitLog fajl se čuva u ${Rocket_Home}/store/commitlog direktorijumu; sa slike se jasno vidi ofset u nazivu fajla, svaki fajl je podrazumevano 1G, nakon što se napuni automatski se generiše novi fajl.

- ConsumeQueue: red konzumacije poruka, uveden radi poboljšanja performansi konzumcije. Pošto je RocketMQ zasnovan na temi Topic, konzumpcija je tematska; ako bi se prolazilo kroz commitlog i tražile poruke po temi, vrlo je neefikasno.
Konzumer može na osnovu ConsumeQueue da pronađe poruke za konzumiranje. ConsumeQueue (logički red konzumcije) kao indeks konzumacije čuva početni fizički ofset u CommitLog-u za poruke iz određene Topic-e, veličinu poruke i HashCode vrednost Tag-a.
ConsumeQueue fajl se može smatrati CommitLog indeksnim fajlom na osnovu Topic-a, organizacija ConsumeQueue foldera je tri sloja: topic/queue/file, konkretan putanja je: $HOME/store/consumequeue/{topic}/{queueId}/{fileName}. Takođe ConsumeQueue fajl ima fiksnu dužinu, svaka stavka ima 20 bajtova: 8 bajtova CommitLog fizičkog ofseta, 4 bajta dužine poruke, 8 bajtova tag hashcode-a; jedan fajl ima 300.000 stavki, može se kao niz pristupiti svakoj stavci, svaki ConsumeQueue fajl je oko 5.72M;

- IndexFile: IndexFile (indeksni fajl) pruža način za pretragu poruka prema ključu ili vremenskom intervalu. Indeksni fajl se nalazi na: {fileName}, naziv fajla je vremenska oznaka vremena kreiranja, fiksna veličina pojedinačnog IndexFile fajla je oko 400M, jedan IndexFile može da sačuva 20.000.000 indeksa; donji sloj skladištenja IndexFile-a je realizacija HashMap strukture u fajl sistemu, pa je indeksni fajl RocketMQ-a u osnovi hash indeks.

Ukratko: RocketMQ koristi hibridnu strukturu skladištenja, odnosno svi redovi jednogBroker-a instance dele jedan log fajl (CommitLog) za skladištenje.
Hibridna struktura RocketMQ-a (sadržaji poruka više Topic-a čuvaju se u jednom CommitLog-u) za Proizvođača i Konzumera razdvaja podatke i indeksne delove. Proizvođač šalje poruku na Broker, Broker sinhrono ili asinhrono čisti disk i čuva u CommitLog.
Dok se poruka perzistira na disk (CommitLog), Proizvođačeva poruka se neće izgubiti. Zato Konzumer sigurno ima priliku da konzumira tu poruku. Kada ne može da izvuci, može sledeći put pokušati ponovo; istovremeno, server podržava dugo ankete (long polling), ako se jedna zahtev za izvlačenje ne uspeš izvući, Broker dozvoli čekanje do 30s; ako u tom periodu stigne nova poruka, odmah je vraća konzumeru.
Ovde RocketMQ koristi pozadinsku servisnu nit na Broker strani — ReputMessageService, koja neprestano raspoređuje zahteve i asinhrono gradi ConsumeQueue (logički red konzumcije) i IndexFile (indeksni fajl).

21.Kažite RocketMQ kako čita i piše fajlove?
RocketMQ verno koristi neke efikasne načine čitanja i pisanja fajlova operativnog sistema — PageCache, sekvencijalno čitanje i pisanje, nulti kopiranje.
- PageCache, sekvencijalno čitanje
U RocketMQ-u, ConsumeQueue logički red konzumcije čuva malo podataka i čita se sekvencijalno; uz mehanizam predčitavanja PageCache-a, brzina čitanja ConsumeQueue fajlova gotovo je jednaka čitanju iz memorije, čak i u slučaju nakupljanja poruka performance ne pada. Za CommitLog, fajlove za čuvanje logova poruka, čitanje sadržaja poruka generiše mnogo slučajnih pristupa, što jako utiče na performanse. Ako se odabari odgovarajući IO algoritam, na primer "Deadline" (ako Block storage koristi SSD), performanse slučajnog čitanja se takođe mogu poboljšati.
PageCache je OS keš fajlova, ubrzava čitanje i pisanje. Generalno, brzina sekvencijalnog čitanja i pisanja fajlova gotovo je jednaka brzini čitanja i pisanja memorije, uglavnom zato što OS koristi PageCache mehanizam da optimizuje IO operacije, deo memorije se koristi kao PageCache. Kod pisanja, OS prvo piše u Cache, a zatim asinhrono pdflush nit smešta podatke sa Cache-a na disk. Kod čitanja, ako jedan pristup fajlu ne pogodi PageCache, OS čita fajl sa diska i istovremeno predčitava susedne blokove.
- Nulti kopiranje
Pored toga, RocketMQ koristi MappedByteBuffer za čitanje i pisanje fajlova. To iskorišćava NIO FileChannel model da direktno mapira fizičke fajlove na disku u adresni prostor korisničkog procesa (ovaj mmap način smanjuje troškove tradicionalnog IO, gde se podaci fajla u kernel adresnom prostoru i korisničkom adresnom prostoru kopiraju unazad), pretvara operacije nad fajlovima u direktne operacije nad memorijom, drastično poboljšavajući efikasnost čitanja i pisanja (zato što je potrebna mehanizma mapiranja memorije, fajlovi skladištenja RocketMQ-a koriste fiksnu strukturu, olakšavajući mapiranje celog fajla u memoriju).
Šta je nulti kopiranje?
U operativnom sistemu, tradicionalnim načinom podaci prolaze kroz nekoliko kopiranja i nekoliko promena konteksta između korisničkog i kernel moda.

- Kopiranje podataka sa diska u memoriju kernel moda;
- Kopiranje iz kernel memorije u korisničku memoriju;
- Zatim kopiranje iz korisničke memorije u memoriju kernel moda mrežnog drajvera;
- Konačno, kopiranje iz memorije kernel moda mrežnog drajvera na mrežnu karticu za prenos.
Zato se može kroz "nulti kopiranje" smanjiti kontekstna promena korisničkih/kernel moda i broj kopiranja memorije, poboljšavajući IO performanse. Uobičajena implementacija "nultog kopiranja" je mmap, u Javi realizovano kroz MappedByteBuffer.

22.Kako se realizuje čišćenje diska?
RocketMQ nudi dve strategije čišćenja diska: sinhrono i asinhrono čišćenje
- Sinhrono čišćenje: kada poruka stigne do memorije Broker-a, mora se očistiti u commitLog fajlu pre nego se smatra uspešnim, zatim se vraća Proizvođaču poruka "slanje uspešno".
- Asinhrono čišćenje: asinhrono čišćenje znači da kada poruka stigne do memorije Broker-a, vraća Proizvođaču poruku "slanje uspešno", budi se nit koja će perzistirati podatke u CommitLog fajl.
Broker direktno operiše memorijom (memorijski mapirani fajl) pri skladištenju poruka, to može povećati propusnost sistema, ali ne može izbegnuti gubitak podataka pri nestanku struje, zato treba perzistirati na disk.
Konačno čišćenje je uvek korišćenje NIO MappedByteBuffer.force() da se podaci iz mapiranog regiona upišu na disk; ako je sinhrono čišćenje, Broker će nakon što piše poruku u CommitLog mapirani region čekati završetak pisanja.
Asinhrono čisti samo budi odgovarajuću nit, ne garantuje vreme izvršenja, tok je kao na slici.

23.Možeš reći kako se RocketMQ balansiranje opterećenja realizuje?
Balansiranje opterećenja u RocketMQ-u se vrši na klijentskoj strani, konkretno može se podeliti na balansiranje pri slanju poruka na strani Proizvođača i balansiranje pri konzumiranju na strani Konzumera.
Balansiranje opterećenja Proizvođača
Proizvođač pri slanju poruka prvo na osnovu Topic-a pronađe odgovarajući TopicPublishInfo, nakon što dobija rutirane informacije TopicPublishInfo, klijent RocketMQ-a u podrazumevanom modu selectOneMessageQueue() metodom izabere jedan red (MessageQueue) iz messageQueueList u TopicPublishInfo i šalje poruku. Ovde postoji promenljiva sendLatencyFaultEnable; ako je uključena, na osnovu random inkrementalnog modula dodatno filtrira "not available" Broker-e.

Takozvana "latencyFaultTolerance" znači da se za ranije neuspešne, za određeno vreme povlači. Na primer, ako je prethodni latency prešao 550Lms, povlači se 3000Lms; preko 1000L, povlači se 60000L; ako je isključena, koristi se random inkrementalni mod za izbor MessageQueue-a za slanje, mehanizam latencyFaultTolerance je ključni za visoku dostupnost slanja poruka.
Balansiranje opterećenja Konzumera
U RocketMQ-u, oba konzumna moda (Push/Pull) na strani Konzumera se temelje na pull modu; Push mod je samo omotač pull moda, u suštini implementacija je nit za izvlačenje poruka koja nakon što sa servera izvleti jednu grupu poruka, prosledi ih u bazen za konzumaciju, i potom "zaustavlja" i ponovo pokušava da izvuku sa servera. Ako se ne uspeš izvući, malo pauzira i ponovo pokušava. U oba pull moda (Push/Pull) Konzumer mora da zna iz kojeg MessageQueue Broker-a da izvlači poruke. Zato je potrebno uraditi balansiranje opterećenja na strani Konzumera, odnosno dodeliti više MessageQueue iz Broker-a određenim Konzumerima u istoj ConsumerGroup.
- Slanje heartbeat-a Konzumer strane
Nakon pokretanja Konzumer-a, kroz periodični zadatak neprestano šalje heartbeat-e svim Broker instancama u RocketMQ klasteru (uključujući naziv konzumer grupe, skup pretplatničkih odnosa, komunikacioni mod i ID klijenta itd.). Broker strana nakon prijema heartbeat-a Konzumer-a čuva ga u lokalnom kešu ConsumerManager-a — consumerTable, a zatim sačuva omotane informacije o mrežnom kanalu klijenta u lokalni keš — channelInfoTable, čime se pruža metapodatak za balansiranje opterećenja Konzumer strane.
- Ključna klasa za balansiranje opterećenja na Konzumer strani — RebalanceImpl
U procesu pokretanja Konzumer instance, prilikom pokretanja MQClientInstance instance, pokreće se nit za balansiranje opterećenja — RebalanceService (izvršava se svakih 20s).
Proverom izvornog koda može se videti da metoda run() nit RebalanceService na kraju poziva metodu rebalanceByTopic() klase RebalanceImpl, i to je ključna metoda za balansiranje opterećenja na Konzumer strani.
rebalanceByTopic() metoda vrši različitu logiku u zavisnosti od toga da li je komunikacioni tip konzumera "difuzni mod" ili "klaster mod". Pogledajmo glavni tok klaster moda:

(1) Iz lokalnog keša rebalanceImpl instance — topicSubscribeInfoTable dobija skup MessageQueue pod temom (mqSet);
(2) Na osnovu topic i consumerGroup kao parametre poziva mQClientFactory.findConsumerIdList() metodu i šalje komunikacioni zahtev Broker strani, dobija listu ID-eva konzumera iz te konzumer grupe;
(3) Prvo sortira MessageQueue i ID konzumera pod temom, a zatim algoritmom dodele redova (podrazumevano: algoritam prosečne dodele) izračunava MessageQueue koje treba izvlačiti. Ovaj algoritam prosečne dodele sličan je algoritmu straničavanja — sve MessageQueue sortiramo kao zapise, sve Konzumer sortiramo kao stranice, izračunamo prosečnu veličinu svake stranice i opseg svake stranice, na kraju prolazimo kroz ceo opseg i izračunavamo koju MessageQueue trenutni Konzumer treba da dobije.

(4) Zatim poziva updateProcessQueueTableInRebalance() metodu, konkretan korak je da prvo uporedi dodeljeni skup MessageQueue (mqSet) sa processQueueTable.

- Crveni deo processQueueTable na gornjoj slici označava delove koji se ne preklapaju sa dodeljenim mqSet. Te redove postavlja Dropped=true, a zatim proverava da li te redove može ukloniti iz processQueueTable keša, konkretno izvršava removeUnnecessaryMessageQueue() metodu, odnosno svake 1s proverava da li može da dobije zaključavanje trenutnog reda konzumcije; ako može, vraća true. Ako čeka 1s i još uvek ne može dobiti zaključavanje, vraća false. Ako vrati true, uklanja odgovarajući Entry iz processQueueTable keša;
- Zeleni deo processQueueTable na gornjoj slici označava presek sa dodeljenim mqSet. Proverava da li je ProcessQueue istekao; u Pull modu se ne mari, ako je Push mod, postavlja Dropped=true i poziva removeUnnecessaryMessageQueue() metodu, kao gore, pokušava ukloniti Entry;
- Konačno, za svaki MessageQueue u filtriranom mqSet kreira ProcessQueue objekat i smešta u red RebalanceImpl processQueueTable (prilikom čega poziva computePullFromWhere(MessageQueue mq) metodu RebalanceImpl instance da dobije sledeću vrednost ofseta konzumcije tog MessageQueue-a, zatim popunjava u pullRequest objekat koji će se kreirati), i kreira zahtev za izvlačenje — pullRequest dodaje u listu za izvlačenje — pullRequestList, na kraju izvršava dispatchPullRequest() metodu, pullRequest objekte za izvlačenje PullMessageRequest redom stavlja u blokirajući red pullRequestQueue servisne nit PullMessageService, koji ih izvlači i šalje zahtev za Pull poruku ka Broker strani. Ovde se može fokusirati na razliku između metoda dispatchPullRequest() klasa RebalancePushImpl i RebalancePullImpl; u RebalancePullImpl ova metoda je prazna.
Balansiranje opterećenja konzumacije istog MessageQueue između različitih konzumera iste konzumer grupe se temelji na ključnom dizajnu: u istom trenutku jedan MessageQueue može konzumirati samo jedan konzumer iz iste konzumer grupe, a jedan konzumer može istovremeno konzumirati više MessageQueue.
24.RocketMQ dugo anketiranje je poznato?
Dugo anketiranje znači da Konzumer izvlači poruke; ako odgovarajući Queue nema podataka, Broker ne vraća odmah, već zadržava PullRequest, čeka da Queue ima poruke ili istekne vreme dlogog anketiranja, a zatim ponovo obrađuje sve PullRequest na tom redu.

- PullMessageProcessor#processRequest
//ako se ništa ne izvuce
case ResponseCode.PULL_NOT_FOUND:
// broker i consumer dozvoljavaju suspend, podrazumevano uključeno
if (brokerAllowSuspend && hasSuspendFlag) {
long pollingTimeMills = suspendTimeoutMillisLong;
if (!this.brokerController.getBrokerConfig().isLongPollingEnable()) {
pollingTimeMills = this.brokerController.getBrokerConfig().getShortPollingTimeMills();
}
String topic = requestHeader.getTopic();
long offset = requestHeader.getQueueOffset();
int queueId = requestHeader.getQueueId();
//kapsuliramo PullRequest
PullRequest pullRequest = new PullRequest(request, channel, pollingTimeMills,
this.brokerController.getMessageStore().now(), offset, subscriptionData, messageFilter);
//zadržavamo PullRequest
this.brokerController.getPullRequestHoldService().suspendPullRequest(topic, queueId, pullRequest);
response = null;
break;
}Za zadržane zahteve postoji servisna nit koja neprestano proverava da li queue ima podatke ili je isteklo vreme.
- PullRequestHoldService#run()
@Override
public void run() {
log.info("{} service started", this.getServiceName());
while (!this.isStopped()) {
try {
if (this.brokerController.getBrokerConfig().isLongPollingEnable()) {
this.waitForRunning(5 * 1000);
} else {
this.waitForRunning(this.brokerController.getBrokerConfig().getShortPollingTimeMills());
}
long beginLockTimestamp = this.systemClock.now();
//proverava zadržane zahteve
this.checkHoldRequest();
long costTime = this.systemClock.now() - beginLockTimestamp;
if (costTime > 5 * 1000) {
log.info("[NOTIFYME] check hold request cost {} ms.", costTime);
}
} catch (Throwable e) {
log.warn(this.getServiceName() + " service has exception. ", e);
}
}
log.info("{} service end", this.getServiceName());
}Detaljno objašnjene česte teme za intervju o RocketMQ-u, ovaj put "prebićemo" intervjutera, mislim da je sigurno (manualni dog). Priredio: Chenmo Wang Er, link za preuzimanje, autor: Sanfen E, link do originalnog teksta.
Ništa me ne može zaustaviti — osim cilja, čak i ako na obali ima ruža, hladovine, mirne luke, ja sam brod koji nije vezan.
Serijski sadržaj:
- Priprema za intervju Java SE deo 👍
- Priprema za intervju Java kolekcijski framework 👍
- Priprema za intervju Java konkurentno programiranje 👍
- Priprema za intervju JVM deo 👍
- Priprema za intervju Spring deo 👍
- Priprema za intervju Redis deo 👍
- Priprema za intervju MyBatis deo 👍
- Priprema za intervju MySQL deo 👍
- Priprema za intervju Operativni sistemi deo 👍
- Priprema za intervju Računarske mreže deo 👍
- Priprema za intervju RocketMQ deo 👍
- Priprema za intervju Distribuirani sistemi deo 👍
- Priprema za intervju Mikroservisi deo 👍
- Priprema za intervju Dizajn obrasci deo 👍
- Priprema za intervju Linux deo 👍
- Priprema za intervju OpenClaw deo 👍
- Priprema za intervju Skills deo 👍
GitHub repozitorijum sa 17000+ zvezdica "Priprema za intervju" drugo PDF izdanje je konačno stiglo! Uključuje osnovu Jave, kolekcijski framework Jave, Java konkurentno programiranje, JVM, Spring, Redis, MyBatis, MySQL, operativni sistemi, računarske mreže, RocketMQ, distribuirani sistemi, mikroservisi, dizajn obrasci, Linux, OpenClaw itd. — ukupno preko 320.000 reči, 500+ crteža, zaista lako razumljivo, zabavno... Detalji: Priprema za intervju 2.0 izdanje PDF objavljen, Java pozadinski programeri moraju da znaju, možda najbolje učne teme 2026. godine
