feat(deps): update module github.com/twmb/franz-go ( v1.21.1 → v1.22.1 ) #3

Open
renovatebot wants to merge 1 commit from renovate/github.com-twmb-franz-go-1.x into main
Member

This PR contains the following updates:

Package Change Age Confidence
github.com/twmb/franz-go v1.21.1 → v1.22.1 age confidence

Release Notes

twmb/franz-go (github.com/twmb/franz-go)

v1.22.1

Compare Source

===

This patch has a few bug fixes in the share consumer (one via bug report,
the others via a corresponding targeted audit), the 848 consumer (these found
only via an audit), and some fixes in the RecordFormatter and RecordReader.

This patch also has produce and consume performance improvements that come
with two minor behavior changes:

  • MaxBufferedRecords now has a default of 50K, up from 10K: 10K was chosen when
    I initially wrote this library and is a very low default limit for average sized
    records. The librdkafka default is 100K; 50K increases producer throughput while
    still keeping producer memory low.

  • PoolKRecords is deprecated and unused. The client now decodes fetched
    records straight into Records, so there is no need for pooling kmsg.Records.

Thanks to @​kmrgirish for the share consumer
bug report (#​1474), and to
@​ajavanma and @​jakezwang
for RecordReader and RecordFormatter fixes.

v1.22.0

Compare Source

===

This release supports Kafka 4.3 and 4.4, has a few new APIs, and has a few
big internal improvements. In particular, I recommend checking out the new
StreamingCompression option, as well as evaluating if you'd like to use
RackAwarePartitioning. There are some behavior changes that you should read
about below. The "next gen" rebalancer is now usable via the new
ServerSideBalancer option. It's had a few releases to shake out bugs
internally (via integration tests and LLM audits), but if you do experience a
bug, please open an issue straightaway.

Some minor bug fixes (that were never reported) were found during the
implementation that are not worth mentioning.

kfake has also been significantly extended and I recommend checking out the
new APIs, in particular:

  • A new Fault type to make it easier to inject errors without Control functions
  • Group introspection cluster APIs
  • BlackholeProduce and SyntheticFetch APIs for benchmarking / play testing

My kcl CLI has been significantly expanded as well and is worth checking
out. It supports essentially everything you can do with a cluster, and now
allows you to run a full broker locally via kcl fake (in memory or a dumb
disk backed localhost broker) - as well as setup the fake broker with fault
injection. I've been running LLM audits and extensions to kcl in particular
to try to shape it up to a "finalized" CLI shape. If you use it and have ideas
for improvements, please open an issue.

Behavior changes

  • Rack aware group partition assignment (KIP-881) now requires BalanceRacks.
    v1.21.0 enabled group balancers to assign partitions based on the rack that
    members were in if you used the range or sticky/cooperative-sticky balancers.
    Well, Rack is also used to opt into preferred read replica assignment
    when fetching by the broker itself. These two decisions conflict with each
    other. Now, BalanceRacks() is required to opt into group balancers using
    the rack while balancing. The client warns when balancing if BalanceRacks
    is on and the brokers have preferred read replicas enabled.

  • ConsumeResetOffset defaults to RewindOffset(time.Minute) rather
    than NewOffset().AtStart(). Setting only ConsumeStartOffset no longer
    sets ConsumeResetOffset
    . I introduced ConsumeStartOffset a while back
    because it was really weird IMO to use a reset offset for both how a consumer
    starts and for how it recovers in the event of data loss or falling behind.
    They were bidirectional since introduction, but since start is newer and much
    less commonly used and you often don't want to recover from the start, I've
    removed the start -> reset mapping when you only set the start. I recommend
    reading the docs on both options for an updated understanding of when and
    how they apply. As well, I've introduced RewindOffset(d) which is only
    relevant to the reset offset (rewind by d duration from the last consumed
    offset on data loss we cannot exactly recover from) and LookbackOffset(d)
    which is relevant to both options but more useful for the start offset
    (start consuming d before the newest record; before Kafka 3.0 it is d
    before the current time). If a committed offset has fallen below the log
    start, the first fetch answers OFFSET_OUT_OF_RANGE and the reset offset
    decides where to resume. Before, a start offset of AtEnd was copied into
    the reset offset, so the consumer skipped to the end. Now, with the
    defaults, it resumes at the log start.

  • Topic recreation is now a hard failure. The client always
    produces to and consumes from the first instance of a topic. If you delete
    and recreate a topic, the client refuses the new version: buffered records
    fail with UNKNOWN_TOPIC_ID, fetches stop, offsets from the old topic cannot
    be committed to the new one, and transactions on the old topic fail. This
    needs a broker that reports topic IDs (Kafka 2.8+). Previously, some things
    in the client continued to accidentally work, and the behavior was
    unreliable and usually not good. If you want your application to stay alive
    across topic recreations, you can PurgeTopicsFromClient and, for
    consumers, AddConsumeTopics. More details about topic recreation are now in
    a new section in the README.

  • MaxDecompressBatchBytes now blocks decompression if a batch would
    decompress too large (default 1GiB)
    . Fetches when consuming can only
    specify to the broker "give me X bytes of batches", but they cannot control
    how large those batches decompress into. A hostile or buggy batch could OOM
    your program. Now, a batch over the limit causes the partition to enter
    a fatal state and return ErrDecompressTooLarge once from polling.
    The application can recover by manually skipping the batch with SetOffsets
    (with the fields in the error; see the docs), or by restarting the client
    with a higher limit. This option does not apply to custom decompressors,
    but, custom decompressors can still return ErrMaxDecompress to stop
    the partition. This option is also closely related to streaming compression,
    which is described below.

Improvements

  • gzip now uses klauspost/compress (same format). Its default level is
    1.7x faster than stdlib's with a slightly better ratio; klauspost's default
    maps to its level 5 where stdlib's mapped to 6, and level for level it is
    1.1x to 1.2x faster. WithLevel(n) now selects klauspost's level n, so
    the bytes a given level produces differ from before.

  • The sticky balancers are now exactly optimal on balance, then rack
    placement (with BalanceRacks), then stickiness. Balancing was already
    load optimal but had some very niche edge cases where maximal stickiness
    was not preserved, especially if balancing used racks. Rack placement
    outranks stickiness: turning BalanceRacks on in a running group
    reassigns, at its next rebalance, every partition held by a member in a
    different zone from the partition's leader.

  • Sticky balancing is much faster, most of all on rejoins and on groups whose
    members subscribe to different topics. Against v1.21.7: a rejoin of 100
    members over 1600 topics of 100 partitions goes from 351ms to 24ms; a regex
    shaped group of 500 members over 20,000 topics from 176ms and 810MB to 12ms
    and 9MB; 2001 members over 500 topics of 2000 partitions with one narrow
    subscriber from 3.4s to 0.3s. Fresh uniform balances are unchanged.

Features

Streaming compression

StreamingCompression is an opt-in producer option that compresses a
partition's backlog of batches together, bounded by their compressed size.
By default a batch is cut at ProducerBatchMaxBytes measured on uncompressed
records. Streaming compression will help reduce traffic to the broker and
increase how effective compression actually is (by pulling more data in at
once). A custom compressor makes this option a no-op.

The client is implemented such that each compression codec's worst case
overhead is tracked internally, which should avoid a compressed batch ever
exceeding ProducerBatchMaxBytes. If this ever does happen, the client
discards the merge, logs a warning, disables streaming compression for the
client going forward (records are still compressed batch by batch), and asks
you to file an issue.

The client has a new option MaxDecompressBatchBytes to bound both (a) how
much the producer can stuff into a merged batch (i.e. how much it will
decompress into), and (b) the maximum size a consumer will decompress a batch
to; the consumer never decompresses past the bound (preventing a zip bomb).
The default is 1GiB.

Rack aware producer partitioning (KIP-1123)

RackAwarePartitioning sends unkeyed records to partitions whose leader is in
the client's Rack (which must also be set), falling back to all partitions
when no leader is. Keyed records are never affected. Unlike the Java client,
this works with any partitioner, since the eligible-broker filtering happens
before your partitioner is consulted. Note that this option skews which
partitions receive records if your producers are not spread across racks in
proportion to partition leaders.

ServerSideBalancer (KIP-848)

ServerSideBalancer opts into KIP-848 "next-gen" consumer groups, where the
broker's group coordinator assigns partitions rather than the client. This
requires Kafka 4.0+ and either a range or sticky / cooperative-sticky
balancer. This replaces the hidden opt_in_kafka_next_gen_balancer_beta
context key from v1.19.0; the key still works in this release but will be removed
in the next. The default remains the classic protocol, matching the Java
client. I still think the classic client side balancers are better (and this
client's implementation is way faster than the Java client), but if you want
to use server side balancing, it is strongly recommended to only use it if
your cluster is Kafka 4.3+. Before 4.3 (before KIP-1251), an offset commit
that races with a heartbeat epoch bump can fail with STALE_MEMBER_EPOCH,
which the client cannot detect nor handle.

BalanceInfo for custom balancers

A balancer that implements GroupMemberBalancerInfo receives a BalanceInfo
before balancing: the group, generation, leader member ID, and lazily built
topic and broker metadata. ConsumerBalancer implements it, so balancers
built on NewConsumerBalancer can call Info(). This allows, for example, a
balancer that assigns every partition to the leader with the other members as
hot standbys. Thanks @​michaelwilner!

API additions

// Producing
func StreamingCompression() ProducerOpt
func RackAwarePartitioning() ProducerOpt

// Consuming
func BalanceRacks() ConsumerOpt
func ServerSideBalancer() GroupOpt
func RewindOffset(d time.Duration) Offset
func LookbackOffset(d time.Duration) Offset

// Decompression bound
func MaxDecompressBatchBytes(n int) Opt
var ErrMaxDecompress error
type ErrDecompressTooLarge struct {
    Topic      string
    Partition  int32
    Offset     int64
    Epoch      int32
    NextOffset int64
}

// Custom balancers
type BalanceInfo struct {
    Group      string
    Generation int32
    LeaderID   string
    Topics     func() map[string]TopicMetadata
    Brokers    func() map[int32]BrokerMetadata
}
type GroupMemberBalancerInfo interface {
    GroupMemberBalancer
    SetBalanceInfo(BalanceInfo)
}
func (*ConsumerBalancer) Info() BalanceInfo
type TopicMetadata struct { ... }
type PartitionMetadata struct { ... }

// Records
type RecordAttrsOpts struct {
    Codec         CompressionCodecType
    TimestampType int8
    Transactional bool
    Control       bool
}
func NewRecordAttrs(RecordAttrsOpts) RecordAttrs

// kversion
func (*Versions) EachSupportedFeature(fn func(name string, min, max int16))
func (*Versions) EachFinalizedFeature(fn func(name string, level int16))
func FeatureLevelDescription(name string, level int16) string

Relevant commits

There are many commits, but some of the more notable ones:

  • 27d11286 feature kversion: FeatureLevelDescription
  • 73358f62 feature kversion: supported and finalized feature levels per release
  • 7be0be16 behavior change kgo: add MaxDecompressedBatchBytes
  • b37f1041 feature kgo: add ServerSideBalancer to opt into KIP-848
  • 033a46c7 improvement kgo: begin ApiVersions at the max a broker told us, for an hour
  • 46a9b2ad behavior change kgo: use ConsumeResetOffset when the broker loses data we cannot locate
  • 8b33e43d improvement kgo: speed up compression on both the legacy and the merge path
  • 9de0fa36 feature kgo: add StreamingCompression, compressed-size-bound batch merging
  • de7327e6 feature kgo: detect misrouted connections (KIP-1242)
  • 8ad36ec7 feature kgo: support TxnOffsetCommit v6
  • e4f7bc43 feature kgo: add rack-aware producer partitioning (KIP-1123)
  • 123f2ffa improvement kgo: repair the sticky plan to the best balance, rack, and stickiness
  • 23ab9a0e behavior change kgo: add BalanceRacks, gate rack aware balancing behind it
  • d4f6db2f improvement kgo: drop reassigned partitions in one pass in AdjustCooperative
  • 35efafc8 behavior change kgo: fail records for a recreated topic instead of producing by name
  • 4f10346a feature kgo: expose BalanceInfo for custom balancer implementations (thanks @​michaelwilner!)
  • cd7f9b4e feature kgo: add NewRecordAttrs constructor (thanks @​pracucci!)

v1.21.7

Compare Source

===

A handful of bug fixes and improvements found by users and while working on
v1.22. Rather than enumerating the relevant commits, you can check the git log
between v1.21.6 and this release - there are many minor commits. As well, kfake
has been improved significantly and has more API surface to aid in writing
tests.

  • A rare, very niche panic while producing has been fixed. Thanks
    @​PumpkinDemo for the report, see
    #​1385 for more details.

  • If retention deleted the segment a consumer was reading, the
    OffsetOutOfRange reset listed by the last consumed timestamp and could skip
    surviving records, or jump to the log end and skip everything. The reset
    now resumes at the log start when below it.

  • Improved KIP-951 handling (the broker returning where a partition should move
    with the produce response if the partition changed leadership). Previously, a
    broker could return NotLeaderForPartition and hint the leader the client
    was already using, at the same or an older epoch. These hints are now
    ignored and the client backs off, rather than spinning. Thanks
    @​3AceShowHand for the report and
    @​jjj-n for a fix, see
    #​1412.

  • Rack aware balancers ignored rack for any topic the group leader did not
    itself consume. A client now loads the rack for all partitions in the
    group, even if the leader does not consume some of the topics.

  • The metadata cache has been improved (there were a few cases where it was
    emptied erroneously).

  • Regex consuming now consistently never matches internal topics such as
    __consumer_offsets. As well, the regex log no longer reports an excluded
    topic as both added and skipped (thanks @​lahsivjar).

  • Decompression allocates less (thanks @​scunningham).
    If you use pools, slices are now reliably returned if decompression errors.

  • The client now starts at a random seed broker rather than always the
    first, so many clients starting at once no longer all hit the same seed.
    Thanks @​chailuecha!

  • A producer that receives RequestTimedOut or NotEnoughReplicasAfterAppend now
    retries after the produce backoff rather than waiting for a metadata refresh.

  • A few other minor improvements and bug fixes.

v1.21.6

Compare Source

===

Some bug fixes (mostly minor - hence the delay for the release) found by users
and further Claude audits. I am gearing up for a 1.22 release but some of the
features I am planning for are more complicated to review, so it may take a bit
of time. Anyway:

  • Previously, rollback from a cooperative group to an eager group was
    deliberately not supported and there was a data race condition if this
    happened. It is now technically supported, although you will experience
    duplicate data. If you want a safe non-duplicate-causing rollback, you need
    to turn off the entire group, remove the cooperative consumer, and swap the
    whole group to eager rebalancing.

  • Fixed a panic: close of closed channel on an acks=0 produce connection
    in a specific edge case (a broker connection dying before the connection
    was fully established caused the panic).

  • If EndTransaction failed with an unconfirmed outcome (a transport error,
    exhausted retries, or UNKNOWN_SERVER_ERROR), the documented abort retry
    was a wire no-op and the next transaction could silently commit the prior
    "failed" transaction's records under KIP-890 part 2. The producer ID is now
    flagged for reload, which fence-aborts anything still ongoing broker-side.

  • GroupTransactSession.End could hang forever, ignoring its context, if the
    group had never joined (e.g. the consumed topic did not exist yet) and the
    transaction committed no offsets.

  • Previously, if a broker replied to ApiVersions with an error, we ignored it
    and you would eventually see an unclear error (usually a bare io.EOF, since
    anything that rejects ApiVersions hangs up right after replying). These
    errors are now handled correctly.

  • Some niche edge case bugs that are only worth reading about if you're super
    interested were found in repeated Claude audits and were fixed (check the PR
    / git history). This includes further KIP-848 "next gen consumer group" fixes.

Relevant commits

  • 582e0f21 bugfix kgo: surface error codes in ApiVersions responses
  • 67ef4c61 bugfix kgo: fix double close of a connection's deadCh on acks=0 produce
  • 3ac2fff1 bugfix kgo: revoke everything when the group protocol downgrades from cooperative to eager
  • 795d5b61 improvement kgo: flatten topic/partition maps in group rebalance logs (thanks @​constanca-m!)
  • 70addc1e improvement kgo: classify retired broker reads as broker dead (thanks @​tomplarge!)
  • 18f9a10f improvement deps: replace golang.org/x/crypto/pbkdf2 with stdlib crypto/pbkdf2 (thanks @​macdewee!)
  • 6ecd2f9f bugfix kgo: recover when an attempted EndTxn outcome is unconfirmed
  • 821f879e bugfix kgo: fix GroupTransactSession.End hanging when the group never joined

v1.21.5

Compare Source

===

Three bug fixes:

  • Fixed a nil-pointer panic when building a group OffsetFetch: if the group
    was assigned a topic that was no longer in the client's tracked set --
    reachable when a topic is purged from consuming while still assigned, for
    example PurgeTopicsFromConsuming overlapping AddConsumeTopics, or the
    automatic regex missing-topic purge -- loadTopic returned nil and
    dereferencing it for the topic ID crashed the client. The topic ID is now
    only set when the topic is known. Thanks @​iwittkau!

  • Fixed a data race on a coordinator's cached node ID. When a broker
    disconnected while a FindCoordinator load for that broker was still in
    flight, deleteStaleCoordinatorsByNode could read the in-flight load's
    node field before the loading goroutine published it (via closing the
    load's wait channel), which go test -race flagged. The node read now
    happens only after the load has been observed as complete. Thanks
    @​nikolauspschuetz!

  • A share partition that was listed in a ShareFetch only to carry a
    piggybacked acknowledgement -- for a cursor that was revoked, paused, or
    migrated to a new leader after its records were drained -- was added to the
    broker's share session but never tracked client-side, so it could never be
    forgotten. The broker would re-acquire and redeliver that partition's
    records indefinitely while the client discarded them ("broker returned
    partition ... we did not ask for"), spinning the share fetch loop. The
    client now tracks every partition it sends, matching the broker's session
    bookkeeping.

Relevant commits

v1.21.4

Compare Source

===

This release is a "large" (many commits) release that has many small or
hard to encounter bugs fixed. I pointed Claude's Fable at this repo and
ran some audit rounds while available and thankfully got through the highest
value audit rounds before Fable was removed.

For once, I will not be describing every bug fixed nor calling out every
relevant commit. Instead, if you are curious, look at
#​1348. Some worthwhile
description is below.

Three important bug fixes to call out:

  • In transactional exactly-once consuming, a SetOffsets seek (which happens
    during GroupTransactSession.End after an aborted transaction) could be
    undone by a concurrent offset load (via a background list or epoch load) that
    completed slightly later. This could happen when the client discovers a
    partition leader moved while you are aborting, which could result in missed
    records.

  • Consuming with read_committed against a broker that returns a partition's
    aborted-transaction list out of offset order could surface aborted,
    rolled-back records as if they were committed. Apache Kafka always returns
    them in order so this was never observed there, but Redpanda does not (when
    an aborted transaction is still in memory and an earlier one is already on
    disk). The list is now sorted client-side, matching the Java client,
    librdkafka, and Sarama.

  • GroupTransactSession.End no longer reports a successful commit when the
    broker answers EndTxn with UNKNOWN_SERVER_ERROR (seen from Redpanda in
    some older versions). Previously the consumer's offsets were advanced past a
    transaction that may have aborted; now the commit is reported as failing
    and the session rewinds for reprocessing.

Beyond those, by area:

  • Many transaction-path fixes for coordinator churn and KIP-890 part 2 that
    would have resulted in not-working (hard client fail) or hung transactions:
    InitProducerID retries CONCURRENT_TRANSACTIONS when taking over a
    crashed producer's transaction, retriable producer-id load failures are no
    longer treated as fatal, KIP-890p2 is opted into only when the negotiated
    versions actually support it (fixing spurious INVALID_TXN_STATE on 4.0+
    clusters running older semantics), a transaction whose every produce failed
    now aborts instead of hanging until the transaction timeout, and a failed
    AddPartitionsToTxn no longer drops partitions added by an earlier request.

  • GzipCompression().WithLevel(...) was completely broken and would panic.

  • More KIP-848 (next-gen consumer group) robustness fixes under coordinator
    and leader churn.

  • Stale consumer-group member rejoining fixes: a member that rejoins claiming a
    partition at an old generation no longer panics the group leader or causes
    two members to consume the same partition. Malformed member metadata,
    duplicate member ids, and negative claimed partitions in a join are now
    rejected or sanitized rather than mis-balancing or panicking the leader.

  • Metadata and topic recreation: a stale per-broker metadata view that
    momentarily omits a just-added partition (the window right after
    CreatePartitions) no longer fails buffered producer records or leaves a
    newly assigned consumer / share partition silently unconsumed; both heal
    once metadata catches up.

  • Share consumer: fetch errors are now classified like the classic consumer
    (retriable errors stripped, a metadata refresh triggered to heal a leader
    move, top-level errors backed off) instead of stalling for up to
    MetadataMaxAge or hot-looping, and leader-move migrations are tracked so
    that leaving or closing cannot strand un-acked records.

  • SASL: KIP-368 re-authentication no longer races the connection's other
    reader, which could corrupt pipelined traffic on brokers that set a session
    lifetime (e.g. AWS MSK IAM); requests now park and replay across a re-auth.
    The Azure Event Hubs ApiVersions reset retry no longer leaks the abandoned
    connection or silently downgrades it to v0.

  • KIP-714 client telemetry: the terminating push is now actually delivered on
    Close, the .rate and .avg rollups are computed correctly (they were
    constant / wrong before), and an unsupported user-metric attribute no longer
    corrupts the OTLP payload (which had disabled metrics for the rest of the
    client's life).

  • Smaller consumer fixes: overlapping manual CommitOffsets no longer reopen
    autocommit early (which could rewind the committed offset), a conformant
    UNDEFINED_EPOCH_OFFSET epoch response no longer raises a false
    ErrDataLoss, and Fetches.EachTopic now preserves TopicID across
    multi-broker responses (it was zero whenever more than one broker replied,
    i.e. normally).

  • RecordReader / RecordFormatter no longer panic on truncated or malformed
    layouts, accept \xNN escapes for bytes above 0x7f, and reject layouts that
    would read nothing and loop forever.

  • WithPools: decompression no longer produces garbage when a pool hands back
    a non-zero-length sized slice, and pooled slices are no longer leaked for
    batches that keep no records (e.g. aborted-transaction data under
    read_committed).

  • Other producer fixes: EnsureProduceConnectionIsOpen dials the right broker
    for filtered ids and no longer breaks an acks=0 connection, the adaptive
    LeastBackupPartitioner now actually picks the least-backed-up partition,
    and producing during or after Close fails cleanly instead of hanging a
    later Flush.

  • A broad set of guards against malformed or hostile broker responses that
    could previously panic the fetcher, hot-loop, or mis-consume: negative or
    oversized record counts and batch lengths, decompression bombs, duplicate or
    omitted partitions, negative offsets, and unexpected top-level fetch errors.
    The client is also more resilient when its own API contracts are violated
    (e.g. AllowRebalance called while a poll is in flight).

v1.21.3

Compare Source

===

This patch release contains a few bug fixes and a few internal improvements.

  • PollRecords / PollFetches could permanently hang since v1.21.0 if a
    consumer session stopped (usually via metadata updates) while
    fetches to more than four brokers were pending and no poll was in
    flight. This could only affect users that deliberately set MaxConcurrentFetches(0),
    or that were using ShareMaxRecordsStrict.

  • Producing to a topic whose partitions ALL have a retriable load error
    (e.g. a rolling restart of an RF=1 broker briefly leaving every
    partition leaderless) no longer fails records up front with "unable to
    partition record due to no usable partitions". Instead, the records
    remain buffered and retried as metadata reloads.

  • Classic consumer groups now rejoin immediately when an offset commit
    returns UNKNOWN_MEMBER_ID or ILLEGAL_GENERATION (the broker lost
    the member, e.g. a session expired during a network blip), rather than
    consuming as a zombie until the heartbeat loop notices the dead session.

  • DescribeShareGroupOffsets, AlterShareGroupOffsets, and
    DeleteShareGroupOffsets are now routed to the group coordinator
    rather than the share coordinator (which would reject the requests
    for being misrouted).

  • The client-internal metadata cache now deeply clones the cached response
    before putting it into the cache and before returning it via
    RequestCachedMetadata (which is now used by default in kadm), eliminating
    data race possibilities.

  • Various next-gen rebalancer session improvements.

Relevant commits

  • f8842170 improvement kgo: fall back to all partitions when no partition is writable (thanks @​ericsg666!)
  • 824e34d2 improvement kgo: rejoin a classic group when a commit returns a fatal member error (thanks @​v14dis14v!)
  • f520e820 bugfix kgo: do not exit manageFetchConcurrency while sources are pending in wantFetch (thanks @​SLoeuillet!)
  • 8d9c836b bugfix kgo: isolate metadata cache from broker response
  • 19f7dbb2 bugfix kgo,kfake: route share group offset RPCs to the group coordinator

v1.21.2

Compare Source

===

This patch release contains two narrow bug fixes and one small feature.
Deps are also bumped so that you are force-pinned to a klauspost/compress
version that has a stack-splitting bugfix that sometimes affected franz-go.

  • PurgeTopicsFromConsuming now correctly persists deleted topics if
    you also had specific topics paused. Previously, when a topic had
    partition-level pauses, unpausing the topic itself (while keeping
    specific partitions paused) was bugged and the topic was stuck in
    an "all paused" state (thanks @​gorakdev!).

  • PollFetches no longer surfaces a spurious UNSTABLE_OFFSET_COMMIT
    fetch error when an OffsetFetch retry is canceled mid-wait by a
    rebalance or client close. Observed flaking TestTxnEtl/sticky/848
    on KIP-848 consumer groups under transactional load.

  • The MSK IAM SASL mechanism now honors AWS_REGION when the broker
    hostname does not match the standard MSK URL format, allowing
    connections through custom DNS names (e.g. private link endpoints)
    (thanks @​janmoritzmeyer0210!).

Relevant commits


Configuration

📅 Schedule: (in timezone Australia/Melbourne)

  • Branch creation
    • At any time (no schedule defined)
  • Automerge
    • At any time (no schedule defined)

🚦 Automerge: Disabled by config. Please merge this manually once you are satisfied.

♻ Rebasing: Whenever PR becomes conflicted, or you tick the rebase/retry checkbox.

🔕 Ignore: Close this PR and you won't be reminded about this update again.


  • If you want to rebase/retry this PR, check this box

This PR has been generated by Mend Renovate CLI.

This PR contains the following updates: | Package | Change | [Age](https://docs.renovatebot.com/merge-confidence/) | [Confidence](https://docs.renovatebot.com/merge-confidence/) | |---|---|---|---| | [github.com/twmb/franz-go](https://github.com/twmb/franz-go) | `v1.21.1` → `v1.22.1` | ![age](https://developer.mend.io/api/mc/badges/age/go/github.com%2ftwmb%2ffranz-go/v1.22.1?slim=true) | ![confidence](https://developer.mend.io/api/mc/badges/confidence/go/github.com%2ftwmb%2ffranz-go/v1.21.1/v1.22.1?slim=true) | --- ### Release Notes <details> <summary>twmb/franz-go (github.com/twmb/franz-go)</summary> ### [`v1.22.1`](https://github.com/twmb/franz-go/blob/HEAD/CHANGELOG.md#v1221) [Compare Source](https://github.com/twmb/franz-go/compare/v1.22.0...v1.22.1) \=== This patch has a few bug fixes in the share consumer (one via bug report, the others via a corresponding targeted audit), the 848 consumer (these found only via an audit), and some fixes in the `RecordFormatter` and `RecordReader`. This patch also has produce and consume performance improvements that come with two minor behavior changes: - `MaxBufferedRecords` now has a default of 50K, up from 10K: 10K was chosen when I initially wrote this library and is a very low default limit for average sized records. The librdkafka default is 100K; 50K increases producer throughput while still keeping producer memory low. - `PoolKRecords` is deprecated and unused. The client now decodes fetched records straight into `Record`s, so there is no need for pooling `kmsg.Record`s. Thanks to [@&#8203;kmrgirish](https://github.com/kmrgirish) for the share consumer bug report ([#&#8203;1474](https://github.com/twmb/franz-go/issues/1474)), and to [@&#8203;ajavanma](https://github.com/ajavanma) and [@&#8203;jakezwang](https://github.com/jakezwang) for `RecordReader` and `RecordFormatter` fixes. ### [`v1.22.0`](https://github.com/twmb/franz-go/blob/HEAD/CHANGELOG.md#v1220) [Compare Source](https://github.com/twmb/franz-go/compare/v1.21.7...v1.22.0) \=== This release supports Kafka 4.3 and 4.4, has a few new APIs, and has a few big internal improvements. In particular, I recommend checking out the new `StreamingCompression` option, as well as evaluating if you'd like to use `RackAwarePartitioning`. There are some behavior changes that you should read about below. The "next gen" rebalancer is now usable via the new `ServerSideBalancer` option. It's had a few releases to shake out bugs internally (via integration tests and LLM audits), but if you do experience a bug, please open an issue straightaway. Some minor bug fixes (that were never reported) were found during the implementation that are not worth mentioning. kfake has also been significantly extended and I recommend checking out the new APIs, in particular: - A new Fault type to make it easier to inject errors without Control functions - Group introspection cluster APIs - BlackholeProduce and SyntheticFetch APIs for benchmarking / play testing My `kcl` CLI has been *significantly* expanded as well and is worth checking out. It supports essentially everything you can do with a cluster, and now allows you to run a full broker locally via `kcl fake` (in memory or a dumb disk backed localhost broker) - as well as setup the fake broker with fault injection. I've been running LLM audits and extensions to `kcl` in particular to try to shape it up to a "finalized" CLI shape. If you use it and have ideas for improvements, please open an issue. #### Behavior changes - **Rack aware group partition assignment (KIP-881) now requires `BalanceRacks`.** v1.21.0 enabled group balancers to assign partitions based on the rack that members were in if you used the range or sticky/cooperative-sticky balancers. Well, `Rack` is also used to opt into preferred read replica assignment when fetching by the broker itself. These two decisions conflict with each other. Now, `BalanceRacks()` is required to opt into group balancers using the rack while balancing. The client warns when balancing if `BalanceRacks` is on and the brokers have preferred read replicas enabled. - **`ConsumeResetOffset` defaults to `RewindOffset(time.Minute)`** rather than `NewOffset().AtStart()`. **Setting only `ConsumeStartOffset` no longer sets `ConsumeResetOffset`**. I introduced `ConsumeStartOffset` a while back because it was really weird IMO to use a reset offset for both how a consumer starts *and* for how it recovers in the event of data loss or falling behind. They were bidirectional since introduction, but since start is newer and much less commonly used and you often don't want to recover from the start, I've removed the start -> reset mapping when you only set the start. I recommend reading the docs on both options for an updated understanding of when and how they apply. As well, I've introduced `RewindOffset(d)` which is *only* relevant to the reset offset (rewind by `d` duration from the last consumed offset on data loss we cannot exactly recover from) and `LookbackOffset(d)` which is relevant to both options but more useful for the start offset (start consuming `d` before the newest record; before Kafka 3.0 it is `d` before the current time). If a committed offset has fallen below the log start, the first fetch answers `OFFSET_OUT_OF_RANGE` and the reset offset decides where to resume. Before, a start offset of `AtEnd` was copied into the reset offset, so the consumer skipped to the end. Now, with the defaults, it resumes at the log start. - **Topic recreation is now a hard failure.** The client always produces to and consumes from the first instance of a topic. If you delete and recreate a topic, the client refuses the new version: buffered records fail with `UNKNOWN_TOPIC_ID`, fetches stop, offsets from the old topic cannot be committed to the new one, and transactions on the old topic fail. This needs a broker that reports topic IDs (Kafka 2.8+). Previously, some things in the client continued to accidentally work, and the behavior was unreliable and usually not good. If you want your application to stay alive across topic recreations, you can `PurgeTopicsFromClient` and, for consumers, `AddConsumeTopics`. More details about topic recreation are now in a new section in the README. - **`MaxDecompressBatchBytes` now blocks decompression if a batch would decompress too large (default 1GiB)**. Fetches when consuming can only specify to the broker "give me X bytes of batches", but they cannot control how large those batches decompress into. A hostile or buggy batch could OOM your program. Now, a batch over the limit causes the partition to enter a fatal state and return `ErrDecompressTooLarge` once from polling. The application can recover by manually skipping the batch with `SetOffsets` (with the fields in the error; see the docs), or by restarting the client with a higher limit. This option does not apply to custom decompressors, but, custom decompressors can still return `ErrMaxDecompress` to stop the partition. This option is also closely related to streaming compression, which is described below. #### Improvements - **gzip now uses klauspost/compress** (same format). Its default level is 1.7x faster than stdlib's with a slightly better ratio; klauspost's default maps to its level 5 where stdlib's mapped to 6, and level for level it is 1.1x to 1.2x faster. `WithLevel(n)` now selects klauspost's level `n`, so the bytes a given level produces differ from before. - **The sticky balancers are now exactly optimal** on balance, then rack placement (with `BalanceRacks`), then stickiness. Balancing was already load optimal but had some very niche edge cases where maximal stickiness was not preserved, especially if balancing used racks. Rack placement outranks stickiness: turning `BalanceRacks` on in a running group reassigns, at its next rebalance, every partition held by a member in a different zone from the partition's leader. - Sticky balancing is much faster, most of all on rejoins and on groups whose members subscribe to different topics. Against v1.21.7: a rejoin of 100 members over 1600 topics of 100 partitions goes from 351ms to 24ms; a regex shaped group of 500 members over 20,000 topics from 176ms and 810MB to 12ms and 9MB; 2001 members over 500 topics of 2000 partitions with one narrow subscriber from 3.4s to 0.3s. Fresh uniform balances are unchanged. #### Features ##### Streaming compression `StreamingCompression` is an opt-in producer option that compresses a partition's backlog of batches together, bounded by their *compressed* size. By default a batch is cut at `ProducerBatchMaxBytes` measured on uncompressed records. Streaming compression will help reduce traffic to the broker and increase how effective compression actually is (by pulling more data in at once). A custom compressor makes this option a no-op. The client is implemented such that each compression codec's worst case overhead is tracked internally, which should avoid a compressed batch ever exceeding `ProducerBatchMaxBytes`. If this ever does happen, the client discards the merge, logs a warning, disables streaming compression for the client going forward (records are still compressed batch by batch), and asks you to file an issue. The client has a new option `MaxDecompressBatchBytes` to bound both (a) how much the producer can stuff into a merged batch (i.e. how much it will decompress into), and (b) the maximum size a consumer will decompress a batch to; the consumer never decompresses past the bound (preventing a zip bomb). The default is 1GiB. ##### Rack aware producer partitioning (KIP-1123) `RackAwarePartitioning` sends unkeyed records to partitions whose leader is in the client's `Rack` (which must also be set), falling back to all partitions when no leader is. Keyed records are never affected. Unlike the Java client, this works with any partitioner, since the eligible-broker filtering happens before your partitioner is consulted. Note that this option skews which partitions receive records if your producers are not spread across racks in proportion to partition leaders. ##### ServerSideBalancer (KIP-848) `ServerSideBalancer` opts into KIP-848 "next-gen" consumer groups, where the broker's group coordinator assigns partitions rather than the client. This requires Kafka 4.0+ and either a range or sticky / cooperative-sticky balancer. This replaces the hidden `opt_in_kafka_next_gen_balancer_beta` context key from v1.19.0; the key still works in this release but will be removed in the next. The default remains the classic protocol, matching the Java client. I still think the classic client side balancers are better (and this client's implementation is way faster than the Java client), but if you want to use server side balancing, it is strongly recommended to only use it if your cluster is Kafka 4.3+. Before 4.3 (before KIP-1251), an offset commit that races with a heartbeat epoch bump can fail with `STALE_MEMBER_EPOCH`, which the client cannot detect nor handle. ##### BalanceInfo for custom balancers A balancer that implements `GroupMemberBalancerInfo` receives a `BalanceInfo` before balancing: the group, generation, leader member ID, and lazily built topic and broker metadata. `ConsumerBalancer` implements it, so balancers built on `NewConsumerBalancer` can call `Info()`. This allows, for example, a balancer that assigns every partition to the leader with the other members as hot standbys. Thanks [@&#8203;michaelwilner](https://github.com/michaelwilner)! #### API additions ```go // Producing func StreamingCompression() ProducerOpt func RackAwarePartitioning() ProducerOpt // Consuming func BalanceRacks() ConsumerOpt func ServerSideBalancer() GroupOpt func RewindOffset(d time.Duration) Offset func LookbackOffset(d time.Duration) Offset // Decompression bound func MaxDecompressBatchBytes(n int) Opt var ErrMaxDecompress error type ErrDecompressTooLarge struct { Topic string Partition int32 Offset int64 Epoch int32 NextOffset int64 } // Custom balancers type BalanceInfo struct { Group string Generation int32 LeaderID string Topics func() map[string]TopicMetadata Brokers func() map[int32]BrokerMetadata } type GroupMemberBalancerInfo interface { GroupMemberBalancer SetBalanceInfo(BalanceInfo) } func (*ConsumerBalancer) Info() BalanceInfo type TopicMetadata struct { ... } type PartitionMetadata struct { ... } // Records type RecordAttrsOpts struct { Codec CompressionCodecType TimestampType int8 Transactional bool Control bool } func NewRecordAttrs(RecordAttrsOpts) RecordAttrs // kversion func (*Versions) EachSupportedFeature(fn func(name string, min, max int16)) func (*Versions) EachFinalizedFeature(fn func(name string, level int16)) func FeatureLevelDescription(name string, level int16) string ``` #### Relevant commits There are many commits, but some of the more notable ones: - [`27d11286`](https://github.com/twmb/franz-go/commit/27d11286) **feature** kversion: FeatureLevelDescription - [`73358f62`](https://github.com/twmb/franz-go/commit/73358f62) **feature** kversion: supported and finalized feature levels per release - [`7be0be16`](https://github.com/twmb/franz-go/commit/7be0be16) **behavior change** kgo: add MaxDecompressedBatchBytes - [`b37f1041`](https://github.com/twmb/franz-go/commit/b37f1041) **feature** kgo: add ServerSideBalancer to opt into KIP-848 - [`033a46c7`](https://github.com/twmb/franz-go/commit/033a46c7) **improvement** kgo: begin ApiVersions at the max a broker told us, for an hour - [`46a9b2ad`](https://github.com/twmb/franz-go/commit/46a9b2ad) **behavior change** kgo: use ConsumeResetOffset when the broker loses data we cannot locate - [`8b33e43d`](https://github.com/twmb/franz-go/commit/8b33e43d) **improvement** kgo: speed up compression on both the legacy and the merge path - [`9de0fa36`](https://github.com/twmb/franz-go/commit/9de0fa36) **feature** kgo: add StreamingCompression, compressed-size-bound batch merging - [`de7327e6`](https://github.com/twmb/franz-go/commit/de7327e6) **feature** kgo: detect misrouted connections (KIP-1242) - [`8ad36ec7`](https://github.com/twmb/franz-go/commit/8ad36ec7) **feature** kgo: support TxnOffsetCommit v6 - [`e4f7bc43`](https://github.com/twmb/franz-go/commit/e4f7bc43) **feature** kgo: add rack-aware producer partitioning (KIP-1123) - [`123f2ffa`](https://github.com/twmb/franz-go/commit/123f2ffa) **improvement** kgo: repair the sticky plan to the best balance, rack, and stickiness - [`23ab9a0e`](https://github.com/twmb/franz-go/commit/23ab9a0e) **behavior change** kgo: add BalanceRacks, gate rack aware balancing behind it - [`d4f6db2f`](https://github.com/twmb/franz-go/commit/d4f6db2f) **improvement** kgo: drop reassigned partitions in one pass in AdjustCooperative - [`35efafc8`](https://github.com/twmb/franz-go/commit/35efafc8) **behavior change** kgo: fail records for a recreated topic instead of producing by name - [`4f10346a`](https://github.com/twmb/franz-go/commit/4f10346a) **feature** kgo: expose BalanceInfo for custom balancer implementations (thanks [@&#8203;michaelwilner](https://github.com/michaelwilner)!) - [`cd7f9b4e`](https://github.com/twmb/franz-go/commit/cd7f9b4e) **feature** kgo: add NewRecordAttrs constructor (thanks [@&#8203;pracucci](https://github.com/pracucci)!) ### [`v1.21.7`](https://github.com/twmb/franz-go/blob/HEAD/CHANGELOG.md#v1217) [Compare Source](https://github.com/twmb/franz-go/compare/v1.21.6...v1.21.7) \=== A handful of bug fixes and improvements found by users and while working on v1.22. Rather than enumerating the relevant commits, you can check the git log between v1.21.6 and this release - there are many minor commits. As well, kfake has been improved significantly and has more API surface to aid in writing tests. - A rare, very niche panic while producing has been fixed. Thanks [@&#8203;PumpkinDemo](https://github.com/PumpkinDemo) for the report, see [#&#8203;1385](https://github.com/twmb/franz-go/issues/1385) for more details. - If retention deleted the segment a consumer was reading, the OffsetOutOfRange reset listed by the last consumed timestamp and could skip surviving records, or jump to the log end and skip everything. The reset now resumes at the log start when below it. - Improved KIP-951 handling (the broker returning where a partition should move with the produce response if the partition changed leadership). Previously, a broker could return NotLeaderForPartition *and* hint the leader the client was already using, at the same or an older epoch. These hints are now ignored and the client backs off, rather than spinning. Thanks [@&#8203;3AceShowHand](https://github.com/3AceShowHand) for the report and [@&#8203;jjj-n](https://github.com/jjj-n) for a fix, see [#&#8203;1412](https://github.com/twmb/franz-go/issues/1412). - Rack aware balancers ignored rack for any topic the group leader did not itself consume. A client now loads the rack for *all* partitions in the group, even if the leader does not consume some of the topics. - The metadata cache has been improved (there were a few cases where it was emptied erroneously). - Regex consuming now consistently never matches internal topics such as `__consumer_offsets`. As well, the regex log no longer reports an excluded topic as both added and skipped (thanks [@&#8203;lahsivjar](https://github.com/lahsivjar)). - Decompression allocates less (thanks [@&#8203;scunningham](https://github.com/scunningham)). If you use pools, slices are now reliably returned if decompression errors. - The client now starts at a random seed broker rather than always the first, so many clients starting at once no longer all hit the same seed. Thanks [@&#8203;chailuecha](https://github.com/chailuecha)! - A producer that receives RequestTimedOut or NotEnoughReplicasAfterAppend now retries after the produce backoff rather than waiting for a metadata refresh. - A few other minor improvements and bug fixes. ### [`v1.21.6`](https://github.com/twmb/franz-go/blob/HEAD/CHANGELOG.md#v1216) [Compare Source](https://github.com/twmb/franz-go/compare/v1.21.5...v1.21.6) \=== Some bug fixes (mostly minor - hence the delay for the release) found by users and further Claude audits. I am gearing up for a 1.22 release but some of the features I am planning for are more complicated to review, so it may take a bit of time. Anyway: - Previously, rollback from a cooperative group to an eager group was deliberately not supported and there was a data race condition if this happened. It is now *technically* supported, although you will experience duplicate data. If you want a safe non-duplicate-causing rollback, you need to turn off the entire group, remove the cooperative consumer, and swap the whole group to eager rebalancing. - Fixed a `panic: close of closed channel` on an `acks=0` produce connection in a specific edge case (a broker connection dying before the connection was fully established caused the panic). - If `EndTransaction` failed with an unconfirmed outcome (a transport error, exhausted retries, or `UNKNOWN_SERVER_ERROR`), the documented abort retry was a wire no-op and the next transaction could silently commit the prior "failed" transaction's records under KIP-890 part 2. The producer ID is now flagged for reload, which fence-aborts anything still ongoing broker-side. - `GroupTransactSession.End` could hang forever, ignoring its context, if the group had never joined (e.g. the consumed topic did not exist yet) and the transaction committed no offsets. - Previously, if a broker replied to ApiVersions with an error, we ignored it and you would eventually see an unclear error (usually a bare io.EOF, since anything that rejects ApiVersions hangs up right after replying). These errors are now handled correctly. - Some niche edge case bugs that are only worth reading about if you're super interested were found in repeated Claude audits and were fixed (check the PR / git history). This includes further KIP-848 "next gen consumer group" fixes. #### Relevant commits - [`582e0f21`](https://github.com/twmb/franz-go/commit/582e0f21) **bugfix** kgo: surface error codes in ApiVersions responses - [`67ef4c61`](https://github.com/twmb/franz-go/commit/67ef4c61) **bugfix** kgo: fix double close of a connection's deadCh on acks=0 produce - [`3ac2fff1`](https://github.com/twmb/franz-go/commit/3ac2fff1) **bugfix** kgo: revoke everything when the group protocol downgrades from cooperative to eager - [`795d5b61`](https://github.com/twmb/franz-go/commit/795d5b61) **improvement** kgo: flatten topic/partition maps in group rebalance logs (thanks [@&#8203;constanca-m](https://github.com/constanca-m)!) - [`70addc1e`](https://github.com/twmb/franz-go/commit/70addc1e) **improvement** kgo: classify retired broker reads as broker dead (thanks [@&#8203;tomplarge](https://github.com/tomplarge)!) - [`18f9a10f`](https://github.com/twmb/franz-go/commit/18f9a10f) **improvement** deps: replace golang.org/x/crypto/pbkdf2 with stdlib crypto/pbkdf2 (thanks [@&#8203;macdewee](https://github.com/macdewee)!) - [`6ecd2f9f`](https://github.com/twmb/franz-go/commit/6ecd2f9f) **bugfix** kgo: recover when an attempted EndTxn outcome is unconfirmed - [`821f879e`](https://github.com/twmb/franz-go/commit/821f879e) **bugfix** kgo: fix GroupTransactSession.End hanging when the group never joined ### [`v1.21.5`](https://github.com/twmb/franz-go/blob/HEAD/CHANGELOG.md#v1215) [Compare Source](https://github.com/twmb/franz-go/compare/v1.21.4...v1.21.5) \=== Three bug fixes: - Fixed a nil-pointer panic when building a group `OffsetFetch`: if the group was assigned a topic that was no longer in the client's tracked set -- reachable when a topic is purged from consuming while still assigned, for example `PurgeTopicsFromConsuming` overlapping `AddConsumeTopics`, or the automatic regex missing-topic purge -- `loadTopic` returned nil and dereferencing it for the topic ID crashed the client. The topic ID is now only set when the topic is known. Thanks [@&#8203;iwittkau](https://github.com/iwittkau)! - Fixed a data race on a coordinator's cached node ID. When a broker disconnected while a `FindCoordinator` load for that broker was still in flight, `deleteStaleCoordinatorsByNode` could read the in-flight load's `node` field before the loading goroutine published it (via closing the load's wait channel), which `go test -race` flagged. The node read now happens only after the load has been observed as complete. Thanks [@&#8203;nikolauspschuetz](https://github.com/nikolauspschuetz)! - A share partition that was listed in a ShareFetch only to carry a piggybacked acknowledgement -- for a cursor that was revoked, paused, or migrated to a new leader after its records were drained -- was added to the broker's share session but never tracked client-side, so it could never be forgotten. The broker would re-acquire and redeliver that partition's records indefinitely while the client discarded them ("broker returned partition ... we did not ask for"), spinning the share fetch loop. The client now tracks every partition it sends, matching the broker's session bookkeeping. #### Relevant commits - [`ab185e42`](https://github.com/twmb/franz-go/commit/ab185e42) **bugfix** kgo: add nil check when loading topics in groupConsumer (thanks [@&#8203;iwittkau](https://github.com/iwittkau)!) - [`aca084ed`](https://github.com/twmb/franz-go/commit/aca084ed) **bugfix** kgo: fix data race on coordinatorLoad.node (thanks [@&#8203;nikolauspschuetz](https://github.com/nikolauspschuetz)!) - [`754bc349`](https://github.com/twmb/franz-go/commit/754bc349) **bugfix** kgo: forget piggyback-only partitions from the share session ### [`v1.21.4`](https://github.com/twmb/franz-go/blob/HEAD/CHANGELOG.md#v1214) [Compare Source](https://github.com/twmb/franz-go/compare/v1.21.3...v1.21.4) \=== This release is a "large" (many commits) release that has many small or hard to encounter bugs fixed. I pointed Claude's Fable at this repo and ran some audit rounds while available and thankfully got through the highest value audit rounds before Fable was removed. For once, I will not be describing every bug fixed nor calling out every relevant commit. Instead, if you are curious, look at [#&#8203;1348](https://github.com/twmb/franz-go/pull/1348). Some worthwhile description is below. Three important bug fixes to call out: - In transactional exactly-once consuming, a `SetOffsets` seek (which happens during `GroupTransactSession.End` after an aborted transaction) could be undone by a concurrent offset load (via a background list or epoch load) that completed slightly later. This could happen when the client discovers a partition leader moved while you are aborting, which could result in missed records. - Consuming with `read_committed` against a broker that returns a partition's aborted-transaction list out of offset order could surface aborted, rolled-back records as if they were committed. Apache Kafka always returns them in order so this was never observed there, but Redpanda does not (when an aborted transaction is still in memory and an earlier one is already on disk). The list is now sorted client-side, matching the Java client, librdkafka, and Sarama. - `GroupTransactSession.End` no longer reports a successful commit when the broker answers `EndTxn` with `UNKNOWN_SERVER_ERROR` (seen from Redpanda in some older versions). Previously the consumer's offsets were advanced past a transaction that may have aborted; now the commit is reported as failing and the session rewinds for reprocessing. Beyond those, by area: - Many transaction-path fixes for coordinator churn and KIP-890 part 2 that would have resulted in not-working (hard client fail) or hung transactions: `InitProducerID` retries `CONCURRENT_TRANSACTIONS` when taking over a crashed producer's transaction, retriable producer-id load failures are no longer treated as fatal, KIP-890p2 is opted into only when the negotiated versions actually support it (fixing spurious `INVALID_TXN_STATE` on 4.0+ clusters running older semantics), a transaction whose every produce failed now aborts instead of hanging until the transaction timeout, and a failed `AddPartitionsToTxn` no longer drops partitions added by an earlier request. - `GzipCompression().WithLevel(...)` was completely broken and would panic. - More KIP-848 (next-gen consumer group) robustness fixes under coordinator and leader churn. - Stale consumer-group member rejoining fixes: a member that rejoins claiming a partition at an old generation no longer panics the group leader or causes two members to consume the same partition. Malformed member metadata, duplicate member ids, and negative claimed partitions in a join are now rejected or sanitized rather than mis-balancing or panicking the leader. - Metadata and topic recreation: a stale per-broker metadata view that momentarily omits a just-added partition (the window right after `CreatePartitions`) no longer fails buffered producer records or leaves a newly assigned consumer / share partition silently unconsumed; both heal once metadata catches up. - Share consumer: fetch errors are now classified like the classic consumer (retriable errors stripped, a metadata refresh triggered to heal a leader move, top-level errors backed off) instead of stalling for up to `MetadataMaxAge` or hot-looping, and leader-move migrations are tracked so that leaving or closing cannot strand un-acked records. - SASL: KIP-368 re-authentication no longer races the connection's other reader, which could corrupt pipelined traffic on brokers that set a session lifetime (e.g. AWS MSK IAM); requests now park and replay across a re-auth. The Azure Event Hubs ApiVersions reset retry no longer leaks the abandoned connection or silently downgrades it to v0. - KIP-714 client telemetry: the terminating push is now actually delivered on `Close`, the `.rate` and `.avg` rollups are computed correctly (they were constant / wrong before), and an unsupported user-metric attribute no longer corrupts the OTLP payload (which had disabled metrics for the rest of the client's life). - Smaller consumer fixes: overlapping manual `CommitOffsets` no longer reopen autocommit early (which could rewind the committed offset), a conformant `UNDEFINED_EPOCH_OFFSET` epoch response no longer raises a false `ErrDataLoss`, and `Fetches.EachTopic` now preserves `TopicID` across multi-broker responses (it was zero whenever more than one broker replied, i.e. normally). - `RecordReader` / `RecordFormatter` no longer panic on truncated or malformed layouts, accept `\xNN` escapes for bytes above 0x7f, and reject layouts that would read nothing and loop forever. - `WithPools`: decompression no longer produces garbage when a pool hands back a non-zero-length sized slice, and pooled slices are no longer leaked for batches that keep no records (e.g. aborted-transaction data under `read_committed`). - Other producer fixes: `EnsureProduceConnectionIsOpen` dials the right broker for filtered ids and no longer breaks an `acks=0` connection, the adaptive `LeastBackupPartitioner` now actually picks the least-backed-up partition, and producing during or after `Close` fails cleanly instead of hanging a later `Flush`. - A broad set of guards against malformed or hostile broker responses that could previously panic the fetcher, hot-loop, or mis-consume: negative or oversized record counts and batch lengths, decompression bombs, duplicate or omitted partitions, negative offsets, and unexpected top-level fetch errors. The client is also more resilient when its own API contracts are violated (e.g. `AllowRebalance` called while a poll is in flight). ### [`v1.21.3`](https://github.com/twmb/franz-go/blob/HEAD/CHANGELOG.md#v1213) [Compare Source](https://github.com/twmb/franz-go/compare/v1.21.2...v1.21.3) \=== This patch release contains a few bug fixes and a few internal improvements. - `PollRecords` / `PollFetches` could permanently hang since v1.21.0 if a consumer session stopped (usually via metadata updates) while fetches to more than four brokers were pending and no poll was in flight. This could only affect users that deliberately set `MaxConcurrentFetches(0)`, or that were using `ShareMaxRecordsStrict`. - Producing to a topic whose partitions ALL have a retriable load error (e.g. a rolling restart of an RF=1 broker briefly leaving every partition leaderless) no longer fails records up front with "unable to partition record due to no usable partitions". Instead, the records remain buffered and retried as metadata reloads. - Classic consumer groups now rejoin immediately when an offset commit returns `UNKNOWN_MEMBER_ID` or `ILLEGAL_GENERATION` (the broker lost the member, e.g. a session expired during a network blip), rather than consuming as a zombie until the heartbeat loop notices the dead session. - `DescribeShareGroupOffsets`, `AlterShareGroupOffsets`, and `DeleteShareGroupOffsets` are now routed to the group coordinator rather than the share coordinator (which would reject the requests for being misrouted). - The client-internal metadata cache now deeply clones the cached response before putting it into the cache and before returning it via `RequestCachedMetadata` (which is now used by default in kadm), eliminating data race possibilities. - Various next-gen rebalancer session improvements. #### Relevant commits - [`f8842170`](https://github.com/twmb/franz-go/commit/f8842170) **improvement** kgo: fall back to all partitions when no partition is writable (thanks [@&#8203;ericsg666](https://github.com/ericsg666)!) - [`824e34d2`](https://github.com/twmb/franz-go/commit/824e34d2) **improvement** kgo: rejoin a classic group when a commit returns a fatal member error (thanks [@&#8203;v14dis14v](https://github.com/v14dis14v)!) - [`f520e820`](https://github.com/twmb/franz-go/commit/f520e820) **bugfix** kgo: do not exit manageFetchConcurrency while sources are pending in wantFetch (thanks [@&#8203;SLoeuillet](https://github.com/SLoeuillet)!) - [`8d9c836b`](https://github.com/twmb/franz-go/commit/8d9c836b) **bugfix** kgo: isolate metadata cache from broker response - [`19f7dbb2`](https://github.com/twmb/franz-go/commit/19f7dbb2) **bugfix** kgo,kfake: route share group offset RPCs to the group coordinator ### [`v1.21.2`](https://github.com/twmb/franz-go/blob/HEAD/CHANGELOG.md#v1212) [Compare Source](https://github.com/twmb/franz-go/compare/v1.21.1...v1.21.2) \=== This patch release contains two narrow bug fixes and one small feature. Deps are also bumped so that you are force-pinned to a klauspost/compress version that has a stack-splitting bugfix that sometimes affected franz-go. - `PurgeTopicsFromConsuming` now correctly persists deleted topics if you also had specific topics paused. Previously, when a topic had partition-level pauses, unpausing the topic itself (while keeping specific partitions paused) was bugged and the topic was stuck in an "all paused" state (thanks [@&#8203;gorakdev](https://github.com/gorakdev)!). - `PollFetches` no longer surfaces a spurious `UNSTABLE_OFFSET_COMMIT` fetch error when an OffsetFetch retry is canceled mid-wait by a rebalance or client close. Observed flaking `TestTxnEtl/sticky/848` on KIP-848 consumer groups under transactional load. - The MSK IAM SASL mechanism now honors `AWS_REGION` when the broker hostname does not match the standard MSK URL format, allowing connections through custom DNS names (e.g. private link endpoints) (thanks [@&#8203;janmoritzmeyer0210](https://github.com/janmoritzmeyer0210)!). #### Relevant commits - [`22a17320`](https://github.com/twmb/franz-go/commit/22a17320) **bugfix** kgo: do not inject fake fetch error when OffsetFetch retry is canceled - [`95d74ab3`](https://github.com/twmb/franz-go/commit/95d74ab3) **bugfix** kgo: fix delTopics not writing back modified pausedPartitions struct (thanks [@&#8203;gorakdev](https://github.com/gorakdev)!) - [`40e5a0e5`](https://github.com/twmb/franz-go/commit/40e5a0e5) **feature** sasl/aws: support custom AWS MSK DNS names (thanks [@&#8203;janmoritzmeyer0210](https://github.com/janmoritzmeyer0210)!) </details> --- ### Configuration 📅 **Schedule**: (in timezone Australia/Melbourne) - Branch creation - At any time (no schedule defined) - Automerge - At any time (no schedule defined) 🚦 **Automerge**: Disabled by config. Please merge this manually once you are satisfied. ♻ **Rebasing**: Whenever PR becomes conflicted, or you tick the rebase/retry checkbox. 🔕 **Ignore**: Close this PR and you won't be reminded about this update again. --- - [ ] <!-- rebase-check -->If you want to rebase/retry this PR, check this box --- This PR has been generated by [Mend Renovate CLI](https://github.com/renovatebot/renovate). <!--renovate-debug:eyJjcmVhdGVkSW5WZXIiOiI0NC40OC4wIiwidXBkYXRlZEluVmVyIjoiNDQuNDguMCIsInRhcmdldEJyYW5jaCI6Im1haW4iLCJsYWJlbHMiOlsidHlwZS9taW5vciJdfQ==-->
fix(deps): update module github.com/twmb/franz-go ( v1.21.1 → v1.21.6 )
All checks were successful
loader-spf-kafka CI / build-and-push (pull_request) Successful in 2m4s
Go lint / lint (pull_request) Successful in 4m4s
5d8942ea19
Author
Member

ℹ️ Artifact update notice

File name: go.mod

In order to perform the update(s) described in the table above, Renovate ran the go get command, which resulted in the following additional change(s):

  • 3 additional dependencies were updated
  • The go directive was updated for compatibility reasons

Details:

Package Change
go 1.25.0 -> 1.26.0
github.com/klauspost/compress v1.18.5 -> v1.20.0
github.com/pierrec/lz4/v4 v4.1.26 -> v4.1.30
github.com/twmb/franz-go/pkg/kmsg v1.13.1 -> v1.14.0
### ℹ️ Artifact update notice ##### File name: go.mod In order to perform the update(s) described in the table above, Renovate ran the `go get` command, which resulted in the following additional change(s): - 3 additional dependencies were updated - The `go` directive was updated for compatibility reasons Details: | **Package** | **Change** | | :---------------------------------- | :--------------------- | | `go` | `1.25.0` -> `1.26.0` | | `github.com/klauspost/compress` | `v1.18.5` -> `v1.20.0` | | `github.com/pierrec/lz4/v4` | `v4.1.26` -> `v4.1.30` | | `github.com/twmb/franz-go/pkg/kmsg` | `v1.13.1` -> `v1.14.0` |
renovatebot force-pushed renovate/github.com-twmb-franz-go-1.x from 5d8942ea19
All checks were successful
loader-spf-kafka CI / build-and-push (pull_request) Successful in 2m4s
Go lint / lint (pull_request) Successful in 4m4s
to 8eff556be1
All checks were successful
loader-spf-kafka CI / build-and-push (pull_request) Successful in 2m48s
Go lint / lint (pull_request) Successful in 3m30s
2026-09-15 04:22:36 +00:00
Compare
renovatebot changed title from fix(deps): update module github.com/twmb/franz-go ( v1.21.1 → v1.21.6 ) to fix(deps): update module github.com/twmb/franz-go ( v1.21.1 → v1.21.7 ) 2026-09-15 04:22:37 +00:00
renovatebot force-pushed renovate/github.com-twmb-franz-go-1.x from 8eff556be1
All checks were successful
loader-spf-kafka CI / build-and-push (pull_request) Successful in 2m48s
Go lint / lint (pull_request) Successful in 3m30s
to 6fed0272cb
Some checks failed
loader-spf-kafka CI / build-and-push (pull_request) Failing after 45s
Go lint / lint (pull_request) Successful in 2m31s
2026-09-18 08:22:18 +00:00
Compare
renovatebot changed title from fix(deps): update module github.com/twmb/franz-go ( v1.21.1 → v1.21.7 ) to feat(deps): update module github.com/twmb/franz-go ( v1.21.1 → v1.22.0 ) 2026-09-18 08:22:19 +00:00
renovatebot force-pushed renovate/github.com-twmb-franz-go-1.x from 6fed0272cb
Some checks failed
loader-spf-kafka CI / build-and-push (pull_request) Failing after 45s
Go lint / lint (pull_request) Successful in 2m31s
to 085bbf9255
Some checks failed
loader-spf-kafka CI / build-and-push (pull_request) Failing after 47s
Go lint / lint (pull_request) Successful in 2m38s
2026-09-28 00:21:53 +00:00
Compare
renovatebot changed title from feat(deps): update module github.com/twmb/franz-go ( v1.21.1 → v1.22.0 ) to feat(deps): update module github.com/twmb/franz-go ( v1.21.1 → v1.22.1 ) 2026-09-28 00:21:54 +00:00
Some checks failed
loader-spf-kafka CI / build-and-push (pull_request) Failing after 47s
Go lint / lint (pull_request) Successful in 2m38s
This pull request can be merged automatically.
You are not authorized to merge this pull request.
View command line instructions

Checkout

From your project repository, check out a new branch and test the changes.
git fetch -u origin renovate/github.com-twmb-franz-go-1.x:renovate/github.com-twmb-franz-go-1.x
git switch renovate/github.com-twmb-franz-go-1.x

Merge

Merge the changes and update on Forgejo.

Warning: The "Autodetect manual merge" setting is not enabled for this repository, you will have to mark this pull request as manually merged afterwards.

git switch main
git merge --no-ff renovate/github.com-twmb-franz-go-1.x
git switch renovate/github.com-twmb-franz-go-1.x
git rebase main
git switch main
git merge --ff-only renovate/github.com-twmb-franz-go-1.x
git switch renovate/github.com-twmb-franz-go-1.x
git rebase main
git switch main
git merge --no-ff renovate/github.com-twmb-franz-go-1.x
git switch main
git merge --squash renovate/github.com-twmb-franz-go-1.x
git switch main
git merge --ff-only renovate/github.com-twmb-franz-go-1.x
git switch main
git merge renovate/github.com-twmb-franz-go-1.x
git push origin main
Sign in to join this conversation.
No reviewers
No labels
No milestone
No project
No assignees
1 participant
Notifications
Due date
The due date is invalid or out of range. Please use the format "yyyy-mm-dd".

No due date set.

Dependencies

No dependencies set

Reference
dmarc-ing/loader-spf-kafka!3
No description provided.