DEV Community

Firat Celik
Firat Celik

Posted on

Kafka Just Learned to Be a Queue. Postgres Has Been One All Along.

On February 17, 2026, Apache Kafka 4.2 shipped with this line in the release announcement: "Kafka Queues (Share Groups) is now production-ready."

That sentence quietly reopens one of backend engineering's oldest arguments. For years, "should my job queue be Postgres or Kafka?" had a lazy answer from both camps. The Kafka camp said Kafka scales. The Postgres camp said Kafka is not a queue: your consumer count is capped by your partition count, there is no per-message ack, and one poison message blocks everything behind it in its partition.

The second argument just lost most of its teeth. So this post rebuilds the comparison from primary sources and two small demos you can run in a few minutes: a Postgres SKIP LOCKED queue in Python, and a Kafka share group in Go.

The thesis

If a job is born inside a database transaction, Postgres is still the better queue in 2026.

Kafka share groups are real, well designed, and production-ready on the broker. But they answer a different question: "I already run Kafka, and now I also need queue semantics." They do not answer "I need a queue." Those two questions have different winners, and a team that does not already run Kafka would be adopting a whole platform to get one feature.

The antithesis, at full strength

Here is the best case against that thesis. It is a good one.

KIP-932 adds a new kind of group, the share group, and with it:

  • Consumers beyond partition count. "The number of consumers in a share group can exceed the number of partitions in a topic." The old trick of over-partitioning just to get parallelism goes away.
  • Per-record acknowledgement. A consumer can accept, release (retry), or reject (unprocessable) each record. Kafka 4.2 added a fourth option, RENEW, to extend the lock for slow work (KIP-1222).
  • Delivery counting with a poison-message limit. After the configured number of attempts (default 5, per the 4.3 broker configs), a record is archived instead of redelivered forever.
  • No queue depth limit and replay. In the KIP's own words: "queues done in a Kafka way, with no maximum queue depth and the ability to reset to a specific time for point-in-time recovery."
  • Many readers, one log. Several share groups can consume the same topic independently, next to ordinary consumer groups.

And the Postgres side has a well-known way to rot: MVCC. In the demo below, a single idle session holding an old snapshot made the median claim latency go from about 0.27 ms to about 15 ms on the same queue. That is not a theoretical concern. It is one forgotten BEGIN away.

So the antithesis is: Kafka now does queues properly, scales past one box, keeps history, and does not degrade because someone left a psql session open. The answer to that lives in the failure modes on both sides.

What actually shipped (facts, as of October 5, 2026)

Kafka share groups (KIP-932) Postgres FOR UPDATE SKIP LOCKED queue
Unit of work Record in a topic partition Row in a table
Claim Broker-side acquisition lock, default 30 s (share.record.lock.duration.ms) Row lock held by your transaction
Ack ACCEPT / RELEASE / REJECT / RENEW COMMIT (and delete or update the row)
Crash behavior Lock expires (or the share session closes), record becomes available again Transaction aborts, lock is released, row is visible again
Delivery guarantee At-least-once (KIP-932) At-least-once for external side effects; atomic with any writes in the same database
Ordering Not guaranteed across batches, especially on redelivery Whatever your ORDER BY gives you, best effort under concurrency
Poison messages Delivery count limit, default 5 You write it (an attempts column)
Dead-letter queue KIP-1191 listed as complete for 4.4.0, which is not released yet You write it (a table)

Client support matters more than broker support, because your code talks to a client:

  • Java: KafkaShareConsumer is the reference client (javadoc).
  • Go: franz-go added "full support for KIP-932 share groups" in v1.21.0 (changelog). The demo below uses v1.22.1.
  • Python and C: librdkafka and confluent-kafka-python 2.15.0 (June 30, 2026) ship the share consumer as Preview, with an explicit note that it "should not be used in production environments" (confluent-kafka-python v2.15.0, librdkafka v2.15.0). Confluent's own GA post says non-Java client support is "targeted for the second half of 2026" (Confluent blog).

If your workers are Python on the mainstream Confluent client, Kafka queues are explicitly not for production today. That alone settles the question for a lot of data teams.

The Postgres queue is one SQL statement

Here is the entire claim step:

