From 6c6acdae72fc57997e2b4e5baedaf18fa4224161 Mon Sep 17 00:00:00 2001 From: Thomas Flament Date: Wed, 19 Aug 2026 18:30:05 +0200 Subject: [PATCH] Ignore a deferred un-assign superseded by a later rebalance On ERR__REVOKE_PARTITIONS the un-assign is deferred until the processing queue and the offset ledger have drained. If the next rebalance granted the partitions back before that happened, the deferred callback still ran and un-assigned them: the consumer then owned partitions at the broker with no local assignment, and since group membership had not changed, nothing triggered another rebalance to rescue it. Track a rebalance id, bumped on every rebalance event, and give up on a deferred un-assign whose id no longer matches. The check runs again before un-assigning, as publishing offsets to zookeeper in between is asynchronous and leaves a second window for the partitions to come back. The drain watchdog is now cleared on every rebalance rather than only on assignment, so a timer armed by a superseded revoke can no longer disconnect a consumer that is not stuck. Issue: BB-835 --- lib/BackbeatConsumer.js | 62 +++++++++++ lib/constants.js | 1 + tests/unit/backbeatConsumer.js | 184 +++++++++++++++++++++++++++++++++ 3 files changed, 247 insertions(+) diff --git a/lib/BackbeatConsumer.js b/lib/BackbeatConsumer.js index 8df60c538..e0fd50d8a 100644 --- a/lib/BackbeatConsumer.js +++ b/lib/BackbeatConsumer.js @@ -176,6 +176,9 @@ class BackbeatConsumer extends EventEmitter { } }); + // a deferred un-assign gives up when this no longer matches + this._rebalanceId = 0; + this._messagesConsumed = 0; // this variable represents how many kafka messages have been // requested without having been received yet, i.e. still @@ -752,15 +755,45 @@ class BackbeatConsumer extends EventEmitter { } } + /** + * Run a shutdown/rebalance step that must not abort the sequence it + * belongs to, logging rather than throwing. + * + * @param {string} op - librdkafka operation name, for the log line + * @param {function} fn - the call to attempt + * @returns {undefined} + */ + _bestEffort(op, fn) { + try { + fn(); + } catch (e) { + // ERR__STATE just means the client moved on without us + const logger = this._consumer.isConnected() && + e.code !== kafka.CODES.ERRORS.ERR__STATE ? this._log.error : this._log.info; + logger.bind(this._log)(`rdkafka.${op} failed`, { + e: e.toString(), + topic: this._topic, + groupId: this._groupId, + }); + } + } + /** * @param {kafka.KafkaError} err Rebalance event * @param {TopicPartition[]} assignment List of (un)assigned partitions * @returns {void} */ _onRebalance(err, assignment) { + const rebalanceId = ++this._rebalanceId; + + clearTimeout(this._drainProcessQueueTimeout); + this._drainProcessQueueTimeout = null; + if (err.code === kafka.CODES.ERRORS.ERR__ASSIGN_PARTITIONS) { this._log.info('rdkafka.assign', { assignment }); + this._setDrain(null); + try { this._consumer.assign(assignment); if (this._circuitBreaker.state !== BreakerState.Nominal) { @@ -779,7 +812,26 @@ class BackbeatConsumer extends EventEmitter { ledger: this._offsetLedger.getProcessingCount(this._topic), }); + const isSuperseded = () => rebalanceId !== this._rebalanceId; + const skipSuperseded = status => { + this._log.info('skipping superseded un-assign', { + status, + rebalanceId, + currentRebalanceId: this._rebalanceId, + topic: this._topic, + groupId: this._groupId, + }); + KafkaBacklogMetrics.onRebalance( + this._topic, this._groupId, unassignStatus.SUPERSEDED); + }; + const unassign = jsutil.once(status => { + // before touching state that now belongs to a later rebalance + if (isSuperseded()) { + skipSuperseded(status); + return; + } + this._log.info(`processing queue ${status}, un-assigning`, { queueLen: this._processingQueue.length(), running: this._processingQueue.running(), @@ -804,6 +856,12 @@ class BackbeatConsumer extends EventEmitter { } const doUnassign = () => { + // re-checked: publishing offsets above is asynchronous + if (isSuperseded()) { + skipSuperseded(status); + return; + } + this._resumePausedPartitions(); try { @@ -889,6 +947,10 @@ class BackbeatConsumer extends EventEmitter { }, this._maxPollIntervalMs - 1000); // 1 second earlier, to be within the limit } else { this._log.error('rdkafka.rebalance', { err, assignment }); + // the bump above just superseded whatever revoke was pending and + // dropped its watchdog, so nothing else will answer this callback, + // and librdkafka requires one: assign(NULL) synchronises the state + this._bestEffort('unassign', () => this._consumer.unassign()); } } diff --git a/lib/constants.js b/lib/constants.js index 4148feb2a..183359b29 100644 --- a/lib/constants.js +++ b/lib/constants.js @@ -25,6 +25,7 @@ const constants = { IDLE: 'idle', DRAINED: 'drained', TIMEOUT: 'timeout', + SUPERSEDED: 'superseded', }, statusReady: 'READY', statusUndefined: 'UNDEFINED', diff --git a/tests/unit/backbeatConsumer.js b/tests/unit/backbeatConsumer.js index 951fad035..afc7d7238 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,186 @@ 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 drainCallbacks; + let queueIdle; + let ledgerCount; + + 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; + drainCallbacks = []; + consumer._processingQueue = { + length: () => 0, + running: () => (queueIdle ? 0 : 1), + idle: () => queueIdle, + setDrain: func => drainCallbacks.push(func), + }; + consumer._offsetLedger.getProcessingCount = () => ledgerCount; + + sinon.stub(KafkaBacklogMetrics, 'onRebalance'); + }); + + afterEach(() => { + clearTimeout(consumer._drainProcessQueueTimeout); + sinon.restore(); + }); + + const completeDrain = () => { + queueIdle = true; + ledgerCount = 0; + consumer._drainCallback(); + }; + + it('should un-assign once the drain completes', done => { + consumer.on('unassign', status => { + assert.strictEqual(status, unassignStatus.DRAINED); + assert(consumer._consumer.unassign.calledOnce); + done(); + }); + + consumer._onRebalance(REVOKE, partitions); + assert(consumer._consumer.unassign.notCalled); + completeDrain(); + }); + + it('should un-assign immediately when nothing is in flight', done => { + queueIdle = true; + ledgerCount = 0; + + consumer.on('unassign', status => { + assert.strictEqual(status, unassignStatus.IDLE); + assert(consumer._consumer.unassign.calledOnce); + done(); + }); + + consumer._onRebalance(REVOKE, partitions); + }); + + it('should not un-assign when a new assignment arrived while draining', () => { + consumer._onRebalance(REVOKE, partitions); + const deferredUnassign = drainCallbacks[drainCallbacks.length - 1]; + + // the next generation grants the partitions back mid-drain + consumer._onRebalance(ASSIGN, partitions); + assert(consumer._consumer.assign.calledOnce); + + queueIdle = true; + ledgerCount = 0; + deferredUnassign(); + + assert(consumer._consumer.unassign.notCalled); + assert(KafkaBacklogMetrics.onRebalance.calledWith( + 'my-test-topic', 'unittest-group', unassignStatus.SUPERSEDED)); + }); + + it('should synchronise the assignment on an arbitrary rebalance error', + () => { + // the bump above superseded whatever revoke was pending and + // dropped its watchdog, so nothing else answers this callback + consumer._onRebalance({ code: -1 }, partitions); + + assert(consumer._consumer.unassign.calledOnce); + }); + + it('should not un-assign when a later revoke superseded the drain', () => { + consumer._onRebalance(REVOKE, partitions); + const firstUnassign = drainCallbacks[drainCallbacks.length - 1]; + + consumer._onRebalance(REVOKE, partitions); + + queueIdle = true; + ledgerCount = 0; + firstUnassign(); + + assert(consumer._consumer.unassign.notCalled); + }); + + 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(KafkaBacklogMetrics.onRebalance.calledWith( + 'my-test-topic', 'unittest-group', unassignStatus.SUPERSEDED)); + }); + + it('should not leave a superseded revoke watchdog armed', () => { + const clock = sinon.useFakeTimers(); + try { + consumer._onRebalance(REVOKE, partitions); + consumer._onRebalance(REVOKE, partitions); + + clock.tick(consumer._maxPollIntervalMs + 1000); + assert(consumer._consumer.disconnect.calledOnce); + } finally { + clock.restore(); + } + }); + + it('should leave the current drain and timeout armed when a superseded ' + + 'un-assign fires', () => { + consumer._onRebalance(REVOKE, partitions); + const supersededUnassign = drainCallbacks[drainCallbacks.length - 1]; + + consumer._onRebalance(ASSIGN, partitions); + consumer._onRebalance(REVOKE, partitions); + + const currentDrain = consumer._drainCallback; + const currentTimeout = consumer._drainProcessQueueTimeout; + assert.notStrictEqual(currentDrain, null); + assert.notStrictEqual(currentTimeout, null); + + // or the callback returns before reaching the guard + queueIdle = true; + ledgerCount = 0; + supersededUnassign(); + + assert(KafkaBacklogMetrics.onRebalance.calledWith( + 'my-test-topic', 'unittest-group', unassignStatus.SUPERSEDED)); + + assert.strictEqual(consumer._drainCallback, currentDrain); + assert.strictEqual(consumer._drainProcessQueueTimeout, currentTimeout); + assert(consumer._consumer.unassign.notCalled); + }); + }); });