Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
217 changes: 187 additions & 30 deletions lib/BackbeatConsumer.js
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;
Expand All @@ -20,6 +19,8 @@ const {
backbeatConsumer: {
CONCURRENCY_DEFAULT,
MAX_QUEUED_DEFAULT,
DISCONNECT_TIMEOUT_MS,
SHUTDOWN_DRAIN_GRACE_MS,
}
} = require('./constants');

Expand All @@ -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
Expand Down Expand Up @@ -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;
Comment on lines +202 to +210

Copy link
Copy Markdown
Contributor

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


/** @type {ConsumerStats} */
this.consumerStats = { lag: {} };
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -759,6 +779,19 @@ class BackbeatConsumer extends EventEmitter {
*/
_onRebalance(err, assignment) {
if (err.code === kafka.CODES.ERRORS.ERR__ASSIGN_PARTITIONS) {
if (this._closing) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

not needed, since we unsubscribe() first ?

// 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 {
Expand All @@ -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.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The 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)?
e.g. most of our tasks are idempotent, and replaying them is fine: for exemple aborting replication wastes some bandwidth, while replaying it after it completes is mostly no-op...

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(),
Expand Down Expand Up @@ -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);
};

Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
this._close().then(() => cb(), () => cb());
this._close().then(cb, () => cb());

}

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.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

why would be it delayed more than 1s?
(what is the cron timer period)

const publishDeadline = Date.now() + SHUTDOWN_DRAIN_GRACE_MS;
while (this._publishOffsetsCronActive && Date.now() < publishDeadline) {
await delay(1000);
Comment on lines +1326 to +1327

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The 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?
(and this is probably a different PR than handling close during rebalance)

}
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',

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: this bestEffort function really does not help readability....

() => 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());
});
}
}

Expand Down
10 changes: 10 additions & 0 deletions lib/constants.js
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ const constants = {
IDLE: 'idle',
DRAINED: 'drained',
TIMEOUT: 'timeout',
SHUTDOWN: 'shutdown',
},
statusReady: 'READY',
statusUndefined: 'UNDEFINED',
Expand Down Expand Up @@ -87,6 +88,15 @@ const constants = {
CONCURRENCY_DEFAULT: 1,
// controls the max number of messages to queue for processing.
MAX_QUEUED_DEFAULT: 1000,
// Bounds close()'s disconnect step: rd_kafka_consumer_close() only
// returns once its LeaveGroup is acknowledged, which can stall past
// the pod's grace period when that request is queued behind an
// in-flight JoinGroup. Return anyway after this so shutdown completes.
DISCONNECT_TIMEOUT_MS: 5000,
// Grace for an unsubscribe to produce its revoke before close()
// gives up and disconnects; only bounds the wait for a revoke to
// start, one in progress drains under the full budget.
SHUTDOWN_DRAIN_GRACE_MS: 5000,
},
// Objects below this size skip the Expect: 100-continue handshake as the
// extra round-trip is a bigger overhead than the bandwidth saved
Expand Down
Loading
Loading