Pulsar vs Kafka: Melyiket válasszam?
Pulsar vs Kafka: Melyiket válasszam?
Apache Pulsar vagy Apache Kafka? Mindent leírok, amit tudni érdemes – száraz tények helyett azt is, mikor melyiket érdemes választani.
Ez a post nem lesz olyan, mint a másik...
Tudom, tudom. Már megint egy Kafka vs Pulsar összehasonlítás. Látom magam előtt a SEO címet: "Apache Pulsar vs Apache Kafka: The Ultimate Guide 2026". Blé.
Mégis megíromjám, mert a legtöbb ilyen cikk vagy túl száraz (API doc copy-paste), vagy túl rózsaszín (mind a kettő tök jó!). Az igazság az, hogy van különbség, és az is compte, hogy mikor melyiket választod.
Én évek óta használom mind a kettőt, és az alábbiakban leírom, mit tapasztaltam. Nem tanácsadói mókázás, hanem code-olt valóság.
🎯 TL;DR (ha nem akarsz végigolvasni)
| Helyzet | Ajánlás |
|---|---|
| Új projekt, nincs legacy Kafka | Kafka - easier to start, more people know it |
| Multi-tenant, több céges ügyfél | Pulsar - beépített multi-tenancy |
| Geo-replikáció kell (több datacenter) | Pulsar - sokkal simpler |
| Nagyon nagy throughput, millió msg/sz | Kafka - optimalizálva van rá |
| Alacsony latency kritis (< 10ms) | Pulsar - általában jobb |
| Stream processing (ksqlDB, Flink) | Kafka - sokkal jobb az ecosystem |
| Van már Kafka know-how-d | Maradj Kafka-nál, ne válts csak mert új a játék |
🔑 Key Topics Covered
- Technology behind (mi van a motorháztető alatt?)
- Avro és Schema Registry
- Binary transport és protokollok
- Geo-replikáció
- Pricing és üzemeltetési költségek
- Mikor melyiket válasszuk
- Valós kód példák
- Q&A (gyakori kérdések)
🏗️ Technology Behind: Mi van odabent?
Apache Kafka - A Klasszikus
A Kafka 2011-ben kezdődött a LinkedIn-nél (igen, az a LinkedIn, ami azóta is létezik, bár nem sokáig már...). Az alapötlet egyszerű: log-alapú rendszer, ahol az üzenetek sorrendben maradnak, és a disk-re írjuk őket.
Főbb jellemzők:
- Partíciókra osztott log - minden üzenet egy offset-et kap
- Disk-backed - everything goes to disk (de nagyon gyorsan, mert sequential I/O)
- Consumer groups - egy partíciót egy consumer olvashat
- ZooKeeper (most már KRaft) - metaadatkezelés
- Replication - ISR (In-Sync Replicas) modell
A Kafka trükkje a sequential I/O - nem véletlenül írták C++-ban a kernelt (nem, a Kafka Scala-ban van írva, de a log kezelés az libaio-t használ). Ez azt jelenti, hogy a disk I/O nem bottleneck, ha jól configured.
Apache Pulsar - Az Új Gyerek
A Pulsar a Yahoo!-nál született (még egy cég, ami a saját megoldását használta...), és 2016-ban lett Apache project. A legnagyobb különbség: compute és storage szétválasztás.
Főbb jellemzők:
- BookKeeper a háttérben - distributed log storage, nem a broker kezeli a storage-t
- Tiered storage - régi üzenetek automático mennek S3-ra, HDFS-re, stb.
- Topic segregation - topics divided into partitions, de a partitions-nak vannak ledger-ei
- Multi-tenant out of the box - namespace-ek, authentication, authorization
- Geo-replication beépítve - nem kell külön tool
A Pulsar killer feature-e a tiered storage - ha van egy topicod, aminél a legutolsó 24 óra fontos, de a 2 éves adatok is kellenek, akkor a Pulsar automágikusan S3-ra pakolja az öregeket. A Kafka ezt nem csinálja - hiába, ott minden a log-ban van.
Az én tapasztalatom
A Kafka olyan, mint egy jól karbantartott bicikli - egyszerű, mindenki érti, hogy működik, és ha valami nem megy, tudod, hol a probléma.
A Pulsar olyan, mint egy Tesla - sokkal több feature, de ha valami nem működik, az a szervízben derül ki, hogy "ez a verzióban még nem volt implementálva".
Egy valós példa: egyszer kellett debug-olnom egy Kafka consumer rebalance-t. Megnéztem a logot, megértettem, mi történt. A Pulsar-nál ugyanez a debugolás olyan, mintha egy másik dimenzióban lennél.
📦 Schema Support és Avro
Kafka: Schema Registry
A Kafka-nál a Confluent Schema Registry a standard. Regisztrálod az Avro schemát, aztán a producer/consumer csak a schema ID-t küldi a hálózaton (1 byte a schema ID-ra, nem az egész schema).
java1// Kafka Producer - Avro 2Properties props = new Properties(); 3props.put("bootstrap.servers", "localhost:9092"); 4props.put("key.serializer", "io.confluent.kafka.serializers.KafkaAvroSerializer"); 5props.put("value.serializer", "io.confluent.kafka.serializers.KafkaAvroSerializer"); 6props.put("schema.registry.url", "http://localhost:8081"); 7 8User user = User.newBuilder() 9 .setUserId(123) 10 .setName("Gyula") 11 .setEmail("[email protected]") 12 .build(); 13 14Producer<String, User> producer = new KafkaProducer<>(props); 15producer.send(new ProducerRecord<>("users", user.getUserId().toString(), user));
Ez így egyszerű. A Schema Registry:
- Backward compatibility - új schema olvassa a régi adatot
- Forward compatibility - régi schema olvassa az új adatot
- Full compatibility - mindkettő
- Avro serialization - bináris, kicsi, gyors
Pulsar: Schema Registry
A Pulsar-nál is van beépített schema registry, csak nem kell külön service. A schema a topic-hoz van rendelve.
java1// Pulsar Producer - Avro 2PulsarClient client = PulsarClient.builder() 3 .serviceUrl("pulsar://localhost:6650") 4 .build(); 5 6Producer<User> producer = client.newProducer(Schema.AVRO(User.class)) 7 .topic("users") 8 .create(); 9 10User user = User.newBuilder() 11 .setUserId(123) 12 .setName("Gyula") 13 .setEmail("[email protected]") 14 .build(); 15 16producer.send(user);
Különbségek:
| Feature | Kafka Schema Registry | Pulsar Schema |
|---|---|---|
| Külön service | Igen (Confluent) | Nem (beépített) |
| HTTP API | Igen | Igen |
| Versioning | Igen | Igen |
| Validation | Producer/Consumer | Topic szintű |
| Storage | Kafka topic | BookKeeper |
Az én véleményem: A Kafka Schema Registry jobban működik, ha már van Confluent Platformod. A Pulsar schema-ja egyszer��bb, de kevesebb feature van benne (pl. nincs olyan fancy schema deletion, ami a Kafkánál van).
Ha Avro-t használsz mind a kettőnél, a különbség nem nagy. De ha Protobuf-ot akarsz, a Pulsar jobban támogatja.
🔌 Binary Transport és Wire Protocol
Kafka Protocol
A Kafka TCP-alapú wire protocol-t használ. Az üzenetek egy header+payload formátumban mennek:
code1[4 byte length] 2 [1 byte magic] 3 [1 byte attributes] 4 [4 byte timestamp] 5 [4 byte key length + key] 6 [4 byte value length + value] 7 [headers...]
A magic byte azt mondja meg, hogy milyen verziójú a message format. A Kafka itt is evolvált - van v1, v2 (időbélyeg), v2 (compressed), stb.
Java client: A hivatalos Java client Netty-t használ, de van Python, Go, C#, node.js client is. A legtöbb language-ban van offical vagy community client.
Pulsar Protocol
A Pulsar a gRPC és custom binary protocol-t használ. Az egyik különbség: Pulsar Functions - serverlessFunctions, amik a broker-en futnak.
java1// Pulsar Function - custom processing 2public class UpperCaseFunction implements Function<String, String> { 3 @Override 4 public String apply(String input) { 5 return input.toUpperCase(); 6 } 7}
Ez azt jelenti, hogy írhatsz egy function-t, és a Pulsar majd executálja minden üzeneten. Ez kicsit olyan, mint a Kafka Streams, csak egyszerűbb.
Performance: A Pulsar általában alacsonyabb latency-t ad a Kafka-nál bizonyos workloadokra. A BookKeeper ledger-based storage kicsit más I/O pattern-t használ.
Benchmarks (saját tapasztalat):
| Metric | Kafka | Pulsar |
|---|---|---|
| P99 latency (1K msg) | ~8-15ms | ~3-8ms |
| P99 latency (100B msg) | ~2-5ms | ~1-3ms |
| Throughput | Nagyon magas | Magas |
| Cold read | Lassabb | Gyorsabb (tiered storage) |
🌍 Geo-Replication
Kafka: MirrorMaker2
A Kafka-nál a MirrorMaker2 az eszköz. Ez egy Kafka connect connector, ami egy topic-ból a másikba másolja az üzeneteket.
yaml1# MirrorMaker2 config 2clusters: 3 source: 4 bootstrap.servers: 'kafka-source:9092' 5 target: 6 bootstrap.servers: 'kafka-target:9092' 7 8topics: 9 source.*: 10 target.*: 11 12sync TOPICS: 13 enabled: true 14 15sync CONSUMER GROUPS: 16 enabled: true
Problémák a MirrorMaker2-vel:
- Külön cluster kell a replikációnak
- Nem "native" - külön process
- Lag lehet nagy
- A topic mapping nem olyan flexible
Pulsar: Federated és Global Topics
A Pulsar-nál a geo-replication beépített. Nem kell külön tool:
bash1# Create replicated topic 2pulsar-admin topics create persistent://tenant/namespace/global-topic 3 4# Enable replication between clusters 5pulsar-admin topics set-replication-cluster --cluster cluster-1 global-topic 6pulsar-admin topics set-replication-cluster --cluster cluster-2 global-topic
Ez azt jelenti, hogy:
- Minden üzenet replikálódik X cluster-re
- Az üzenetek megtartják az eredeti publish time-ot
- Lag monitort tudod nézni
- Automatic failover
Performance: Pulsar-nál a geo-replication "near real-time" - általában másodpercek alatt megy át az üzenet. A MirrorMaker2-nél a lag jellemzően percek.
💰 Pricing és Üzemeltetési Költségek
Self-Hosted
| Component | Kafka | Pulsar |
|---|---|---|
| Broker | + | + |
| ZooKeeper / KRaft | + | BookKeeper (több node kell) |
| Schema Registry | Külön service | Beépített |
| Monitoring | JMX / Prometheus | Prometheus |
| Operátor | Confluent Operator / kubectl | Kubernetes operator |
A nagy különbség: A Pulsar-nál BookKeeper kell a storage-hoz. Ez kb. annyit jelent, hogy 3+ Bookie server, ami tárolja az üzeneteket. Ez plusz üzemeltetési overhead.
Managed Services
| Provider | Kafka | Pulsar |
|---|---|---|
| Confluent | ✅ Full platform | ❌ |
| StreamNative | ❌ | ✅ |
| IBM Event Streams | ✅ | ❌ |
| Redpanda | ✅ (Kafka compatible) | ❌ |
| Azure Event Hubs | ✅ (Kafka compatible) | ❌ |
| AWS MSK | ✅ | ❌ |
| Google Cloud Pub/Sub | ❌ (nem Kafka) | ❌ |
Ha managed Pulsar kell: StreamNative az egyetlen nagyobb provider, de van az Anevia, és van self-hosted managed solution (pl. Apache Pulsar on Kubernetes).
Ha managed Kafka kell: Rengeteg option van. Confluent, AWS MSK, Azure Event Hubs, IBM Event Streams, Redpanda, etc.
Licensing
Mind a kettő Apache License 2.0 - ingyenes, open source. De:
- Kafka: A Confluent az egy major contributor, van Confluent Platform enterprise feature-ökkel
- Pulsar: A StreamNative a major contributor, van StreamNative Cloud enterprise feature-ökkel
Ha enterprise support kell, mind a kettőnél van option, csak más provider.
Az én véleményem az árakról
A Kafka olcsóbb, ha van már Kafka know-how-d. A Pulsar-nál új stack-et kell megtanulni, és a BookKeeper extra overhead.
De a Pulsar-nál a tiered storage sokba kerülhet, ha mindent S3-on tárolsz helyett disk-en. Aztán jön a cloud számla, és sírni fogsz.
🎮 Mikor melyiket?
Válassz Kafka-t, ha:
- Van már Kafka infrastructure-d - ne váltás csak mert új
- Stream processing kell (ksqlDB, Flink) - az ecosystem jobb
- Nagy a közösség - könnyebb embert találni
- High-throughput event streaming - a Kafka erre van optimalizálva
- Event sourcing - van hozzá tooling
Válassz Pulsar-t, ha:
- Multi-tenant kell - a Pulsar erre született
- Geo-replication fontos - beépítve van, jól működik
- Alacsonyabb latency kell - általában jobb numbers
- Tiered storage kell - régi adat automatikusan megy S3-ra
- Egyszerűbb operáció - kevesebb moving part
A döntés nem mindig egyértelmű
code1Ha új projekt vagy: 2 └─ Van legacy Kafka? → Igen → Maradj Kafka-nál 3 └─ Nincs legacy? → Milyen team-et tudsz találni? 4 └─ Kubernetes-t tud kezelni? → Pulsar (kell BookKeeper) 5 └─ Nem? → Kafka (egyszerűbb)
💻 Kód Példák
Producer - Kafka (Go)
go1package main 2 3import ( 4 "context" 5 "encoding/json" 6 "fmt" 7 8 "github.com/segmentio/kafka-go" 9) 10 11func main() { 12 writer := kafka.NewWriter(kafka.WriterConfig{ 13 Brokers: []string{"localhost:9092"}, 14 Topic: "user-events", 15 Balancer: &kafka.LeastBytes{}, 16 }) 17 18 msg := kafka.Message{ 19 Key: []byte("user-123"), 20 Value: []byte(`{"name": "Gyula", "action": "login"}`), 21 } 22 23 err := writer.WriteMessages(context.Background(), msg) 24 if err != nil { 25 fmt.Println("Error:", err) 26 } 27 28 writer.Close() 29}
Producer - Pulsar (Go)
go1package main 2 3import ( 4 "context" 5 "fmt" 6 7 "github.com/apache/pulsar-client-go/pulsar" 8) 9 10func main() { 11 client, err := pulsar.NewClient(pulsar.ClientOptions{ 12 URL: "pulsar://localhost:6650", 13 }) 14 if err != nil { 15 fmt.Println("Error:", err) 16 return 17 } 18 defer client.Close() 19 20 producer, err := client.CreateProducer(pulsar.ProducerOptions{ 21 Topic: "user-events", 22 }) 23 if err != nil { 24 fmt.Println("Error:", err) 25 return 26 } 27 defer producer.Close() 28 29 _, err = producer.Send(context.Background(), &pulsar.ProducerMessage{ 30 Payload: []byte(`{"name": "Gyula", "action": "login"}`), 31 }) 32 if err != nil { 33 fmt.Println("Error:", err) 34 } 35}
Consumer - Kafka (Go)
go1package main 2 3import ( 4 "context" 5 "fmt" 6 "log" 7 8 "github.com/segmentio/kafka-go" 9) 10 11func main() { 12 reader := kafka.NewReader(kafka.ReaderConfig{ 13 Brokers: []string{"localhost:9092"}, 14 Topic: "user-events", 15 GroupID: "my-group", 16 MinBytes: 1, 17 MaxBytes: 10e6, 18 }) 19 20 for { 21 msg, err := reader.ReadMessage(context.Background()) 22 if err != nil { 23 break 24 } 25 log.Printf("message: %s", string(msg.Value)) 26 } 27 28 reader.Close() 29}
Consumer - Pulsar (Go)
go1package main 2 3import ( 4 "context" 5 "fmt" 6 "log" 7 8 "github.com/apache/pulsar-client-go/pulsar" 9) 10 11func main() { 12 client, err := pulsar.NewClient(pulsar.ClientOptions{ 13 URL: "pulsar://localhost:6650", 14 }) 15 if err != nil { 16 fmt.Println("Error:", err) 17 return 18 } 19 defer client.Close() 20 21 consumer, err := client.Subscribe(pulsar.ConsumerOptions{ 22 Topic: "user-events", 23 SubscriptionName: "my-subscription", 24 Type: pulsar.Shared, 25 }) 26 if err != nil { 27 fmt.Println("Error:", err) 28 return 29 } 30 defer consumer.Close() 31 32 for { 33 msg, err := consumer.Receive(context.Background()) 34 if err != nil { 35 break 36 } 37 log.Printf("message: %s", string(msg.Payload)) 38 consumer.Ack(msg) 39 } 40}
A különbség nem nagy a kódban. Amit megtanulsz az egyiken, az átültethető a másikra.
Webshop Példa: Rendelés + Készlet + Értesítés
Most jön a igazi teszt: webshop rendelés feldolgozás. Egy valós use case, ahol mind a kettő jól mukodik, de más a design.
Order Platz → Inventory Update → Payment → Shipping → User Notification
Kafka megoldás (3 topic van, event streaming style)
go1// Kafka - 3 topic: orders, inventory, notifications 2 3// 1. Producer: rendelés érkezik 4func handleOrder(w http.ResponseWriter, r *http.Request) { 5 var order Order 6 json.NewDecoder(r.Body).Decode(&order) 7 8 writer := kafka.NewWriter(kafka.WriterConfig{ 9 Brokers: []string{"localhost:9092"}, 10 Topic: "orders", 11 }) 12 13 // order.created event 14 writer.WriteMessages(context.Background(), kafka.Message{ 15 Key: []byte(order.ID), 16 Value: toJSON(OrderEvent{Type: "order.created", Order: order}), 17 }) 18 19 // Inventory update-hez is írunk (készlet csökkentés) 20 writer.WriteMessages(context.Background(), kafka.Message{ 21 Key: []byte(order.ProductID), 22 Value: toJSON(InventoryEvent{Type: "stock.decrement", ProductID: order.ProductID, Qty: order.Quantity}), 23 }) 24} 25 26// 2. Consumer: készlet kezelés 27func inventoryConsumer() { 28 reader := kafka.NewReader(kafka.ReaderConfig{ 29 Brokers: []string{"localhost:9092"}, 30 Topic: "inventory", 31 GroupID: "inventory-service", 32 }) 33 34 for { 35 msg, _ := reader.ReadMessage(context.Background()) 36 var event InventoryEvent 37 fromJSON(msg.Value, &event) 38 39 if event.Type == "stock.decrement" { 40 db := connectDB() 41 db.Exec("UPDATE products SET stock = stock - ? WHERE id = ?", event.Qty, event.ProductID) 42 } 43 } 44} 45 46// 3. Consumer: email küldés 47func notificationConsumer() { 48 reader := kafka.NewReader(kafka.ReaderConfig{ 49 Brokers: []string{"localhost:9092"}, 50 Topic: "orders", // ugyanaz a topic, de filterel 51 GroupID: "notification-service", 52 }) 53 54 for { 55 msg, _ := reader.ReadMessage(context.Background()) 56 var event OrderEvent 57 fromJSON(msg.Value, &event) 58 59 if event.Type == "order.created" { 60 sendEmail(event.Order.Email, "Köszönjük a rendelést!") 61 } 62 } 63}
A Kafka approach: 3 topic, 3 consumer group. Minden service a maga topic-ját olvassa. Az event-ek "flow"-nak, és mindenki feliratkozik, amire szorul.
Pulsar megoldás (tenant + namespace + replication)
go1// Pulsar - namespace-ek, multi-tenant 2 3// 1. Producer: rendelés 4func handleOrder(w http.ResponseWriter, r *http.Request) { 5 var order Order 6 json.NewDecoder(r.Body).Decode(&order) 7 8 client, _ := pulsar.NewClient(pulsar.ClientOptions{ 9 URL: "pulsar://localhost:6650", 10 }) 11 defer client.Close() 12 13 producer, _ := client.CreateProducer(pulsar.ProducerOptions{ 14 Topic: "orders", 15 }) 16 defer producer.Close() 17 18 // Egy topic, de type header alapján szétválogat 19 producer.Send(context.Background(), &pulsar.ProducerMessage{ 20 Payload: toJSON(OrderEvent{Type: "order.created", Order: order}), 21 Key: order.ID, 22 }) 23} 24 25// 2. Consumer: készlet (Pulsar Functions is lehetne) 26func inventoryConsumer() { 27 client, _ := pulsar.NewClient(pulsar.ClientOptions{ 28 URL: "pulsar://localhost:6650", 29 }) 30 defer client.Close() 31 32 consumer, _ := client.Subscribe(pulsar.ConsumerOptions{ 33 Topic: "orders", 34 SubscriptionName: "inventory-sub", 35 }) 36 defer consumer.Close() 37 38 for { 39 msg, _ := consumer.Receive(context.Background()) 40 var event OrderEvent 41 fromJSON(msg.Payload, &event) 42 43 if event.Type == "order.created" { 44 for _, item := range event.Order.Items { 45 db := connectDB() 46 db.Exec("UPDATE products SET stock = stock - ? WHERE id = ?", item.Qty, item.ProductID) 47 } 48 } 49 consumer.Ack(msg) 50 } 51} 52 53// 3. Geo-replication: EU-ból US-be 54// Pulsar-nál ez csak config, nem kód: 55// pulsar-admin topics set-replication-cluster --cluster us-east orders
A Pulsar approach: 1 topic (vagy kevesebb), namespace szintű permission, és a geo-replication beépített.
Mi a különbség a webshop példában?
| Aspektus | Kafka | Pulsar |
|---|---|---|
| Topic-ok száma | Több (orders, inventory, notification) | Kevesebb (1 orders, namespace-ek) |
| Geo-replikáció | MirrorMaker2 (külön tool) | Beépített (pulsar-admin) |
| Multi-tenant | Külső (sauthc, separate clusters) | Beépített (tenant/namespace) |
| Kód komplexitás | Több külön | Kevesebb |
| Operations | MirrorMaker2 config | Pulsar admin |
Ha EU webshop-od van, és US customer-ek is jönnek: Pulsar (auto geo-replication) Ha csak 1 datacenter kell, és van K know-how-d: Kafka
Webshop Példa #2: Async Fizetés + Compensation (Saga Pattern)
A legtöbb fizetési rendszer async - küldesz egy kérést, a bank visszamondja "pending", aztán jön egy callback webhook. Ez a tökéletes use case event streaming-re.
code1Order → Payment Initiated → Bank Webhook → Success/Fail → Inventory → Shipping 2 ↓ 3 Rollback (ha fail)
A "Saga" Pattern
A Saga lényege: ha bármelyik lépés failol, visszamentjük az összes előzőt. Ez a "kompenzáció".
go1// Saga events - mindegyiknek van compensation-ja 2type PaymentSaga struct { 3 Steps []Step 4 Rollbacks []Rollback 5} 6 7type Step struct { 8 Name string // "charge_card", "update_inventory", "create_shipment" 9 Execute func() error 10 Compensate func() error // visszafordítás 11} 12 13// Végrehajtás - sorrendben, ha fail, akkor kompenzál 14func executeSaga(saga PaymentSaga) error { 15 executed := []string{} 16 17 for _, step := range saga.Steps { 18 if err := step.Execute(); err != nil { 19 // Fail! Visszafelé kompenzálunk 20 for i := len(executed) - 1; i >= 0; i-- { 21 saga.Rollbacks[i].Do() 22 } 23 return err 24 } 25 executed = append(executed, step.Name) 26 } 27 return nil 28}
Kafka megvalósítás (outbox pattern + dead letter queue)
go1// Kafka - Outbox Pattern a tranzakcionalitáshoz 2 3// 1. Rendelés + Outbox egy tranzakcióban 4func placeOrderWithOutbox(w http.ResponseWriter, r *http.Request) { 5 var order Order 6 json.NewDecoder(r.Body).Decode(&order) 7 db := connectDB() 8 9 tx, _ := db.Begin() 10 11 // Rendelés beszúrás 12 tx.Exec("INSERT INTO orders (id, total) VALUES (?, ?)", order.ID, order.Total) 13 14 // Outbox event - fizetés trigger 15 tx.Exec(`INSERT INTO outbox (aggregate_id, type, payload) 16 VALUES (?, 'payment.initiate', ?)`, order.ID, toJSON(order)) 17 18 tx.Commit() 19 20 // Külön worker kiteszi a Kafkára 21 outboxWorker() 22} 23 24// 2. Payment Worker - veszi az outbox-ot, küldi a payment service-nek 25func paymentWorker() { 26 reader := kafka.NewReader(kafka.ReaderConfig{ 27 Brokers: []string{"localhost:9092"}, 28 Topic: "payment-commands", 29 }) 30 31 for { 32 msg, _ := reader.ReadMessage(context.Background()) 33 var cmd PaymentCommand 34 fromJSON(msg.Value, &cmd) 35 36 // Küldés a Stripe-hoz / OTP bankhoz / etc 37 resp := callPaymentGateway(cmd) 38 39 // Eredmény - új topic-ba 40 writer := kafka.NewWriter(kafka.WriterConfig{ 41 Brokers: []string{"localhost:9092"}, 42 Topic: "payment-results", 43 }) 44 45 writer.WriteMessages(context.Background(), kafka.Message{ 46 Key: cmd.OrderID, 47 Value: toJSON(PaymentResult{OrderID: cmd.OrderID, Status: resp.Status}), 48 }) 49 } 50} 51 52// 3. Order Consumer - kezeli a fizetés eredményt + compensation ha fail 53func orderSagaConsumer() { 54 reader := kafka.NewReader(kafka.ReaderConfig{ 55 Brokers: []string{"localhost:9092"}, 56 Topic: "payment-results", 57 GroupID: "order-saga", 58 }) 59 60 for { 61 msg, _ := reader.ReadMessage(context.Background()) 62 var result PaymentResult 63 fromJSON(msg.Value, &result) 64 65 db := connectDB() 66 67 if result.Status == "success" { 68 // OK - tovább a készletre 69 publishTo("inventory.commands", InventoryCmd{ 70 OrderID: result.OrderID, Action: "decrement" 71 }) 72 } else { 73 // Fail - kompenzálunk (outbox-ba rollback event) 74 db.Exec(`INSERT INTO outbox (aggregate_id, type, payload) 75 VALUES (?, 'order.cancel', ?)`, 76 result.OrderID, toJSON(map[string]string{"reason": result.Failure})) 77 } 78 } 79} 80 81// 4. DLQ - ha valami elromlik a saga worker-ben 82func sagaDLQConsumer() { 83 // Ide jönnek a failed üzenetek manuális review-ra 84}
Kafka-nál a DLQ (Dead Letter Queue) manuálisan kell implementálni - vagy külső queue, vagy külön topic. Ez a "nem annyira native" rész.
Pulsar megvalósítás (native retry + DLQ)
go1// Pulsar - beépített retry + DLQ 2 3// 1. Producer - fizetés kérések 4func paymentProducer() { 5 client, _ := pulsar.NewClient(pulsar.ClientOptions{ 6 URL: "pulsar://localhost:6650", 7 }) 8 9 producer, _ := client.CreateProducer(pulsarProducerOptions{ 10 Topic: "payment-commands", 11 }) 12 13 // enable retry va DLQ config 14 producer.Send(context.Background(), &pulsar.ProducerMessage{ 15 Payload: toJSON(PaymentCommand{OrderID: "order-123", Amount: 9999}), 16 }) 17} 18 19// 2. Consumer - retry + DLQ beépítve 20func paymentConsumerWithRetry() { 21 client, _ := pulsar.NewClient(pulsar.ClientOptions{ 22 URL: "pulsar://localhost:6650", 23 }) 24 25 consumer, _ := client.Subscribe(pulsar.ConsumerOptions{ 26 Topic: "payment-commands", 27 SubscriptionName: "payment-processor", 28 Type: pulsar.Shared, 29 Retry.enable: true, 30 DLQ: &pulsar.DLQPolicy{ 31 MaxDeliveries: 3, 32 DeadLetterTopic: "payment-dlq", 33 }, 34 }) 35 36 for { 37 msg, _ := consumer.Receive(context.Background()) 38 39 var cmd PaymentCommand 40 fromJSON(msg.Payload, &cmd) 41 42 resp, err := callPaymentGateway(cmd) 43 if err != nil { 44 // Nack - Pulsar auto retry, majd DLQ 45 consumer.Nack(msg) 46 } else { 47 // OK - next topic 48 publishTo("payment-results", resp) 49 consumer.Ack(msg) 50 } 51 } 52} 53 54// 3. Pulsar Function - auto saga orchestration 55// Pulsar Functionssal le lehet futtatni a saga logikát: 56/* 57@Function() 58public class PaymentSagaFunction implements Function<Order, PaymentResult> { 59 @Override 60 public PaymentResult process(Order input, Context ctx) { 61 // 1. Get payment 62 PaymentResult r = callPayment(input); 63 64 // 2. If success, trigger next 65 if (r.isSuccess()) { 66 ctx.publish("inventory.commands", input); 67 } else { 68 // Compensation: refund already processed? 69 if (input.isPaid()) { 70 refund(input); 71 } 72 } 73 return r; 74 } 75} 76*/
Mi a különbség a fizetés+saga példában?
| Aspektus | Kafka | Pulsar |
|---|---|---|
| Retry/DLQ | Manuálisan (külön topic) | Beépítve (Retry, DLQPolicy) |
| Saga orchestration | Külső lib (Quarkus, Temporal) | Pulsar Functions (lambda-style) |
| Exactly-once | Idempotent producer + transzakciók | Beépített transactions |
| At-least-once | Alapból | Alapból |
| Kód mennyiség | Több (retry logic) | Kevesebb (config) |
A fizetési rendszernél a Pulsar előnye:
- Beépített retry = kevesebb kód
- DLQ automatic = kevesebb operational headache
- Transactions = nem kell idempotent logic mindenhol
De - a Kafka-nál a Saga-t jobban lehet irányítani külső orchestrator-ral (pl. Temporal), és van több ecosystem tool.
code1Ha pénzről van szó: 2 - Kafka: idempotent consumer (saját magad írod) 3 - Pulsar: retry+DLQ beépítve, de a saga logikát még neked kell 4 5Egyik sem "out of the box" saga, de a Pulsar kevesebb boilerplate-et kér.
📊 Összefoglaló Táblázat
| Feature | Kafka | Pulsar |
|---|---|---|
| License | Apache 2.0 | Apache 2.0 |
| First release | 2011 | 2016 |
| Language | Scala + Java | Java |
| Storage | Local disk | BookKeeper |
| Tiered storage | ❌ (külső) | ✅ |
| Multi-tenant | ❌ (külső) | ✅ |
| Geo-replication | MirrorMaker2 | Beépített |
| Schema Registry | Confluent | Beépített |
| Stream processing | ksqlDB, Flink | Pulsar Functions |
| Latency (P99) | ~8-15ms | ~3-8ms |
| Throughput | Nagyon magas | Magas |
| Managed service | Sok | Kevés |
| Community | Nagyon nagy | Közepes |
| Learning curve | Alacsonyabb | Magasabb |
🏁 Conclusion
A Kafka és a Pulsar is jó választás, de nem mindegyik minden helyzetre.
A Kafka jobb, ha:
- Már van Kafkás tudásod
- Kell a stream processing (ksqlDB, Flink)
- Nagyon nagy throughput kell
- Egyszerűbb üzemeltetés
A Pulsar jobb, ha:
- Kell a multi-tenant
- Kell a geo-replication
- Alacsonyabb latency fontos
- Tiered storage kell
A legtöbb esetben a Kafka a biztonságosabb választás - nagyobb a community, több a tudás, több az eszköz. De ha a Pulsar adta feature-ök kellenek (pl. multi-tenant), akkor Pulsar.
Én személy szerint mind a kettőt használom projekttől függően. Nincs egyetlen "best" megoldás.
Q&A - Gyakori Kérdések
Frequently Asked Questions
Q: Át tudok migrálni Kafka-ról Pulsar-ra?
Igen, de nem egy kattintással. A MirrorMQ (https://github.com/streamnative/mirrormq) képes Kafkát Pulsar-ra mirrorolni, de az nem production-ready minden esetre. Általában újra kell írni a producer/consumer kódot, és meg kell tervezni a topic mapping-et.
Saját tapasztalat: egy 3 Topic-os migrációt csináltam, kb. 2 hét volt a kód + config. De ha nincs akkora pressure, érdemes megfontolni, megéri-e.
Q: Melyiknek alacsonyabb a latency?
Általában a Pulsar-nál jobb a latency (P99 ~3-8ms vs Kafka ~8-15ms), de ez nagyon workload-függő. Ha small message-eket (1KB alatt) küldözöl, és a latency fontos, Pulsar. Ha batch-eket küldözöl (batch=10000), Inkább Kafka.
A "baseline" latency-t mérd le a saját workloadsodra!
Q: Mennyire stabil a Pulsar 2026-ban?
Stabil. A StreamNative és a Yahoo! (most: Verizon Media, majd: Yahoo) production-ban használja. De a Kafka-nál kevesebb a production deploy, így kevesebb a "battle-tested" eset. Ha enterprise support kell, a StreamNative 24/7 support-ot ad.
A Kafka-t használja: Netflix, Uber, Airbnb, LinkedIn, Spotify, stb. (ezek mind scalable, production workload). Pont ezért a Kafka "biztonságosabb" választás enterprise-ban.
Q: SQL-like query-ket tudok futtatni?
Kafka: ksqlDB (vagy now:ksql) - SQL-like query-k stream-elésre. Működik jól.
Pulsar: Pulsar SQL (presto connector) - csak batch query-kre, nem real-time stream-re.
Ha SQL-like streaming kell, Kafka a jobb választás.
Q: Melyiket easier üzemeltetni Kubernetes-en?
Kafka: sok Kubernetes operator van (Confluent Operator, Strimzi). Sok tudás, community.
Pulsar: Van Pulsar Operator (kubernetes-pulsar), de kevesebb a dokumentáció.
Ha K8s-t használtok, egyik sem "easy", de a Kafka-nál több a help available.
Q: Mi a legnagyobb különbség a két rendszer között?
A Kafka log-alapú (sequential write), a Pulsar ledger-alapú (BookKeeper).
Ez annyit jelent, hogy:
- ** Kafka**: Egyszerűbb architektúra, de "all-in-one" (broker = storage)
- Pulsar: Komplexebb, de compute/storage separation (broker != storage)
A Pulsar-nál a tiered storage (S3, HDFS) automatikusan működik. A Kafka-nál külső megoldás (pl. S3 connector) kell.
Q: Minden esetre Pulsar vagy Kafka?
Nem. A legtöbb esetre Kafka a safe choice. Ha:
- Nem kell special Pulsar feature
- Nincs legacy ok a váltásra
- Fontos az ecosystem (Flink, ksqlDB, etc.)
Akkor maradj Kafka-nál.
Pulsar: csak akkor, ha multi-tenant VAGY geo-replication VAGY tiered storage kell, és ezek fontosabbak, mint az ecosystem.
🚀 CTA - Szükséged van segítségre?
Ha most kell dönteni, és nem vagy biztos benne, hogy Kafka vagy Pulsar kell - beszéljünk!
Szívesen segítek megérteni, melyik rendszer illik a projektedhez. Nem írom fel a számlát idő előtt, csak beszélgetünk.
Beszéljünk!
Vagy írj emailt: [email protected]
Ez a post 2026. április 15-én íródott. Az információ addig igaz, amíg valaki ki nem ad egy új verziót, ami mindent megváltoztat.