GinYM opened a new issue, #20018:
URL: https://github.com/apache/hudi/issues/20018

   ### Bug Description
   
   **What happened:**
   A Flink writer fails permanently on startup while initializing the metadata 
table, with:
   ```
   org.apache.hudi.exception.HoodieException: Failed to initialize metadata 
table
   Caused by: java.lang.IllegalArgumentException: For empty cleans, earliest 
instant to retain can never be null
       at 
org.apache.hudi.common.util.ValidationUtils.checkArgument(ValidationUtils.java:42)
       at 
org.apache.hudi.table.action.clean.CleanActionExecutor.createEmptyCleanMetadata(CleanActionExecutor.java:252)
       at 
org.apache.hudi.table.action.clean.CleanActionExecutor.runClean(CleanActionExecutor.java:226)
       at 
org.apache.hudi.table.action.clean.CleanActionExecutor.runPendingClean(CleanActionExecutor.java:204)
       at 
org.apache.hudi.table.action.clean.CleanActionExecutor.execute(CleanActionExecutor.java:289)
       at 
org.apache.hudi.table.HoodieFlinkCopyOnWriteTable.clean(HoodieFlinkCopyOnWriteTable.java:366)
       at 
org.apache.hudi.client.BaseHoodieTableServiceClient.clean(BaseHoodieTableServiceClient.java:909)
       at 
org.apache.hudi.client.BaseHoodieWriteClient.clean(BaseHoodieWriteClient.java:1072)
       at 
org.apache.hudi.metadata.HoodieBackedTableMetadataWriter.executeClean(HoodieBackedTableMetadataWriter.java:2327)
       at 
org.apache.hudi.metadata.HoodieBackedTableMetadataWriter.cleanIfNecessary(HoodieBackedTableMetadataWriter.java:2322)
       at 
org.apache.hudi.metadata.HoodieBackedTableMetadataWriter.performTableServices(HoodieBackedTableMetadataWriter.java:2172)
       at 
org.apache.hudi.client.HoodieFlinkTableServiceClient.initMetadataTable(HoodieFlinkTableServiceClient.java:210)
   ```
   
   This is thrown from `StreamWriteOperatorCoordinator.start()`, so it is fatal 
to the whole job —
   the JobManager cannot start and the job crash-loops on every restart until 