WITH next AS (
    SELECT id FROM jobs
    ORDER BY id
    LIMIT 50
    FOR UPDATE SKIP LOCKED
)
DELETE FROM jobs j USING next
WHERE j.id = next.id
RETURNING j.id, j.payload;
Enter fullscreen mode Exit fullscreen mode

Run it inside a transaction, do the work, write the result, COMMIT. The commit is the ack. If the worker dies before committing, Postgres rolls the transaction back, the rows were never deleted, and another worker picks them up. That is the same contract as a share group's acquisition lock, except the "lock timeout" is "your connection died."

The Postgres docs are blunt about what SKIP LOCKED is for: "Skipping locked rows provides an inconsistent view of the data, so this is not suitable for general purpose work, but can be used to avoid lock contention with multiple consumers accessing a queue-like table" (SELECT docs).

Demo A: does SKIP LOCKED actually matter?

Eight worker threads drain 20,000 jobs in batches of 50. Each worker writes a row into a done table in the same transaction as the dequeue. The script asserts zero duplicates and zero lost jobs after every run. Each script execution does three alternating runs per cell and reports the median; the table shows the range of those medians across four executions:

Simulated work per batch (lock held) FOR UPDATE SKIP LOCKED plain FOR UPDATE
0 ms 14,400 to 17,600 jobs/s 13,200 to 32,800 jobs/s (about 32,000 in 3 of 4)
5 ms 14,000 to 17,600 jobs/s 7,100 to 7,200 jobs/s

Duplicates: 0. Lost: 0. In every run, for both variants.

The first row is the surprise: with near-zero lock hold time, plain FOR UPDATE was usually about twice as fast. A plausible explanation (interpretation, not measured): each SKIP LOCKED claim steps over the rows the other seven workers hold, while a blocked worker waits very briefly and gets its rows. The second row is the one that matters in production. Once workers hold the claim while doing real work, the blocking version collapses to about 7,000 jobs/s because everyone queues behind the head of the table, and SKIP LOCKED comes out 2 to 2.5 times ahead. Both variants lost nothing. SKIP LOCKED is about throughput under real work, not correctness.

One subtle thing the script surfaced: with SKIP LOCKED, a worker can get zero rows while jobs still exist, because everything left is claimed by someone else. It happened between 2 and 36 times per run. Treat an empty claim as "back off and poll again," never as "the queue is empty, exit."

Demo B: kill a worker mid-batch

A session claims 10 of 100 jobs and never commits. While it holds them, other sessions can claim 90. After pg_terminate_backend() on that session (the database-side equivalent of kill -9 on the worker), all 100 are claimable again. No lease column, no reaper job, as long as jobs are short (more on that below).

The real advantage: the job and the data commit together

This is the part Kafka cannot match without extra machinery. If your API inserts a user and enqueues "send welcome email" in the same transaction, the job exists if and only if the user exists. River, a Go job queue backed by Postgres, makes this its core argument (River docs): with a separate queue store you either enqueue after commit and lose jobs on a crash, or enqueue before commit and have workers race a row they cannot see yet.

The mainstream already voted on this. Solid Queue "is configured by default in new Rails 8 applications" and "leverages the FOR UPDATE SKIP LOCKED clause" (Solid Queue README).

Where the Postgres queue dies: one old snapshot

A queue table is the worst-case MVCC workload: every row is inserted once and deleted once. Postgres does not remove a deleted row immediately; the routine vacuuming docs explain that "the row version must not be deleted while it is still potentially visible to other transactions." If anything holds an old snapshot, vacuum cannot clean, and the dead rows pile up right where your claim query starts scanning.

Demo C: churn 200,000 jobs, then measure

The script pushes 200,000 jobs through the queue, enqueues 5,000 fresh ones, runs VACUUM (VERBOSE), and times 200 claims of 10 jobs each. Then it repeats the whole thing with one extra session sitting idle in transaction in REPEATABLE READ, which is exactly what a forgotten report query or a paused migration looks like.

Median claim latency VACUUM VERBOSE says
No old snapshot 0.269 ms / 0.272 ms 40,293 / 37,680 removed, 0 dead but not yet removable
One idle REPEATABLE READ session 15.754 ms / 14.961 ms 0 removed, 200,000 dead but not yet removable

