[
https://issues.apache.org/jira/browse/KAFKA-21062?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Michael Westerby reassigned KAFKA-21062:
----------------------------------------
Assignee: Michael Westerby
> Potential data loss during KRaft migration due to a lost ZooKeeper
> acknowledgement on /migration write
> ------------------------------------------------------------------------------------------------------
>
> Key: KAFKA-21062
> URL: https://issues.apache.org/jira/browse/KAFKA-21062
> Project: Kafka
> Issue Type: Bug
> Components: kraft, migration, zkclient
> Affects Versions: 3.9.2
> Reporter: Michael Westerby
> Assignee: Michael Westerby
> Priority: Major
>
> A lost ZooKeeper acknowledgment on a successful {{/migration}} znode write
> can cause the KRaft migration driver to mistake its own prior write for a
> conflicting one, permanently stalling it once it reaches the {{DUAL_WRITE}}
> stage. This prevents further writes to ZooKeeper, allowing the state in
> between KRaft and ZooKeeper to diverge during the migration.
> This happens because the client the migration uses to talk to ZooKeeper will
> retry an identical request to {{/migration}} if it sees a connection loss. If
> the original request had actually already succeeded on the server, but only
> its acknowledgment was lost, the retry sends the same expected version of the
> {{/migration}} znode again. Since that version has already moved on,
> ZooKeeper rejects the retry as a version conflict, even though nothing is
> actually wrong.
> The migration driver currently has no way to tell this apart from a genuine
> conflict caused by a second controller writing at the same time, so it treats
> every such rejection as a hard failure. Because the write is treated as a
> failure, the driver's own record of what version the {{/migration}} znode is
> currently at never gets corrected, and stays at the old, stale value. Every
> subsequent write then reuses that same stale value and hits the same
> rejection, indefinitely. The driver becomes permanently stuck, unable to make
> further migration progress, until an unrelated event, such as a new
> controller election, forces an unconditional resynchronization of
> {{{}/migration{}}}.
> While stuck, ZooKeeper stops receiving updates from the controller, and any
> ZK-mode brokers that rely on ZooKeeper for their view of cluster state are
> never notified of the corresponding changes via RPCs either. This causes the
> states in ZooKeeper and KRaft to diverge, which can destabilize the cluster
> further.
> This is of particular concern when some, but not all, brokers have already
> been restarted into their KRaft mode. The remaining ZK-mode brokers are left
> running on an increasingly outdated view of the cluster state. In the worst
> case, this can cause real data loss if the KRaft side elects a new partition
> leader while the driver is stuck. As that change is never propagated to
> ZooKeeper, and the ZK-mode brokers are never notified of it via RPC either,
> the old leader can continue to operate under the stale, lower epoch. The old
> leader can keep accepting and acknowledging produce requests it is no longer
> authorized to serve. When it eventually discovers the true, higher-epoch
> leader, it must truncate its own log to reconcile, silently discarding any
> writes it had already acknowledged to producers in the meantime.
> h2. Walkthrough
> [{{KafkaZkClient}}
> |https://github.com/apache/kafka/blob/3.9.2/core/src/main/scala/kafka/zk/KafkaZkClient.scala]exposes
> two methods used to write to the {{/migration}} znode in ZooKeeper during a
> ZK-to-KRaft migration:
> # the standalone
> [{{updateMigrationState}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/core/src/main/scala/kafka/zk/KafkaZkClient.scala#L1767]
> # the bundled
> [{{{}retryMigrationRequestsUntilConnected{}}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/core/src/main/scala/kafka/zk/KafkaZkClient.scala#L2012].
> Both take a
> [{{ZkMigrationLeadershipState}}|https://github.com/apache/kafka/blob/3.9.2/metadata/src/main/java/org/apache/kafka/metadata/migration/ZkMigrationLeadershipState.java]
> and send its {{migrationZkVersion}} field to ZooKeeper as the expected
> version for a conditional write to {{{}/migration{}}}.
> Both methods ultimately use {{{}KafkaZkClient{}}}'s own retry loop
> [({{{}retryRequestsUntilConnected{}}})|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/core/src/main/scala/kafka/zk/KafkaZkClient.scala#L2123],
> which repeatedly calls into the lower-level {{ZooKeeperClient}} and resends
> whenever it sees {{{}CONNECTIONLOSS{}}}. If the connection is lost
> ({{{}ConnectionLossException{}}}) after a write has already committed
> server-side, but before its acknowledgment reaches the client, this loop
> transparently resends the identical request, still carrying the same
> {{{}migrationZkVersion{}}}, with no way of knowing the original attempt
> already succeeded. Since the version has genuinely moved on, this resend is
> rejected with {{{}BadVersionException{}}}, even though nothing is actually
> wrong.
> The two methods surface this differently:
> # {{updateMigrationState}} issues a single {{{}SetDataRequest{}}}. On the
> resend's {{{}BadVersionException{}}}, this propagates directly out via
> {{{}resp.maybeThrow(){}}}.
> # {{retryMigrationRequestsUntilConnected}} wraps each request in a multi-op
> transaction that includes a {{{}CheckOp{}}}/{{{}SetDataOp{}}} on
> {{{}/migration{}}}. On the resend's {{{}BadVersionException{}}},
> {{handleUnwrappedMigrationResult}} throws an unconditional
> {{{}RuntimeException{}}}, assuming a second KRaft controller must be writing
> to ZooKeeper.
> Both methods are called from the
> [{{{}KRaftMigrationDriver{}}}|https://github.com/apache/kafka/blob/3.9.2/metadata/src/main/java/org/apache/kafka/metadata/migration/KRaftMigrationDriver.java],
> which caches the {{ZkMigrationLeadershipState}} (and its
> {{{}migrationZkVersion{}}}) in {{{}this.migrationLeadershipState{}}}. Once
> the driver reaches {{{}DUAL_WRITE{}}},
> [{{MetadataChangeEvent.run()}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/metadata/src/main/java/org/apache/kafka/metadata/migration/KRaftMigrationDriver.java#L510]
> calls into them continuously, from three places:
> #
> [{{{}handleDelta{}}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/metadata/src/main/java/org/apache/kafka/metadata/migration/KRaftMigrationDriver.java#L574],
> for each incremental metadata change, writes
> topic/config/ACL/delegation-token updates via
> {{{}retryMigrationRequestsUntilConnected{}}}.
> #
> [{{{}handleSnapshot{}}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/metadata/src/main/java/org/apache/kafka/metadata/migration/KRaftMigrationDriver.java#L570],
> the same mechanism, used when a full metadata snapshot is loaded instead of
> an incremental delta.
> # The dedicated checkpoint write at the end of the same event,
> [{{{}zkMigrationClient.setMigrationRecoveryState(zkStateAfterDualWrite){}}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/metadata/src/main/java/org/apache/kafka/metadata/migration/KRaftMigrationDriver.java#L593],
> via {{{}updateMigrationState{}}}.
> Each of these calls is wrapped in {{{}applyMigrationOperation{}}}, which only
> [reassigns {{this.migrationLeadershipState}} when the call returns
> {*}successfully{*}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/metadata/src/main/java/org/apache/kafka/metadata/migration/KRaftMigrationDriver.java#L257].
> As these methods will now always throw an exception, this leaves the
> driver's cached {{migrationZkVersion}} stale, and is never corrected. Every
> subsequent write, from any of the three call sites, reuses that same stale
> version and fails identically, so the driver becomes stuck indefinitely,
> unable to make further migration progress.
> Both exceptions will eventually surface at
> [{{MigrationEvent.handleException}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/metadata/src/main/java/org/apache/kafka/metadata/migration/KRaftMigrationDriver.java#L416]
> :
> {code:java}
> public void handleException(Throwable e) {
> if (e instanceof MigrationClientAuthException) {
>
> KRaftMigrationDriver.this.faultHandler.handleFault("Encountered ZooKeeper
> authentication in " + this, e);
> } else if (e instanceof MigrationClientException) {
> log.info(String.format("Encountered ZooKeeper error during
> event %s. Will retry.", this), e.getCause());
> } else if (e instanceof RejectedExecutionException) {
> log.debug("Not processing {} because the event queue is
> closed.", this);
> } else {
> KRaftMigrationDriver.this.faultHandler.handleFault("Unhandled
> error in " + this, e);
> }
> }
> {code}
> # The {{MigrationClientException}} which wraps {{updateMigrationState}} ‘s
> {{BadVersionException}} (via
> [{{{}ZkMigrationClient.wrapZkExeception{}}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/core/src/main/scala/kafka/zk/ZkMigrationClient.scala#L68])
> is simply logged at INFO level as “will retry”.
> # The unwrapped {{RuntimeException}} from
> {{retryMigrationRequestsUntilConnected}} instead falls through to
> {{KRaftMigrationDriver.this.faultHandler.handleFault(...)}} . The controller
> which builds that fault handler provides it with [{{fatal =
> false}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/core/src/main/scala/kafka/server/ControllerServer.scala#L304]
> in
> [{{ControllerServer.scala}}|https://github.com/apache/kafka/blob/3.9.2/core/src/main/scala/kafka/server/ControllerServer.scala]
> which when provided to the {{StandardFaultHandlerFactory}} will create a
> {{LoggingFaultHandler}} . This will simply log exceptions as errors, without
> terminating the process.
> {code:java}
> val migrationDriver = KRaftMigrationDriver.newBuilder()
> ...
> .setFaultHandler(sharedServer.faultHandlerFactory.build(
> "zk migration",
> fatal = false,
> () => {}
> ))
> {code}
> Given both paths result in errors which are just logged, nothing actually
> causes the process to exit fatally, meaning the current driver will persist
> in this stalled state indefinitely. It only recovers once an unrelated event
> (e.g. a new controller election, which performs an unconditional overwrite of
> {{{}/migration{}}}) happens to resynchronize it.
> This {{ConnectionLossException}} / {{BadVersionException}} pattern has been
> seen before, with an existing precedent on how to address it.
> [{{KafkaZkClient.conditionalUpdatePath}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/core/src/main/scala/kafka/zk/KafkaZkClient.scala#L965]
> supports an optional checker function that, on a {{BadVersionException}} ,
> re-reads the znode and compares its content against what the caller intended
> to write, so a genuine conflict can be told apart from the caller’s own
> retried write having already landed.
> [{{ReplicationUtils.updateLeaderAndIsr}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/core/src/main/scala/kafka/utils/ReplicationUtils.scala#L33-L34]
> already rely on this via
> [{{checkLeaderAndIsrZkData}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/core/src/main/scala/kafka/utils/ReplicationUtils.scala#L38]
> for leader/ISR updates. However, neither the {{updateMigrationState}} nor
> {{retryMigrationRequestsUntilConnected}} currently use it for
> {{{}/migration{}}}.
> h2. ZkWriteBehindLag Metric
> When the controller is in this stuck state, the {{ZkWriteBehindLag}} metric
> can fail to report the actual lag present. The metric works by comparing the
> controller's latest committed offset against the offset it believes it has
> already mirrored to ZooKeeper, but that mirrored offset is recorded before
> the checkpoint write to {{/migration}} is attempted, not after it succeeds.
> On an event that doesn't otherwise involve any topic, config, ACL, quota,
> producer ID, or delegation token changes, such as one driven by a no-op
> record, this recording happens regardless, immediately before the checkpoint
> write that then fails. The metric can therefore report at or near zero lag
> even while the driver is completely stuck, giving operators no indication
> that dual-write has failed.
> The {{ZkWriteBehindLag}} gauge in
> [{{QuorumControllerMetrics}}|https://github.com/apache/kafka/blob/3.9.2/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java]
> is computed as
> [{{{}lastCommittedRecordOffset() -
> dualWriteOffset(){}}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/metadata/src/main/java/org/apache/kafka/controller/metrics/QuorumControllerMetrics.java#L167].
> {{dualWriteOffset}} is updated via
> [{{{}controllerMetrics.updateDualWriteOffset(image.highestOffsetAndEpoch().offset()){}}}|https://github.com/apache/kafka/blob/5e9866f43ab8e7e41ef39e5584ac50019381328d/metadata/src/main/java/org/apache/kafka/metadata/migration/KRaftMigrationDriver.java#L590],
> called in
> {{MetadataChangeEvent.run()}} immediately before the checkpoint write
> described above.
> Whether that update is actually reached, for any given event, depends on what
> the event contains:
> * If the event includes any topic, config, client quota, SCRAM credential,
> producer ID, ACL, or delegation token change,
> {{KRaftMigrationZkWriter.handleDelta}} calls into the bundled write path
> ({{{}retryMigrationRequestsUntilConnected{}}}) for each one. Once the driver
> is already stuck, these calls fail the same way as described above, and the
> exception propagates out of {{handleDelta}} before {{updateDualWriteOffset}}
> is ever reached. The metric correctly stays where it was, and the reported
> lag keeps growing as expected.
> * If the event contains none of those changes (for example, one driven only
> by a broker registration change or a {{{}NoOpRecord{}}}), {{handleDelta}} has
> nothing to write, returns normally, and execution reaches
> {{updateDualWriteOffset}} regardless. {{dualWriteOffset}} advances to the
> latest offset, and only then does the checkpoint write immediately after it
> fail. On these events, the lag metric is misleadingly reset toward zero, even
> though the driver is still stuck.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)