← Back to Blog

Pulsar vs Kafka: Melyiket válasszam?

pulsarkafkamessagingdistributed-systemscomparison

Pulsar vs Kafka: Melyiket válasszam?

April 16, 2026
862 views
4.0
Paál Gyula
Paál Gyula
Founder & Lead Architect

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)

HelyzetAjánlás
Új projekt, nincs legacy KafkaKafka - easier to start, more people know it
Multi-tenant, több céges ügyfélPulsar - beépített multi-tenancy
Geo-replikáció kell (több datacenter)Pulsar - sokkal simpler
Nagyon nagy throughput, millió msg/szKafka - 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-dMaradj 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).

java
1// 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.

java
1// 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:

FeatureKafka Schema RegistryPulsar Schema
Külön serviceIgen (Confluent)Nem (beépített)
HTTP APIIgenIgen
VersioningIgenIgen
ValidationProducer/ConsumerTopic szintű
StorageKafka topicBookKeeper

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:

code
1[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.

java
1// 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):

MetricKafkaPulsar
P99 latency (1K msg)~8-15ms~3-8ms
P99 latency (100B msg)~2-5ms~1-3ms
ThroughputNagyon magasMagas
Cold readLassabbGyorsabb (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.

yaml
1# 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:

bash
1# 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

ComponentKafkaPulsar
Broker++
ZooKeeper / KRaft+BookKeeper (több node kell)
Schema RegistryKülön serviceBeépített
MonitoringJMX / PrometheusPrometheus
OperátorConfluent Operator / kubectlKubernetes 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

ProviderKafkaPulsar
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:

  1. Van már Kafka infrastructure-d - ne váltás csak mert új
  2. Stream processing kell (ksqlDB, Flink) - az ecosystem jobb
  3. Nagy a közösség - könnyebb embert találni
  4. High-throughput event streaming - a Kafka erre van optimalizálva
  5. Event sourcing - van hozzá tooling

Válassz Pulsar-t, ha:

  1. Multi-tenant kell - a Pulsar erre született
  2. Geo-replication fontos - beépítve van, jól működik
  3. Alacsonyabb latency kell - általában jobb numbers
  4. Tiered storage kell - régi adat automatikusan megy S3-ra
  5. Egyszerűbb operáció - kevesebb moving part

A döntés nem mindig egyértelmű

code
1Ha ú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)

go
1package 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)

go
1package 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)

go
1package 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)

go
1package 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)

go
1// 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)

go
1// 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?

AspektusKafkaPulsar
Topic-ok számaTöbb (orders, inventory, notification)Kevesebb (1 orders, namespace-ek)
Geo-replikációMirrorMaker2 (külön tool)Beépített (pulsar-admin)
Multi-tenantKülső (sauthc, separate clusters)Beépített (tenant/namespace)
Kód komplexitásTöbb különKevesebb
OperationsMirrorMaker2 configPulsar 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.

code
1Order → Payment Initiated → Bank Webhook → Success/Fail → Inventory → Shipping
23                                   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ó".

go
1// 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)

go
1// 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)

go
1// 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?

AspektusKafkaPulsar
Retry/DLQManuálisan (külön topic)Beépítve (Retry, DLQPolicy)
Saga orchestrationKülső lib (Quarkus, Temporal)Pulsar Functions (lambda-style)
Exactly-onceIdempotent producer + transzakciókBeépített transactions
At-least-onceAlapbólAlapból
Kód mennyiségTö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.

code
1Ha 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

FeatureKafkaPulsar
LicenseApache 2.0Apache 2.0
First release20112016
LanguageScala + JavaJava
StorageLocal diskBookKeeper
Tiered storage❌ (külső)
Multi-tenant❌ (külső)
Geo-replicationMirrorMaker2Beépített
Schema RegistryConfluentBeépített
Stream processingksqlDB, FlinkPulsar Functions
Latency (P99)~8-15ms~3-8ms
ThroughputNagyon magasMagas
Managed serviceSokKevés
CommunityNagyon nagyKözepes
Learning curveAlacsonyabbMagasabb

🏁 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?

A:

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?

A:

Á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?

A:

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?

A:

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?

A:

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:

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?

A:

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.

Follow us
All Rights Reserved
© 2011-2026
Progressive Innovation
LAB