the pending clean
   instant is manually removed/repaired.
   
   Root cause: `CleanPlanActionExecutor.getEmptyCleanerPlan()` (added in 
#18337, "Adding empty
   clean support to hudi") can build and persist a `HoodieCleanerPlan` with no
   `earliestInstantToRetain` set, when `planner.getEarliestCommitToRetain()` 
returns
   `Option.empty()` (e.g. not enough commits yet under the active cleaner 
policy). The
   request-side code guards *scheduling a new* empty-clean instant on
   `cleanerPlan.getEarliestInstantToRetain() != null`, but 
`CleanActionExecutor` — which executes
   an *already-persisted* pending clean plan — has no equivalent guard and 
hard-asserts non-null in
   `createEmptyCleanMetadata()`:
   
       ValidationUtils.checkArgument(cleanerPlan.getEarliestInstantToRetain() 
!= null,
           "For empty cleans, earliest instant to retain can never be null");
   
   Any pending plan that reaches `runPendingClean()` with a null 
`earliestInstantToRetain` crashes
   here instead of completing gracefully. This is the exact NPE risk @yihua 
flagged in review of
   #18337; the fix only closed the gap on the request side, not the execution 
side.
   **What you expected:**
   
   Executing a pending empty-clean plan that has no `earliestInstantToRetain` 
should complete the
   instant with a valid (if boundary-less) `HoodieCleanMetadata`, not throw and 
permanently block
   job/table startup.
   
   **Steps to reproduce:**
   1. Enable the metadata table on a Hudi table written via Flink's 
`StreamWriteOperatorCoordinator`,
      with incremental cleaning enabled and an empty-clean interval configured
      (`hoodie.write.empty.clean.create.duration.ms` / 
`getIntervalToCreateEmptyCleanHours`).
   2. Get a `.clean.requested`/`.clean.inflight` instant persisted on the 
timeline whose
      `HoodieCleanerPlan.earliestInstantToRetain` is unset — reachable via
      `CleanPlanActionExecutor.getEmptyCleanerPlan()` when 
`planner.getEarliestCommitToRetain()`
      returns `Option.empty()`.
   3. Restart the writer so 
`HoodieBackedTableMetadataWriter.initMetadataTable()` →
      `cleanIfNecessary()` → `executeClean()` → 
`CleanActionExecutor.runPendingClean()` picks up and
      executes that pending plan.
   4. `cleanStats` resolves empty → `createEmptyCleanMetadata()` is invoked →
      `ValidationUtils.checkArgument` throws → operator coordinator fails to 
start → job crash-loops.
   
   
   ### Environment
   
   **Hudi version:**
   1.2.0.x line (post-#18337, merged 2026-04-24)
   **Query engine:** (Spark/Flink/Trino etc)
   Flink (1.18), via `StreamWriteOperatorCoordinator` / metadata table writer 
path
   **Relevant configs:**
   - Metadata table: enabled
   - Incremental cleaner mode: enabled
   - `hoodie.write.empty.clean.create.duration.ms` (empty-clean interval): 
configured/non-default
   - Cleaner policy: any policy where 
`CleanPlanner.getEarliestCommitToRetain()` can return
     `Option.empty()` (e.g. too few commits yet relative to 
`hoodie.cleaner.commits.retained`)
   
   ### Logs and Stack Trace
   
   ```
   2026-09-21 23:15:09.513 [flink-pekko.actor.default-dispatcher-15] ERROR 
org.apache.flink.runtime.entrypoint.ClusterEntrypoint  - Fatal error occurred 
in the cluster entrypoint.
   org.apache.flink.util.FlinkRuntimeException: Failed to start the operator 
coordinators
        at 
org.apache.flink.runtime.scheduler.DefaultOperatorCoordinatorHandler.startOperatorCoordinators(DefaultOperatorCoordinatorHandler.java:170)
        at 
org.apache.flink.runtime.scheduler.DefaultOperatorCoordinatorHandler.startAllOperatorCoordinators(DefaultOperatorCoordinatorHandler.java:82)
        at 
org.apache.flink.runtime.scheduler.adaptive.CreatingExecutionGraph.handleExecutionGraphCreation(CreatingExecutionGraph.java:127)
        at 
org.apache.flink.runtime.scheduler.adaptive.CreatingExecutionGraph.lambda$new$0(CreatingExecutionGraph.java:84)
        at 
org.apache.flink.runtime.scheduler.adaptive.AdaptiveScheduler.runIfState(AdaptiveScheduler.java:1253)
        at 
org.apache.flink.runtime.scheduler.adaptive.AdaptiveScheduler.lambda$runIfState$30(AdaptiveScheduler.java:1268)
        at 
java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539)
        at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
        at 
org.apache.flink.runtime.rpc.pekko.PekkoRpcActor.lambda$handleRunAsync$4(PekkoRpcActor.java:451)
        at 
org.apache.flink.runtime.concurrent.ClassLoadingUtils.runWithContextClassLoader(ClassLoadingUtils.java:68)
        at 
org.apache.flink.runtime.rpc.pekko.PekkoRpcActor.handleRunAsync(PekkoRpcActor.java:451)
        at 
org.apache.flink.runtime.rpc.pekko.PekkoRpcActor.handleRpcMessage(PekkoRpcActor.java:218)
        at 
org.apache.flink.runtime.rpc.pekko.FencedPekkoRpcActor.handleRpcMessage(FencedPekkoRpcActor.java:85)
        at 
org.apache.flink.runtime.rpc.pekko.PekkoRpcActor.handleMessage(PekkoRpcActor.java:168)
        at 
org.apache.pekko.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:33)
        at 
org.apache.pekko.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:29)
        at scala.PartialFunction.applyOrElse(PartialFunction.scala:127)
        at scala.PartialFunction.applyOrElse$(PartialFunction.scala:126)
        at 
org.apache.pekko.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:29)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:175)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:176)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:176)
        at org.apache.pekko.actor.Actor.aroundReceive(Actor.scala:547)
        at org.apache.pekko.actor.Actor.aroundReceive$(Actor.scala:545)
        at 
org.apache.pekko.actor.AbstractActor.aroundReceive(AbstractActor.scala:229)
        at org.apache.pekko.actor.ActorCell.receiveMessage(ActorCell.scala:590)
        at org.apache.pekko.actor.ActorCell.invoke(ActorCell.scala:557)
        at org.apache.pekko.dispatch.Mailbox.processMailbox(Mailbox.scala:280)
        at org.apache.pekko.dispatch.Mailbox.run(Mailbox.scala:241)
        at org.apache.pekko.dispatch.Mailbox.exec(Mailbox.scala:253)
        at 
java.base/java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:373)
        at 
java.base/java.util.concurrent.ForkJoinPool$WorkQueue.topLevelExec(ForkJoinPool.java:1182)
        at 
java.base/java.util.concurrent.ForkJoinPool.scan(ForkJoinPool.java:1655)
        at 
java.base/java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1622)
        at 
java.base/java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:165)
   Caused by: org.apache.hudi.exception.HoodieException: Failed to start 
operator coordinator.
        at 
org.apache.hudi.sink.StreamWriteOperatorCoordinator.start(StreamWriteOperatorCoordinator.java:274)
        at 
org.apache.flink.runtime.operators.coordination.OperatorCoordinatorHolder.start(OperatorCoordinatorHolder.java:185)
        at 
org.apache.flink.runtime.scheduler.DefaultOperatorCoordinatorHandler.startOperatorCoordinators(DefaultOperatorCoordinatorHandler.java:165)
        ... 34 common frames omitted
   Caused by: org.apache.hudi.exception.HoodieException: Failed to initialize 
metadata table
        at 
org.apache.hudi.client.HoodieFlinkTableServiceClient.initMetadataTable(HoodieFlinkTableServiceClient.java:215)
        at 
org.apache.hudi.client.HoodieFlinkWriteClient.initMetadataTable(HoodieFlinkWriteClient.java:435)
        at 
org.apache.hudi.sink.StreamWriteOperatorCoordinator.initMetadataTable(StreamWriteOperatorCoordinator.java:521)
        at 
org.apache.hudi.sink.StreamWriteOperatorCoordinator.start(StreamWriteOperatorCoordinator.java:243)
        ... 36 common frames omitted
   Caused by: java.lang.IllegalArgumentException: For empty cleans, earliest 
instant to retain can never be null
        at 
org.apache.hudi.common.util.ValidationUtils.checkArgument(ValidationUtils.java:42)
        at 
org.apache.hudi.table.action.clean.CleanActionExecutor.createEmptyCleanMetadata(CleanActionExecutor.java:254)
        at 
org.apache.hudi.table.action.clean.CleanActionExecutor.runClean(CleanActionExecutor.java:226)
        at 
org.apache.hudi.table.action.clean.CleanActionExecutor.runPendingClean(CleanActionExecutor.java:204)
        at 
org.apache.hudi.table.action.clean.CleanActionExecutor.execute(CleanActionExecutor.java:289)
        at 
org.apache.hudi.table.HoodieFlinkCopyOnWriteTable.clean(HoodieFlinkCopyOnWriteTable.java:366)
        at 
org.apache.hudi.client.BaseHoodieTableServiceClient.clean(BaseHoodieTableServiceClient.java:909)
        at 
org.apache.hudi.client.BaseHoodieWriteClient.clean(BaseHoodieWriteClient.java:1072)
        at 
org.apache.hudi.metadata.HoodieBackedTableMetadataWriter.executeClean(HoodieBackedTableMetadataWriter.java:2327)
        at 
org.apache.hudi.metadata.HoodieBackedTableMetadataWriter.cleanIfNecessary(HoodieBackedTableMetadataWriter.java:2322)
        at 
org.apache.hudi.metadata.HoodieBackedTableMetadataWriter.performTableServices(HoodieBackedTableMetadataWriter.java:2172)
        at 
org.apache.hudi.metadata.HoodieTableMetadataWriter.performTableServices(HoodieTableMetadataWriter.java:178)
        at 
org.apache.hudi.client.HoodieFlinkTableServiceClient.initMetadataTable(HoodieFlinkTableServiceClient.java:210)
        ... 39 common frames omitted
   ```


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to