From 604b27940f918caf90224e6970ff5a1b9fa3bf8f Mon Sep 17 00:00:00 2001 From: Thomas Flament Date: Wed, 2 Sep 2026 11:26:44 +0200 Subject: [PATCH 1/2] Guard the deferred un-assign against a superseding rebalance On ERR__REVOKE_PARTITIONS the un-assign is deferred until in-flight work drains. librdkafka however answers a subscribed-topic metadata change with an immediate rejoin even while that revoke is still unanswered, so a fresh assignment can be granted and applied before the drain completes -- and the deferred un-assign then discarded it, leaving a live group member owning partitions at the broker with no local assignment, consuming nothing, with nothing left to trigger a rebalance. Disarm the pending drain triggers and watchdog on every rebalance event, so a superseded un-assign can no longer fire; the watchdog disarm also stops a second revoke from leaving an orphaned timer that would later disconnect a healthy consumer. The async offset publish can still be in flight when a newer rebalance lands, so the un-assign continuation also checks a rebalance counter before releasing the assignment. Issue: BB-835 --- lib/BackbeatConsumer.js | 23 ++++++ tests/unit/backbeatConsumer.js | 145 +++++++++++++++++++++++++++++++++ 2 files changed, 168 insertions(+) diff --git a/lib/BackbeatConsumer.js b/lib/BackbeatConsumer.js index 8df60c538..0ae1f67a9 100644 --- a/lib/BackbeatConsumer.js +++ b/lib/BackbeatConsumer.js @@ -170,6 +170,9 @@ class BackbeatConsumer extends EventEmitter { // installed below delegates through this flag, so we don't // need to add/remove ledger listeners on every rebalance. this._drainCallback = null; + // Bumped on every rebalance event; a deferred un-assign + // captured under an older value must not run. + this._rebalanceId = 0; this._offsetLedger.on('empty', topic => { if (topic === this._topic && this._drainCallback) { this._drainCallback(); @@ -758,6 +761,15 @@ class BackbeatConsumer extends EventEmitter { * @returns {void} */ _onRebalance(err, assignment) { + // A new rebalance supersedes whatever a previous revoke left + // pending: disarm its drain triggers and watchdog, and stamp + // this event so a continuation already in flight (the async + // offset publish below) can tell it has been superseded. + const rebalanceId = ++this._rebalanceId; + this._setDrain(null); + clearTimeout(this._drainProcessQueueTimeout); + this._drainProcessQueueTimeout = null; + if (err.code === kafka.CODES.ERRORS.ERR__ASSIGN_PARTITIONS) { this._log.info('rdkafka.assign', { assignment }); @@ -804,6 +816,17 @@ class BackbeatConsumer extends EventEmitter { } const doUnassign = () => { + if (rebalanceId !== this._rebalanceId) { + // a newer rebalance landed while offsets were + // being published: the assignment is no longer + // ours to release + this._log.info('skipping superseded un-assign', { + topic: this._topic, + rebalanceId, + currentRebalanceId: this._rebalanceId, + }); + return; + } this._resumePausedPartitions(); try { diff --git a/tests/unit/backbeatConsumer.js b/tests/unit/backbeatConsumer.js index 951fad035..884ceefdd 100644 --- a/tests/unit/backbeatConsumer.js +++ b/tests/unit/backbeatConsumer.js @@ -2,9 +2,11 @@ const assert = require('assert'); const sinon = require('sinon'); const BackbeatConsumer = require('../../lib/BackbeatConsumer'); +const KafkaBacklogMetrics = require('../../lib/KafkaBacklogMetrics'); const { CODES } = require('node-rdkafka'); const { kafka } = require('../config.json'); +const { unassignStatus } = require('../../lib/constants'); const { BreakerState } = require('breakbeat').CircuitBreaker; class BackbeatConsumerMock extends BackbeatConsumer { @@ -397,4 +399,147 @@ describe('backbeatConsumer', () => { }); }); }); + + describe('_onRebalance deferred un-assign', () => { + const REVOKE = { code: CODES.ERRORS.ERR__REVOKE_PARTITIONS }; + const ASSIGN = { code: CODES.ERRORS.ERR__ASSIGN_PARTITIONS }; + const partitions = [ + { topic: 'my-test-topic', partition: 0 }, + { topic: 'my-test-topic', partition: 1 }, + ]; + + let consumer; + let queueIdle; + let ledgerCount; + let unassigns; + + beforeEach(() => { + consumer = new BackbeatConsumerMock({ + kafka, + groupId: 'unittest-group', + topic: 'my-test-topic', + }); + + consumer._consumer = { + assign: sinon.stub(), + unassign: sinon.stub(), + disconnect: sinon.stub(), + commit: sinon.stub(), + pause: sinon.stub(), + resume: sinon.stub(), + isConnected: () => true, + assignments: () => [], + subscription: () => ['my-test-topic'], + }; + + queueIdle = false; + ledgerCount = 1; + consumer._processingQueue = { + length: () => 0, + running: () => (queueIdle ? 0 : 1), + idle: () => queueIdle, + setDrain: () => {}, + }; + consumer._offsetLedger.getProcessingCount = () => ledgerCount; + + unassigns = []; + consumer.on('unassign', status => unassigns.push(status)); + + sinon.stub(KafkaBacklogMetrics, 'onRebalance'); + }); + + afterEach(() => { + clearTimeout(consumer._drainProcessQueueTimeout); + sinon.restore(); + }); + + // fire the drain through the live callback, which a superseding + // rebalance nulls out + const completeDrain = () => { + queueIdle = true; + ledgerCount = 0; + if (consumer._drainCallback) { + consumer._drainCallback(); + } + }; + + it('should defer the un-assign until the drain completes', () => { + consumer._onRebalance(REVOKE, partitions); + assert(consumer._consumer.unassign.notCalled); + + completeDrain(); + + assert(consumer._consumer.unassign.calledOnce); + assert.deepStrictEqual(unassigns, [unassignStatus.DRAINED]); + }); + + it('should un-assign immediately when nothing is in flight', () => { + queueIdle = true; + ledgerCount = 0; + + consumer._onRebalance(REVOKE, partitions); + + assert(consumer._consumer.unassign.calledOnce); + assert.deepStrictEqual(unassigns, [unassignStatus.IDLE]); + }); + + it('should not un-assign when a new assignment arrived while ' + + 'draining', () => { + consumer._onRebalance(REVOKE, partitions); + + // the next generation grants the partitions back mid-drain + consumer._onRebalance(ASSIGN, partitions); + assert(consumer._consumer.assign.calledOnce); + assert.strictEqual(consumer._drainCallback, null); + assert.strictEqual(consumer._drainProcessQueueTimeout, null); + + completeDrain(); + // the ledger going empty must not fire anything either + consumer._offsetLedger.emit('empty', 'my-test-topic'); + + assert(consumer._consumer.unassign.notCalled); + assert.deepStrictEqual(unassigns, []); + }); + + it('should not un-assign when the partitions were granted back ' + + 'while offsets were being published', () => { + let publishDone; + consumer._kafkaBacklogMetricsConfig = + { zkPath: '/test', intervalS: 5 }; + consumer._publishOffsetsCron = cb => { + publishDone = cb; + }; + + consumer._onRebalance(REVOKE, partitions); + completeDrain(); + assert.strictEqual(typeof publishDone, 'function'); + assert(consumer._consumer.unassign.notCalled); + + consumer._onRebalance(ASSIGN, partitions); + publishDone(); + + assert(consumer._consumer.unassign.notCalled); + assert.deepStrictEqual(unassigns, []); + }); + + it('should not leave a superseded revoke watchdog armed', () => { + const clock = sinon.useFakeTimers(); + try { + consumer._onRebalance(REVOKE, partitions); + consumer._onRebalance(REVOKE, partitions); + + completeDrain(); + assert(consumer._consumer.unassign.calledOnce); + assert.deepStrictEqual(unassigns, [unassignStatus.DRAINED]); + + // the first revoke's watchdog must not fire late and + // disconnect a consumer that already answered + clock.tick(consumer._maxPollIntervalMs + 1000); + assert(consumer._consumer.disconnect.notCalled); + assert(consumer._consumer.unassign.calledOnce); + } finally { + clock.restore(); + } + }); + }); }); From 02a7eb76f37a32ecbbe805115185ee7b53483027 Mon Sep 17 00:00:00 2001 From: Thomas Flament Date: Wed, 2 Sep 2026 11:26:44 +0200 Subject: [PATCH 2/2] Add a functional repro for the superseded deferred un-assign Manufactures the interleaving deterministically against a real broker: a revoke deferred behind an in-flight task, then a partition count change observed by a group-updating metadata request, which makes librdkafka rejoin and apply a fresh assignment while the revoke is still unanswered. The deferred un-assign must then leave that assignment alone. Two variants differ only in which request observes the change: an application getMetadata() call, and librdkafka's own periodic topic.metadata.refresh.interval.ms refresh, which fires with no application metadata call at all. Both are red without the rebalance guard (the consumer ends with an empty local assignment, ready but consuming nothing) and green with it. Issue: BB-835 --- .../lib/BackbeatConsumerRebalanceSupersede.js | 263 ++++++++++++++++++ 1 file changed, 263 insertions(+) create mode 100644 tests/functional/lib/BackbeatConsumerRebalanceSupersede.js diff --git a/tests/functional/lib/BackbeatConsumerRebalanceSupersede.js b/tests/functional/lib/BackbeatConsumerRebalanceSupersede.js new file mode 100644 index 000000000..a91b65ccd --- /dev/null +++ b/tests/functional/lib/BackbeatConsumerRebalanceSupersede.js @@ -0,0 +1,263 @@ +const assert = require('assert'); +const { promisify } = require('util'); +const kafka = require('node-rdkafka'); +const sinon = require('sinon'); + +const BackbeatProducer = require('../../../lib/BackbeatProducer'); +const BackbeatConsumer = require('../../../lib/BackbeatConsumer'); +const { withTopicPrefix } = require('../../../lib/util/topic'); + +const zookeeperConf = { connectionString: 'localhost:2181' }; +const kafkaConf = { hosts: 'localhost:9092' }; + +const { ERR__ASSIGN_PARTITIONS, ERR__REVOKE_PARTITIONS } = kafka.CODES.ERRORS; + +const sleep = ms => new Promise(resolve => setTimeout(resolve, ms)); + +function deferred() { + let resolve; + const promise = new Promise(res => { + resolve = res; + }); + return { promise, resolve }; +} + +class InstrumentedConsumer extends BackbeatConsumer { + _onRebalance(err, assignment) { + if (!this.rebalanceLog) { + this.rebalanceLog = []; + } + this.rebalanceLog.push(err.code); + super._onRebalance(err, assignment); + } +} + +// When a revoke arrives with work in flight, the un-assign is deferred +// until the drain completes. librdkafka however answers a subscribed-topic +// metadata change with an immediate rejoin even while that revoke is +// still unanswered, so a fresh assignment can be granted and applied +// before the drain completes -- and the deferred un-assign then wipes it, +// leaving a live group member that owns partitions at the broker but +// consumes nothing, with no rebalance ever coming to fix it. +// +// Both tests manufacture that interleaving with partition-count bumps, +// differing only in which metadata request observes the change: an +// application getMetadata() call, or librdkafka's own periodic refresh, +// which runs regardless of what the application does. If librdkafka ever +// stops rejoining while a revoke is unanswered, the setup fails on 'no +// assign delivered' rather than on the final consumption assertion. +describe('BackbeatConsumer superseded deferred un-assign', function testSuite() { + this.timeout(120000); + + let admin; + let producer; + let consumer; + let rawTopic; + let fullTopic; + let taskStarted; + let taskGate; + let unassigns; + let consumedAfterRelease; + + function queueProcessor(message, cb) { + const value = message.value.toString(); + if (value === 'hold') { + taskStarted.resolve(); + taskGate.promise.then(() => cb()); + return; + } + consumedAfterRelease.push(value); + process.nextTick(cb); + } + + beforeEach(() => { + rawTopic = `backbeat-supersede-spec-${Date.now()}`; + fullTopic = withTopicPrefix(rawTopic); + taskStarted = deferred(); + taskGate = deferred(); + unassigns = []; + consumedAfterRelease = []; + admin = kafka.AdminClient.create({ + 'client.id': 'supersede-spec-admin', + 'metadata.broker.list': kafkaConf.hosts, + }); + }); + + afterEach(async function teardown() { + this.timeout(40000); + sinon.restore(); + taskGate.resolve(); + if (consumer) { + // a consumer stalled by the bug may misbehave on close(); + // never let teardown mask the test outcome + await Promise.race([ + promisify(consumer.close.bind(consumer))(), + sleep(15000), + ]); + } + if (producer) { + await promisify(producer.close.bind(producer))(); + } + try { + admin.disconnect(); + } catch { + // already disconnected + } + consumer = null; + producer = null; + }); + + async function start({ refreshIntervalMs } = {}) { + await promisify(admin.createTopic.bind(admin))({ + topic: fullTopic, + /* eslint-disable camelcase */ + num_partitions: 1, + replication_factor: 1, + /* eslint-enable camelcase */ + }, 15000); + + if (refreshIntervalMs) { + // BackbeatConsumer has no rdkafka config passthrough, so + // shrink the periodic refresh by wrapping the constructor + const RealKafkaConsumer = kafka.KafkaConsumer; + sinon.replace(kafka, 'KafkaConsumer', + // an arrow function cannot be invoked with `new` + // eslint-disable-next-line prefer-arrow-callback + function kafkaConsumerWithFastRefresh(conf, topicConf) { + return new RealKafkaConsumer({ + ...conf, + 'topic.metadata.refresh.interval.ms': refreshIntervalMs, + }, topicConf); + }); + } + consumer = new InstrumentedConsumer({ + clientId: 'BackbeatConsumer-supersede', + zookeeper: zookeeperConf, + kafka: kafkaConf, + groupId: `supersede-group-${Math.random()}`, + topic: rawTopic, + queueProcessor, + fromOffset: 'earliest', + // rebalance callbacks are delivered by the consume poll, + // which stops while the pipeline is full: a free slot must + // remain next to the held task for the revoke to reach us + concurrency: 2, + }); + sinon.restore(); + consumer.on('unassign', status => unassigns.push(status)); + + producer = new BackbeatProducer({ + kafka: kafkaConf, + topic: rawTopic, + pollIntervalMs: 100, + }); + await Promise.all([ + new Promise(resolve => consumer.on('ready', resolve)), + new Promise(resolve => producer.on('ready', resolve)), + ]); + consumer.subscribe(); + } + + const send = messages => + promisify(producer.send.bind(producer))(messages); + + const sawRevoke = () => + (consumer.rebalanceLog || []).includes(ERR__REVOKE_PARTITIONS); + + const sawAssignAfterRevoke = () => { + const log = consumer.rebalanceLog || []; + const revokeIdx = log.indexOf(ERR__REVOKE_PARTITIONS); + return revokeIdx !== -1 && + log.slice(revokeIdx + 1).includes(ERR__ASSIGN_PARTITIONS); + }; + + // wait for cond; when poking, issue the same full-cluster metadata + // request BackbeatConsumer itself makes, which librdkafka answers + // with a consumer-group subscription re-check + async function driveUntil(cond, poke, what) { + for (let i = 0; i < 30; i++) { + if (cond()) { + return; + } + if (poke) { + await promisify(consumer.getMetadata.bind(consumer))( + { allTopics: true, timeout: 10000 }).catch(() => {}); + } + await sleep(500); + } + assert.fail(what); + } + + async function waitUntil(cond, timeoutMs, what) { + const deadline = Date.now() + timeoutMs; + while (!cond()) { + if (Date.now() > deadline) { + assert.fail(typeof what === 'function' ? what() : what); + } + await sleep(100); + } + } + + const bump = (partitions => + promisify(admin.createPartitions.bind(admin))( + fullTopic, partitions, 15000)); + + async function runScenario({ poke }) { + // hold one task in flight so the next revoke defers its un-assign + await send([{ key: 'k-hold', message: 'hold' }]); + await taskStarted.promise; + + // first partition bump: the next group-updating metadata + // response delivers a normal revoke + await bump(2); + await driveUntil(sawRevoke, poke, + 'no revoke delivered after the first partition bump'); + + await sleep(300); + assert.strictEqual(unassigns.length, 0, + 'the un-assign should be deferred while a task is in flight'); + assert(consumer._drainProcessQueueTimeout, + 'the drain watchdog should be armed while the drain is pending'); + + // second bump while the revoke is unanswered: librdkafka rejoins + // and grants a fresh assignment before the drain completes + await bump(3); + await driveUntil(sawAssignAfterRevoke, poke, + 'no assign delivered while the revoke was parked: the ' + + 'trigger is not reproducible on this librdkafka version'); + + assert.strictEqual(unassigns.length, 0, + 'the drain must still be pending when the fresh assignment lands'); + + // release the held task: the drain completes and the deferred + // un-assign fires -- it must not touch the superseding assignment + taskGate.resolve(); + await waitUntil(() => consumer._processingQueue.idle(), 10000, + 'processing queue drain'); + await sleep(1500); + + await send([ + { key: 'a', message: 'after-1' }, + { key: 'b', message: 'after-2' }, + { key: 'c', message: 'after-3' }, + ]); + await waitUntil(() => consumedAfterRelease.length >= 3, 20000, + () => 'consumer stalled after the deferred un-assign: ' + + `local assignment=[${consumer._consumer.assignments() + .map(a => a.partition)}], ` + + `isReady=${consumer.isReady()}, ` + + `consumed=${JSON.stringify(consumedAfterRelease)}`); + } + + it('keeps the assignment granted while draining, when an application ' + + 'metadata call observes the topic change', async () => { + await start(); + await runScenario({ poke: true }); + }); + + it('keeps the assignment granted while draining, when the periodic ' + + 'metadata refresh observes the topic change', async () => { + await start({ refreshIntervalMs: 1000 }); + await runScenario({ poke: false }); + }); +});