apache / apache/hudi

[BUG] Failure Encountered When Reading Hudi with Flink in Batch Runtime Mode and FlinkOptions.READ_AS_STREAMING=false

Open
#10,576 6 comments 0 reactions 0 assignees View on GitHub
engine:flink priority:high
Dominant language
Java
Stars
6.2k
Forks
2.5k
Avg merge
2d 8h
Merged PRs (30d)
111

Description

I am currently experiencing an issue when attempting to read Hudi with Flink. The problem arises when I configure the Flink RuntimeMode as 'batch' and set the Hudi FlinkOptions.READ_AS_STREAMING to 'false'.

A clear and concise description of the problem.

**To Reproduce**

1. Set Flink RuntimeMode to 'batch'.
2. Set Hudi FlinkOptions.READ_AS_STREAMING to 'false'.
3. Attempt to read Hudi with Flink.

**Expected behavior**

I expected read Hudi table in batch successfully with Flink under these configurations.

** Actual behavior **

A failure occurs when attempting to read Hudi with Flink under these configurations.

**Environment Description**

* Hudi version : From 1.10 ~ 1.14

* Flink version: 1.13

**Additional context**

In the `HoodieTableSource` implementation for Flink's `DynamicTableSource`, a `ScanRuntimeProvider` is provided. This `ScanRuntimeProvider` implements the `produceDataStream` method, which generates a `DataStreamSource`. However, when in Bounded mode, it not explicitly specify the `Boundedness` parameter. By default, Flink uses `Boundedness.CONTINUOUS_UNBOUNDED` as the default parameter, which could potentially be the cause of this issue.

[Code at Hudi HoodieTableSource.java
](https://github.com/apache/hudi/blob/4c7ac6112daab349ebcdd1fbb2216d9d1138ca14/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableSource.java#L227C1-L228C1)

``` java
if (conf.getBoolean(FlinkOptions.READ_AS_STREAMING)) {
...
} else {
...
DataStreamSource source = execEnv.addSource(func, asSummaryString(), typeInfo);
...
}

```
Perhaps the code could be modified as follows:

``` java
if (!isBounded()) {
...
} else {
...
DataStreamSource source = execEnv.addSource(func, asSummaryString(), typeInfo, Boundedness.BOUNDED);
...
}

```

**Stacktrace**

``` java

Caused by: java.lang.IllegalStateException: Detected an UNBOUNDED source with the 'execution.runtime-mode' set to 'BATCH'. This combination is not allowed, please set the 'execution.runtime-mode' to STREAMING or AUTOMATIC

org.apache.flink.client.program.ProgramInvocationException: The main method caused an error: Detected an UNBOUNDED source with the 'execution.runtime-mode' set to 'BATCH'. This combination is not allowed, please set the 'execution.runtime-mode' to STREAMING or AUTOMATIC
at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:381)
at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:223)
at org.apache.flink.client.ClientUtils.executeProgram(ClientUtils.java:114)
at org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:812)
at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:246)
at org.apache.flink.client.cli.CliFrontend.parseAndRun(CliFrontend.java:1054)
at org.apache.flink.client.cli.CliFrontend.lambda$main$10(CliFrontend.java:1132)
at java.security.AccessController.doPrivileged(Native Method)
at javax.security.auth.Subject.doAs(Subject.java:422)
at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1671)
at org.apache.flink.runtime.security.contexts.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)
at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1132)

```

Contributor guide

No contributing guide indexed for this repository

Research direction

Start with hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableSource.java, especially the ScanRuntimeProvider and its produceDataStream implementation. Reproduce with Flink RuntimeMode set to batch and FlinkOptions.READ_AS_STREAMING=false, then verify that reading Hudi succeeds without the unbounded-source failure and uses bounded behavior.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
data-engineering, stream-processing
Issue type
Bug
Difficulty
3/5
Estimated time
1-2 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
42/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.