(Two executions of the script, shown as first / second. In the no-snapshot case autovacuum had already cleaned most of the dead rows before the manual VACUUM ran.)

Same table, same query, same hardware. Roughly 55 to 58 times slower per claim, and the dead rows keep accumulating with every job processed while that session lives. Nothing errors. Your queue just gets slower until someone notices.

Things that hold old snapshots in real systems, all documented:

  • Long-running or idle-in-transaction sessions (vacuum docs tell you to find them via pg_stat_activity and backend_xmin).
  • Old replication slots (same page).
  • A hot standby with hot_standby_feedback on, which "can cause database bloat on the primary for some workloads" (replication config docs).

Mitigations that actually work:

  1. Set idle_in_transaction_session_timeout for application and human roles, and transaction_timeout (added in PostgreSQL 17) where it is safe. Both are documented in client connection defaults, including the warning to be careful with connection poolers.
  2. Alert on "dead but not yet removable" and on the age of the oldest backend_xmin, not just on queue depth.
  3. Keep analytics off the primary. An extra test on the same box: an old snapshot held in a different database on the same cluster did not block cleanup of the queue table (50,000 removed), while the same snapshot in the same database blocked all of it (50,000 dead but not yet removable). That is a real escape hatch, but it has a price: a queue in its own database can no longer commit atomically with your application tables, which was the main reason to use Postgres in the first place.

The trap: your worker can be the long transaction

The claim pattern above holds the transaction open while the job runs. For a 50 ms job that is fine. For a 5 minute job, the worker itself is the old transaction. A quick check on the same box: a worker sitting idle in transaction after its DELETE ... RETURNING (plain READ COMMITTED, nothing exotic) left 50,000 freshly deleted rows in another table of the same database "dead but not yet removable." So long jobs need a lease instead: claim and commit in a short transaction, work outside any transaction, then ack only if you still own the lease.

-- claim: short transaction, commit right away
UPDATE jobs
SET locked_until = now() + interval '5 minutes',
    lease_token  = gen_random_uuid(),
    attempts     = attempts + 1
WHERE id IN (
    SELECT id FROM jobs
    WHERE locked_until IS NULL OR locked_until < now()
    ORDER BY id
    LIMIT 10
    FOR UPDATE SKIP LOCKED
)
RETURNING id, payload, lease_token, attempts;

-- ack: only the current lease holder can delete the job
DELETE FROM jobs WHERE id = $1 AND lease_token = $2;
Enter fullscreen mode Exit fullscreen mode

Tested on the box: when a lease expired and a second worker re-claimed the job, the first worker's late ack deleted 0 rows and the second worker's ack deleted 1. Notice what you just built: an acquisition lock with a timeout and a delivery count. That is a share group, in SQL. You also gave up the "job commits with the data" property for the work itself, so this path needs idempotent jobs, exactly like Kafka.

A note on LISTEN/NOTIFY

If you use LISTEN/NOTIFY to wake workers, know the documented edges (NOTIFY docs): notifications are delivered only when the sending transaction commits, payloads must be shorter than 8000 bytes by default, and the notification queue (8 GB in a standard install) cannot be cleaned while a listening session sits in a long transaction. If that queue fills, "transactions calling NOTIFY will fail at commit." Use NOTIFY as a doorbell, never as the queue itself. PostgreSQL 19, with GA scheduled for October 29, 2026 (release dates), lists an improvement to "only wake up backends that are listening to specified notifications" in its draft release notes.

Where Kafka share groups bite

These are less about slow decay and more about semantics you have to design around.

1. The lock timeout is a duplicate-work timer. The default acquisition lock is 30 seconds. If your job takes 45, the broker hands the record to another consumer while the first one is still working. Use RENEW for long jobs, and make processing idempotent anyway. The franz-go kgo docs put it plainly: "Share group consumers are at-least-once regardless: any TCP hiccup can cause re-delivery, so user processing must be idempotent."

2. Ordering is gone. KIP-932: "The records in a share-partition can be delivered out of order to a consumer, in particular when redeliveries occur." If you need per-key ordering, you want a classic consumer group, not a share group.

