[ 
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)

Reply via email to