This is an automated email from the ASF dual-hosted git repository.
danny0405 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git
The following commit(s) were added to refs/heads/master by this push:
new 07a6635e4990 fix(flink): rethrow StreamWriteOperatorCoordinator
start() failures (#19432)
07a6635e4990 is described below
commit 07a6635e499038467016ec9b830a224769551930
Author: Joy <[email protected]>
AuthorDate: Sun Aug 2 11:33:07 2026 +0800
fix(flink): rethrow StreamWriteOperatorCoordinator start() failures (#19432)
* fix(flink): rethrow coordinator start failures
Avoid using context.failJob when StreamWriteOperatorCoordinator startup
fails, because global failover can keep the half-initialized coordinator
instance alive without invoking start() again.
---------
Co-authored-by: jiangyu84 <[email protected]>
---
.../org/apache/hudi/sink/StreamWriteOperatorCoordinator.java | 10 +++++++---
.../test/java/org/apache/hudi/sink/TestWriteCopyOnWrite.java | 2 +-
2 files changed, 8 insertions(+), 4 deletions(-)
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
index eb3574a77474..e34acac60f98 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
@@ -253,9 +253,13 @@ public class StreamWriteOperatorCoordinator
initClientIds(conf);
}
restoreEvents(Long.MAX_VALUE);
- } catch (Throwable throwable) {
- log.error("Failed to start operator coordinator.", throwable);
- context.failJob(throwable);
+ } catch (Exception exception) {
+ // Rethrow instead of context.failJob(): failJob triggers an in-graph
global failover that
+ // keeps this same coordinator instance alive without calling start()
again, leaving the
+ // half-initialized null fields (executor, writeClient, metaClient ...)
to be reused and NPE
+ // later. Rethrowing surfaces the failure as a JobMaster start failure
so the partially
+ // initialized instance is discarded rather than kept serving.
+ throw new HoodieException("Failed to start operator coordinator.",
exception);
}
}
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestWriteCopyOnWrite.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestWriteCopyOnWrite.java
index 4198ecbb9945..008494fade54 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestWriteCopyOnWrite.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestWriteCopyOnWrite.java
@@ -745,7 +745,7 @@ public class TestWriteCopyOnWrite extends TestWriteBase {
protected void validateNonBlockingConcurrencyControlConditions() {
assertThrows(
- IllegalArgumentException.class,
+ HoodieException.class,
() -> preparePipeline(conf),
"Non-blocking concurrency control requires the MOR table with simple
bucket index");
}