3. A new share group starts at the end. share.auto.offset.reset defaults to latest (group configs). Create a share group on a topic that already has a backlog and, unless you set earliest, your workers will skip all of it.

4. Size-based retention can delete work you have not done. From the KIP: retention based on size "does potentially silently remove records that were eligible for delivery." A Postgres table does not delete your unprocessed jobs because the disk got full.

5. The unit of sharing is the record batch. The KIP says it directly: "the natural unit of sharing is the record batch." In the Go demo below, four consumers share one partition. With default fetch settings, one consumer got 949 of 1,000 jobs and two got nothing. With ShareMaxRecordsStrict() (the client side of the strict mode added by KIP-1206 in 4.2), the split was 246 / 251 / 252 / 251. Caveat: this ran against kfake, franz-go's in-process fake broker, not a real cluster, so treat the exact split as illustrative. The behavior matches what the KIP describes.

6. No atomicity with your database. Writing a row to Postgres and producing to Kafka are two systems. The standard fix is the transactional outbox: write the event to an outbox table in the same transaction, then let CDC (for example the Debezium outbox event router) publish it. Notice what that means: the reliable way to feed a Kafka queue from an OLTP app starts with a Postgres table.

7. In-flight work is capped per partition. group.share.partition.max.record.locks defaults to 2000 per share-partition (4.3 broker configs). Slow consumers holding locks can stall fetching on that partition until locks are acked or expire. And the delivery count "cannot be relied upon to be precise in all situations" (KIP-932), so do not build billing logic on it.

The Go side: a share group on one partition

The core loop, using franz-go v1.22.1:

cl, err := kgo.NewClient(
    kgo.SeedBrokers(addrs...),
    kgo.ConsumeTopics("jobs"),
    kgo.ShareGroup("email-workers"),
    kgo.ShareMaxRecords(25),
    kgo.ShareMaxRecordsStrict(), // without this one member may get a whole record batch
)
if err != nil {
    panic(err)
}
defer cl.Close()

for {
    pollCtx, cancel := context.WithTimeout(ctx, time.Second)
    fetches := cl.PollFetches(pollCtx)
    cancel()
    fetches.EachRecord(func(r *kgo.Record) {
        if err := sendEmail(r.Value); err != nil {
            r.Ack(kgo.AckRelease) // transient: try again, delivery count goes up
            return
        }
        r.Ack(kgo.AckAccept)
    })
}
Enter fullscreen mode Exit fullscreen mode

In the full demo, every job whose number ends in 0 is released on its first delivery to simulate a transient failure. Three runs in a row: 1,000 of 1,000 jobs accepted, 100 on a second delivery, 0 accepted twice, four consumers on one partition. The headline feature works.

The decision rule

Put it on one line first:

Tasks that belong to a transaction go in the database that owns the transaction. Facts that many systems read go in a log. Share groups let the log also hand out tasks, which is great once you already have the log.

Then the checklist.

Choose a Postgres queue when:

  • The job is created in the same transaction as the business data it depends on.
  • Your workers are in Python, or any language without a production share consumer yet.
  • Your measured peak fits your Postgres with your payloads. Demo A did 14,000 to 17,600 tiny jobs/s with SKIP LOCKED on 8 vCPUs; River claims "tens of thousands of jobs per second." Measure, do not inherit someone's number.
  • Jobs are short enough to run inside the claim transaction, or you use leases for the long ones.
  • You can enforce idle_in_transaction_session_timeout, keep analytics off the primary, and alert on dead-but-not-removable rows.

Choose Kafka share groups when:

  • Kafka is already the system of record for these events, and the queue is one more reader.
  • Several independent consumers need the same stream, or you need replay and reset-to-timestamp.
  • Your workers are on Java or Go (franz-go), processing is idempotent, and ordering does not matter.
  • You are fine designing around lock timeouts, retention, and share.auto.offset.reset.

Use both when the job starts in your database but the work fans out to other teams: outbox table in Postgres, CDC into Kafka, share groups downstream.

Red flags that you picked wrong: you are adding an outbox just to feed a Kafka queue that only one service reads (you wanted Postgres), or your Postgres queue needs a second database to stay fast (you are halfway to wanting a log).

Run it yourself

