-
Notifications
You must be signed in to change notification settings - Fork 24
Leave the consumer group promptly on shutdown #2838
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||
|---|---|---|---|---|---|---|
| @@ -1,7 +1,6 @@ | ||||||
| const { EventEmitter } = require('events'); | ||||||
| const kafka = require('node-rdkafka'); | ||||||
| const assert = require('assert'); | ||||||
| const async = require('async'); | ||||||
| const joi = require('joi'); | ||||||
| const jsutil = require('arsenal').jsutil; | ||||||
| const Logger = require('werelogs').Logger; | ||||||
|
|
@@ -20,6 +19,8 @@ const { | |||||
| backbeatConsumer: { | ||||||
| CONCURRENCY_DEFAULT, | ||||||
| MAX_QUEUED_DEFAULT, | ||||||
| DISCONNECT_TIMEOUT_MS, | ||||||
| SHUTDOWN_DRAIN_GRACE_MS, | ||||||
| } | ||||||
| } = require('./constants'); | ||||||
|
|
||||||
|
|
@@ -30,6 +31,9 @@ const { startLinkedSpanFromKafkaEntry } = require('arsenal/build/lib/tracing').k | |||||
| const CLIENT_ID = 'BackbeatConsumer'; | ||||||
| const { withTopicPrefix } = require('./util/topic'); | ||||||
|
|
||||||
| // Resolve after `ms` (call-time setTimeout so test fake timers apply). | ||||||
| const delay = ms => new Promise(resolve => setTimeout(resolve, ms)); | ||||||
|
|
||||||
| /** | ||||||
| * Stats on how we are consuming Kafka | ||||||
| * @typedef {Object} ConsumerStats | ||||||
|
|
@@ -195,6 +199,15 @@ class BackbeatConsumer extends EventEmitter { | |||||
| // since the last consume, and is used to trigger a new consume | ||||||
| // immediately after the current one | ||||||
| this._tasksCompletedSinceLastConsume = false; | ||||||
| // Set at close() entry: stops fetching, declines new grants | ||||||
| this._closing = false; | ||||||
| // Set at the disconnect step: rebalance callbacks are answered | ||||||
| // at once, since the client's close blocks on them | ||||||
| this._disconnecting = false; | ||||||
| // True between receiving a revoke and answering it (draining) | ||||||
| this._waitingDrain = false; | ||||||
| // Bounds close()'s wait for the drain | ||||||
| this._closeDrainTimeout = null; | ||||||
|
|
||||||
| /** @type {ConsumerStats} */ | ||||||
| this.consumerStats = { lag: {} }; | ||||||
|
|
@@ -432,6 +445,13 @@ class BackbeatConsumer extends EventEmitter { | |||||
| this._tryConsumedTimeout = null; | ||||||
| } | ||||||
|
|
||||||
| // stop fetching once closing: new entries would refill the | ||||||
| // pipeline close() is draining, and this ends the | ||||||
| // self-rescheduling consume loop that would poll a closed client | ||||||
| if (this._closing) { | ||||||
| return undefined; | ||||||
| } | ||||||
|
|
||||||
| // use non-flowing mode of consumption to add some flow | ||||||
| // control: explicit consumption of messages is required, | ||||||
| // needs explicit polling to get new messages | ||||||
|
|
@@ -759,6 +779,19 @@ class BackbeatConsumer extends EventEmitter { | |||||
| */ | ||||||
| _onRebalance(err, assignment) { | ||||||
| if (err.code === kafka.CODES.ERRORS.ERR__ASSIGN_PARTITIONS) { | ||||||
| if (this._closing) { | ||||||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. not needed, since we |
||||||
| // Decline the grant, but still answer the callback or the | ||||||
| // client stays parked. assign() is refused mid-disconnect, | ||||||
| // so answer with unassign(). | ||||||
| this._log.info('rdkafka.assign declined, shutting down', | ||||||
| { assignment }); | ||||||
| this._bestEffort('unassign', | ||||||
| () => this._consumer.unassign()); | ||||||
| if (!this._waitingDrain) { | ||||||
| this.emit('unassign', unassignStatus.SHUTDOWN); | ||||||
| } | ||||||
| return; | ||||||
| } | ||||||
| this._log.info('rdkafka.assign', { assignment }); | ||||||
|
|
||||||
| try { | ||||||
|
|
@@ -779,6 +812,19 @@ class BackbeatConsumer extends EventEmitter { | |||||
| ledger: this._offsetLedger.getProcessingCount(this._topic), | ||||||
| }); | ||||||
|
|
||||||
| if (this._disconnecting) { | ||||||
| // The client's close delivers this revoke and blocks until | ||||||
| // answered; entries can no longer commit, so answer at | ||||||
| // once instead of draining. | ||||||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. correct in term of kafka (will be replayed anyway); but it is better to abort the current task or finish them (even if we don't commit)? |
||||||
| KafkaBacklogMetrics.onRebalance(this._topic, this._groupId, | ||||||
| unassignStatus.SHUTDOWN); | ||||||
| this._bestEffort('unassign', | ||||||
| () => this._consumer.unassign()); | ||||||
| this.emit('unassign', unassignStatus.SHUTDOWN); | ||||||
| return; | ||||||
| } | ||||||
|
|
||||||
| this._waitingDrain = true; | ||||||
| const unassign = jsutil.once(status => { | ||||||
| this._log.info(`processing queue ${status}, un-assigning`, { | ||||||
| queueLen: this._processingQueue.length(), | ||||||
|
|
@@ -814,6 +860,7 @@ class BackbeatConsumer extends EventEmitter { | |||||
| logger.bind(this._log)('rdkafka.unassign failed', { e: e.toString(), assignment }); | ||||||
| } | ||||||
|
|
||||||
| this._waitingDrain = false; | ||||||
| this.emit('unassign', status); | ||||||
| }; | ||||||
|
|
||||||
|
|
@@ -1144,6 +1191,47 @@ class BackbeatConsumer extends EventEmitter { | |||||
| return subscription.length === 0 ? null : subscription; | ||||||
| } | ||||||
|
|
||||||
| /** | ||||||
| * Read the consumer's current partition assignment. | ||||||
| * | ||||||
| * @return {Object[]|null} assigned partitions, or null if unavailable | ||||||
| */ | ||||||
| _getAssignments() { | ||||||
| try { | ||||||
| return this._consumer.assignments(); | ||||||
| } catch (err) { | ||||||
| this._log.debug('could not read consumer assignments', { | ||||||
| method: 'BackbeatConsumer._getAssignments', | ||||||
| topic: this._topic, | ||||||
| groupId: this._groupId, | ||||||
| error: err.message, | ||||||
| }); | ||||||
| return null; | ||||||
| } | ||||||
| } | ||||||
|
|
||||||
| /** | ||||||
| * Run a consumer operation that may fail while closing (state-dependent | ||||||
| * calls throw ERR__STATE mid-disconnect). | ||||||
| * | ||||||
| * @param {string} op - operation name, for logging | ||||||
| * @param {function} fn - operation to run | ||||||
| * @returns {undefined} | ||||||
| */ | ||||||
| _bestEffort(op, fn) { | ||||||
| try { | ||||||
| fn(); | ||||||
| } catch (e) { | ||||||
| const logger = this._consumer.isConnected() ? | ||||||
| this._log.error : this._log.debug; | ||||||
| logger.bind(this._log)(`rdkafka.${op} failed`, { | ||||||
| topic: this._topic, | ||||||
| groupId: this._groupId, | ||||||
| e: e.toString(), | ||||||
| }); | ||||||
| } | ||||||
| } | ||||||
|
|
||||||
| /** | ||||||
| * Helper method to wait for consumers to be assigned to partitions before | ||||||
| * proceeding to send canary messages. | ||||||
|
|
@@ -1221,42 +1309,111 @@ class BackbeatConsumer extends EventEmitter { | |||||
| * @return {undefined} | ||||||
| */ | ||||||
| close(cb) { | ||||||
| // Keep the callback signature for existing callers; errors are | ||||||
| // swallowed as before (cb takes no arguments). | ||||||
| this._close().then(() => cb(), () => cb()); | ||||||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Suggested change
|
||||||
| } | ||||||
|
|
||||||
| async _close() { | ||||||
| this._closing = true; | ||||||
| if (this._publishOffsetsCronTimer) { | ||||||
| clearInterval(this._publishOffsetsCronTimer); | ||||||
| this._publishOffsetsCronTimer = null; | ||||||
| } | ||||||
| if (this._publishOffsetsCronActive) { | ||||||
| return setTimeout(() => this.close(cb), 1000); | ||||||
| // Wait briefly for an in-flight offset publish to finish, but never | ||||||
| // block shutdown on it: after the grace, proceed and let it error out. | ||||||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. why would be it delayed more than 1s? |
||||||
| const publishDeadline = Date.now() + SHUTDOWN_DRAIN_GRACE_MS; | ||||||
| while (this._publishOffsetsCronActive && Date.now() < publishDeadline) { | ||||||
| await delay(1000); | ||||||
|
Comment on lines
+1326
to
+1327
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. if we go that road, should we not use some sort of even (or callback), instead of a polling loop? |
||||||
| } | ||||||
| this._circuitBreaker.stop(); | ||||||
| return async.waterfall([ | ||||||
| next => { | ||||||
| if (this._consumer?.isConnected()) { | ||||||
| const subscription = this._getSubscription(); | ||||||
| if (subscription !== null) { | ||||||
| this._consumer.unsubscribe(); | ||||||
| // Wait for partition unassign to complete before | ||||||
| // disconnecting, the rebalance callback will handle | ||||||
| // waiting for current jobs to complete as well as commit | ||||||
| // the latest offsets | ||||||
| this.once('unassign', () => next()); | ||||||
| return; | ||||||
| } | ||||||
| } | ||||||
| process.nextTick(next); | ||||||
| }, | ||||||
| next => { | ||||||
| if (this._zookeeper) { | ||||||
| this._zookeeper.close(); | ||||||
| } | ||||||
| if (this._consumer?.isConnected()) { | ||||||
| this._consumer.disconnect(); | ||||||
| this._consumer.once('disconnected', () => next()); | ||||||
| } else { | ||||||
| process.nextTick(next); | ||||||
|
|
||||||
| if (this._consumer?.isConnected()) { | ||||||
| // Best-effort unsubscribe: a no-op when not subscribed, and the | ||||||
| // trigger for the drain revoke when we hold partitions. | ||||||
| // Unconditional so a transient ERR__STATE on subscription() can't | ||||||
| // skip it. | ||||||
| this._bestEffort('unsubscribe', | ||||||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. nit: this |
||||||
| () => this._consumer.unsubscribe()); | ||||||
| const assignments = this._getAssignments(); | ||||||
| const holdsPartitions = | ||||||
| assignments !== null && assignments.length > 0; | ||||||
| // Wait for the drain only when a revoke is coming: one already | ||||||
| // draining, or triggered by the unsubscribe because we still hold | ||||||
| // partitions. Otherwise no 'unassign' arrives and disconnect() | ||||||
| // leaves the group on its own, so waiting would hang. | ||||||
| if (this._waitingDrain || holdsPartitions) { | ||||||
| await this._waitForShutdownUnassign(); | ||||||
| } | ||||||
| } | ||||||
|
|
||||||
| // From here rebalance callbacks are answered at once: disconnect() | ||||||
| // blocks until every pending one is answered. | ||||||
| this._disconnecting = true; | ||||||
| if (this._zookeeper) { | ||||||
| this._zookeeper.close(); | ||||||
| } | ||||||
| if (this._consumer?.isConnected()) { | ||||||
| await this._disconnectWithTimeout(); | ||||||
| } | ||||||
| } | ||||||
|
|
||||||
| /** | ||||||
| * Wait for the revoke handler to drain current jobs and emit 'unassign'. | ||||||
| * Once a revoke is draining it emits on its own -- on drain completion or | ||||||
| * its watchdog at maxPollIntervalMs - 1000 -- so just wait for it. Only | ||||||
| * bound the wait for a revoke to start: if the unsubscribe was postponed | ||||||
| * by an in-flight rebalance none arrives, so after a short grace give up | ||||||
| * and let disconnect() perform the leave. | ||||||
| * @return {Promise<undefined>} resolves once un-assigned or the grace ends | ||||||
| */ | ||||||
| _waitForShutdownUnassign() { | ||||||
| return new Promise(resolve => { | ||||||
| const done = jsutil.once(() => { | ||||||
| this.removeListener('unassign', done); | ||||||
| clearTimeout(this._closeDrainTimeout); | ||||||
| this._closeDrainTimeout = null; | ||||||
| resolve(); | ||||||
| }); | ||||||
| this.once('unassign', done); | ||||||
| this._closeDrainTimeout = setTimeout(() => { | ||||||
| if (this._waitingDrain) { | ||||||
| return; | ||||||
| } | ||||||
| }, | ||||||
| ], () => cb()); | ||||||
| this._log.warn('closing: no un-assign after unsubscribe, ' | ||||||
| + 'proceeding with disconnect', { | ||||||
| topic: this._topic, | ||||||
| groupId: this._groupId, | ||||||
| }); | ||||||
| this._bestEffort('unassign', | ||||||
| () => this._consumer.unassign()); | ||||||
| done(); | ||||||
| }, SHUTDOWN_DRAIN_GRACE_MS); | ||||||
| }); | ||||||
| } | ||||||
|
|
||||||
| /** | ||||||
| * Disconnect the client, bounded by DISCONNECT_TIMEOUT_MS. | ||||||
| * @return {Promise<undefined>} resolves once disconnected or the bound ends | ||||||
| */ | ||||||
| _disconnectWithTimeout() { | ||||||
| return new Promise(resolve => { | ||||||
| const done = jsutil.once(() => { | ||||||
| clearTimeout(timeout); | ||||||
| resolve(); | ||||||
| }); | ||||||
| const timeout = setTimeout(() => { | ||||||
| this._log.error('closing: disconnect timed out', { | ||||||
| topic: this._topic, | ||||||
| groupId: this._groupId, | ||||||
| timeoutMs: DISCONNECT_TIMEOUT_MS, | ||||||
| }); | ||||||
| done(); | ||||||
| }, DISCONNECT_TIMEOUT_MS); | ||||||
| this._consumer.once('disconnected', done); | ||||||
| this._bestEffort('disconnect', | ||||||
| () => this._consumer.disconnect()); | ||||||
| }); | ||||||
| } | ||||||
| } | ||||||
|
|
||||||
|
|
||||||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
does not seem like 4 states : but a single state/automaton