Hardware and versions for every number above: an 8 vCPU, 15 GB RAM Linux VM (kernel 6.12, Debian 13), PostgreSQL 17.11 with default settings except max_connections=200, client and server on the same machine over a Unix socket, Python 3.13.5 with psycopg 3.3.6, Go 1.26.0 with franz-go v1.22.1 and kfake. These are toy payloads on one box. They show mechanisms and ratios, not capacity for your system.

pg_queue_demo.py (Python, needs a throwaway Postgres)
"""
Postgres-as-a-queue demo: SKIP LOCKED vs blocking FOR UPDATE, crash recovery, and the MVCC failure mode.
Needs a THROWAWAY Postgres (it drops tables named jobs and done).
  pip install "psycopg[binary]"
  DSN="host=/tmp port=54329 user=postgres dbname=postgres" python pg_queue_demo.py
"""
import os, time, threading, statistics
import psycopg

DSN = os.environ.get("DSN", "host=/tmp port=54329 user=postgres dbname=postgres")

CLAIM_SKIP = """
WITH next AS (
    SELECT id FROM jobs
    ORDER BY id
    LIMIT %(n)s
    FOR UPDATE SKIP LOCKED
)
DELETE FROM jobs j USING next
WHERE j.id = next.id
RETURNING j.id, j.payload
"""
CLAIM_BLOCKING = CLAIM_SKIP.replace("FOR UPDATE SKIP LOCKED", "FOR UPDATE")


def setup(conn):
    conn.execute("DROP TABLE IF EXISTS jobs, done")
    conn.execute("CREATE TABLE jobs (id bigserial PRIMARY KEY, payload jsonb NOT NULL)")
    # 'done' stands in for a side effect that lives in the SAME database as the queue
    conn.execute("CREATE TABLE done (job_id bigint NOT NULL, worker int NOT NULL)")
    conn.commit()


def enqueue(conn, n):
    conn.execute("INSERT INTO jobs (payload) SELECT jsonb_build_object('n', g) "
                 "FROM generate_series(1, %s) g", (n,))
    conn.commit()


def worker(wid, claim_sql, batch, stats, empties, work_ms):
    with psycopg.connect(DSN) as conn:
        while True:
            with conn.transaction():
                rows = conn.execute(claim_sql, {"n": batch}).fetchall()
                if rows:
                    if work_ms:
                        time.sleep(work_ms / 1000)  # "work" done while holding the claim
                    with conn.cursor() as cur:      # result commits atomically with the dequeue
                        cur.executemany("INSERT INTO done (job_id, worker) VALUES (%s, %s)",
                                        [(r[0], wid) for r in rows])
            if rows:
                stats[wid] += len(rows)
                continue
            # An empty claim does not always mean an empty queue (see the blocking variant).
            left = conn.execute("SELECT count(*) FROM jobs").fetchone()[0]
            conn.commit()
            if left == 0:
                return
            empties[wid] += 1


def run_workers(claim_sql, workers, batch, work_ms=0.0):
    stats, empties = [0] * workers, [0] * workers
    ts = [threading.Thread(target=worker, args=(i, claim_sql, batch, stats, empties, work_ms))
          for i in range(workers)]
    t0 = time.perf_counter()
    for t in ts:
        t.start()
    for t in ts:
        t.join()
    return time.perf_counter() - t0, sum(empties)


def experiment_throughput(n=20_000, workers=8, batch=50, repeats=3):
    for work_ms in (0.0, 5.0):
        print(f"\n== A. {workers} workers drain {n} jobs, batch={batch}, "
              f"simulated work per batch={work_ms} ms, {repeats} alternating runs")
        results = {"FOR UPDATE SKIP LOCKED": [], "FOR UPDATE (blocking)": []}
        for _ in range(repeats):
            for label, sql in (("FOR UPDATE SKIP LOCKED", CLAIM_SKIP),
                               ("FOR UPDATE (blocking)", CLAIM_BLOCKING)):
                with psycopg.connect(DSN) as conn:
                    setup(conn)
                    enqueue(conn, n)
                secs, empties = run_workers(sql, workers, batch, work_ms)
                with psycopg.connect(DSN) as conn:
                    total, distinct = conn.execute(
                        "SELECT count(*), count(DISTINCT job_id) FROM done").fetchone()
                    left = conn.execute("SELECT count(*) FROM jobs").fetchone()[0]
                assert total == distinct == n and left == 0, (total, distinct, left)
                results[label].append((n / secs, empties))
        for label, rs in results.items():
            print(f"{label:24s} median {statistics.median(r[0] for r in rs):7.0f} jobs/s | "
                  f"empty claims while work remained, per run: {[r[1] for r in rs]} | "
                  f"duplicates 0, lost 0")


def experiment_crash(n=100):
    print("\n== B. a worker dies mid-transaction")
    with psycopg.connect(DSN) as conn:
        setup(conn)
        enqueue(conn, n)
    victim = psycopg.connect(DSN)
    victim.execute(CLAIM_SKIP, {"n": 10}).fetchall()  # claimed 10 jobs, never commits
    with psycopg.connect(DSN, autocommit=True) as admin:
        q = "SELECT count(*) FROM (SELECT id FROM jobs FOR UPDATE SKIP LOCKED) s"
        print(f"victim holds its batch: claimable by others = {admin.execute(q).fetchone()[0]}")
        admin.execute("SELECT pg_terminate_backend(%s)", (victim.info.backend_pid,))
        time.sleep(0.5)
        print(f"victim session killed:  claimable by others = {admin.execute(q).fetchone()[0]}")
    try:
        victim.close()
    except Exception:
        pass


def median_claim_ms(rounds=200):
    lat = []
    with psycopg.connect(DSN) as conn:
        for _ in range(rounds):
            t0 = time.perf_counter()
            with conn.transaction():
                conn.execute(CLAIM_SKIP, {"n": 10}).fetchall()
            lat.append((time.perf_counter() - t0) * 1000)
    return statistics.median(lat)


def experiment_mvcc(churn=200_000):
    print(f"\n== C. churn {churn} jobs through the queue, then time claims, "
          f"with and without one old snapshot held open")
    for hold in (False, True):
        with psycopg.connect(DSN) as conn:
            setup(conn)
        holder = None
        if hold:
            # stands in for a long report query, a stuck migration, or an idle-in-transaction session
            holder = psycopg.connect(DSN, autocommit=True)
            holder.execute("BEGIN ISOLATION LEVEL REPEATABLE READ")
            holder.execute("SELECT 1").fetchone()  # first statement takes the snapshot
            with psycopg.connect(DSN, autocommit=True) as c2:
                st, xmin = c2.execute("SELECT state, backend_xmin FROM pg_stat_activity "
                                      "WHERE pid = %s", (holder.info.backend_pid,)).fetchone()
            print(f"   holder session: state={st!r}, backend_xmin={xmin}")
        with psycopg.connect(DSN) as conn:
            enqueue(conn, churn)
        run_workers(CLAIM_SKIP, 8, 100)            # drain everything, leaving dead rows behind
        with psycopg.connect(DSN) as conn:
            enqueue(conn, 5_000)                   # fresh work arrives behind the dead rows
        notes = []
        with psycopg.connect(DSN, autocommit=True) as conn:
            conn.add_notice_handler(lambda d: notes.append(d.message_primary or ""))
            conn.execute("VACUUM (VERBOSE) jobs")  # give vacuum every chance
        tuples_line = next((l.strip() for n in notes for l in n.splitlines() if "tuples:" in l), "")
        med = median_claim_ms()
        with psycopg.connect(DSN) as conn:
            size = conn.execute("SELECT pg_size_pretty(pg_total_relation_size('jobs'))").fetchone()[0]
        print(f"old snapshot held={str(hold):5s} median claim latency={med:8.3f} ms  "
              f"jobs table+index={size}\n   VACUUM VERBOSE: {tuples_line}")
        if holder:
            holder.execute("ROLLBACK")
            holder.close()


if __name__ == "__main__":
    with psycopg.connect(DSN) as c:
        print(c.execute("SELECT version()").fetchone()[0])
    experiment_throughput()
    experiment_crash()
    experiment_mvcc()
Enter fullscreen mode Exit fullscreen mode

main.go (Go, franz-go share group against kfake)
// Share groups (KIP-932) in Go with franz-go: 4 consumers on ONE partition,
// per-record ack, and a transient failure that gets redelivered.
// Runs against kfake, franz-go's in-process fake broker, so no real cluster is needed.
package main

import (
    "context"
    "fmt"
    "strconv"
    "sync"
    "time"

    "github.com/twmb/franz-go/pkg/kfake"
    "github.com/twmb/franz-go/pkg/kgo"
)

const (
    topic   = "jobs"
    group   = "email-workers"
    total   = 1000
    members = 4
)

func main() {
    c, err := kfake.NewCluster(kfake.NumBrokers(1), kfake.SeedTopics(1, topic)) // ONE partition
    if err != nil {
        panic(err)
    }
    defer c.Close()
    // Default is "latest": a brand-new share group would skip everything already in the topic.
    c.SetGroupConfigs(group, map[string]string{"share.auto.offset.reset": "earliest"})

    producer, err := kgo.NewClient(kgo.SeedBrokers(c.ListenAddrs()...), kgo.DefaultProduceTopic(topic))
    if err != nil {
        panic(err)
    }
    for i := 0; i < total; i++ {
        producer.Produce(context.Background(), kgo.StringRecord(strconv.Itoa(i)), nil)
    }
    if err := producer.Flush(context.Background()); err != nil {
        panic(err)
    }
    producer.Close()

    var (
        mu          sync.Mutex
        accepted    = map[string]int{} // job -> times accepted
        perMember   = make([]int, members)
        redelivered int
    )
    ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
    defer cancel()

    var wg sync.WaitGroup
    for m := 0; m < members; m++ {
        wg.Add(1)
        go func(m int) {
            defer wg.Done()
            cl, err := kgo.NewClient(
                kgo.SeedBrokers(c.ListenAddrs()...),
                kgo.ConsumeTopics(topic),
                kgo.ShareGroup(group),
                kgo.ShareMaxRecords(25),
                kgo.ShareMaxRecordsStrict(), // without this the broker may hand one member a whole record batch
            )
            if err != nil {
                panic(err)
            }
            defer cl.Close()
            for {
                mu.Lock()
                done := len(accepted) == total
                mu.Unlock()
                if done || ctx.Err() != nil {
                    _ = cl.FlushAcks(context.Background())
                    return
                }
                pollCtx, pollCancel := context.WithTimeout(ctx, time.Second)
                fetches := cl.PollFetches(pollCtx)
                pollCancel()
                fetches.EachRecord(func(r *kgo.Record) {
                    n, _ := strconv.Atoi(string(r.Value))
                    if n%10 == 0 && r.DeliveryCount() == 1 {
                        r.Ack(kgo.AckRelease) // transient failure: give it back for another attempt
                        return
                    }
                    time.Sleep(2 * time.Millisecond) // pretend to send the email
                    r.Ack(kgo.AckAccept)
                    mu.Lock()
                    accepted[string(r.Value)]++
                    perMember[m]++
                    if r.DeliveryCount() > 1 {
                        redelivered++
                    }
                    mu.Unlock()
                })
            }
        }(m)
    }
    wg.Wait()

    dups := 0
    for _, k := range accepted {
        if k > 1 {
            dups += k - 1
        }
    }
    fmt.Printf("partitions=1 consumers=%d distinct jobs accepted=%d/%d\n", members, len(accepted), total)
    fmt.Printf("accepted per consumer=%v\n", perMember)
    fmt.Printf("accepted on a 2nd+ delivery (after RELEASE)=%d, accepted more than once=%d\n", redelivered, dups)
}
Enter fullscreen mode Exit fullscreen mode

For the Go demo: go mod init sharedemo && go get github.com/twmb/franz-go@v1.22.1 github.com/twmb/franz-go/pkg/kfake@v0.0.0-20260927204940-b5a45ccfdf7e && go run . (the kfake module requires Go 1.26, which the go command downloads automatically if your toolchain is older).

Takeaway

Kafka 4.2 fixed the strongest technical objection to "Kafka as a queue." It did not change where your jobs are born. If they are born in a Postgres transaction, keep them there, and spend your effort on the one thing that kills Postgres queues: old snapshots. When the events already live in Kafka and many teams read them, share groups are now a legitimate queue, on Java and Go today, and on Python once the client leaves Preview.

Sources

Top comments (0)