Skip to content

[FLINK] Support ORC filesystem sink format#12327

Open
zhanglistar wants to merge 1 commit into
apache:mainfrom
zhanglistar:codex/flink-orc-sink-format
Open

[FLINK] Support ORC filesystem sink format#12327
zhanglistar wants to merge 1 commit into
apache:mainfrom
zhanglistar:codex/flink-orc-sink-format

Conversation

@zhanglistar

@zhanglistar zhanglistar commented Jun 22, 2026

Copy link
Copy Markdown
Contributor

What changes are proposed in this pull request?

Gluten Flink support ORC flilesystem sink format, solves #12203.
Depends on bigo-sg/velox4j#43 and bigo-sg/velox#52.

How was this patch tested?

UT

Was this patch authored or co-authored using generative AI tooling?

Copilot AI review requested due to automatic review settings June 22, 2026 09:07

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

Copilot AI review requested due to automatic review settings June 23, 2026 02:15

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

@zhanglistar
zhanglistar force-pushed the codex/flink-orc-sink-format branch from b9fe77f to 4e42142 Compare June 23, 2026 03:39
Copilot AI review requested due to automatic review settings June 23, 2026 03:55

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

Copilot AI review requested due to automatic review settings June 23, 2026 04:07

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

Comment thread gluten-flink/ut/src/test/resources/nexmark/q10.sql Outdated
Copilot AI review requested due to automatic review settings June 23, 2026 04:29

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

Copilot AI review requested due to automatic review settings June 23, 2026 08:30

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

Copilot AI review requested due to automatic review settings June 24, 2026 05:16

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

@zhanglistar
zhanglistar force-pushed the codex/flink-orc-sink-format branch from 647a12f to 43d0e3f Compare June 29, 2026 09:45
Copilot AI review requested due to automatic review settings June 29, 2026 11:13
Copilot AI review requested due to automatic review settings July 3, 2026 04:16

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 9 out of 9 changed files in this pull request and generated no new comments.

Comments suppressed due to low confidence (1)

gluten-flink/ut/src/test/java/org/apache/gluten/table/runtime/stream/custom/NexmarkTest.java:263

  • This AssertJ call doesn't assert anything because it doesn't chain an assertion (e.g. .isTrue()). As written, the test will pass even if checkJobRunningStatus() returns false, reducing coverage for the Kafka-source path.
    String insertQuery = sqlStatements[sqlStatements.length - 2].trim();

Copilot AI review requested due to automatic review settings July 3, 2026 04:21

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 9 out of 9 changed files in this pull request and generated 1 comment.

Copilot AI review requested due to automatic review settings July 3, 2026 06:11

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 9 out of 9 changed files in this pull request and generated 1 comment.

Comment on lines +226 to +247
// Clean the ORC output directory before running q10_orc to ensure deterministic verification.
if ("q10_orc.sql".equals(queryFileName)) {
Path orcOutputDir = Paths.get("/tmp/data/output/bid_orc");
if (Files.exists(orcOutputDir)) {
try {
try (java.util.stream.Stream<Path> files = Files.walk(orcOutputDir)) {
files
.sorted(java.util.Comparator.reverseOrder())
.forEach(
p -> {
try {
Files.deleteIfExists(p);
} catch (IOException e) {
throw new RuntimeException("Failed to delete " + p, e);
}
});
}
} catch (IOException e) {
throw new RuntimeException("Failed to clean ORC output directory", e);
}
}
}
Copilot AI review requested due to automatic review settings July 6, 2026 08:49
@zhanglistar
zhanglistar force-pushed the codex/flink-orc-sink-format branch from 1c10b05 to b418981 Compare July 6, 2026 08:49

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 14 out of 14 changed files in this pull request and generated 6 comments.

Comment on lines 371 to 376
@Override
public void snapshotState(StateSnapshotContext context) throws Exception {
// TODO: implement it
snapshotNativeState(context.getCheckpointId());
task.snapshotState(0);
super.snapshotState(context);
}
Comment on lines +232 to +245
private void processAvailableElement() {
final StatefulElement statefulElement = task.statefulGet();
try {
if (statefulElement.isWatermark()) {
StatefulWatermark watermark = statefulElement.asWatermark();
output.emitWatermark(new Watermark(watermark.getTimestamp()));
} else {
outputBridge.collect(
output, statefulElement.asRecord(), sessionResource.getAllocator(), outputType);
}
} finally {
statefulElement.close();
}
}
Comment on lines +214 to +227
private void processAvailableElement() {
final StatefulElement element = task.statefulGet();
try {
if (element.isWatermark()) {
StatefulWatermark watermark = element.asWatermark();
output.emitWatermark(new Watermark(watermark.getTimestamp()));
} else {
outputBridge.collect(
output, element.asRecord(), sessionResource.getAllocator(), outputType);
}
} finally {
element.close();
}
}
Comment on lines +254 to +272
private void finishTask() {
while (true) {
UpIterator.State state = task.advance();
switch (state) {
case AVAILABLE:
processAvailableElement();
break;
case BLOCKED:
task.waitFor();
break;
case FINISHED:
return;
default:
// Treat unknown states as terminal, consistent with GlutenSourceFunction.
LOG.warn("Unexpected Velox task state in finishTask: {}", state);
return;
}
}
}
Comment on lines +236 to +254
private void finishTask() {
while (true) {
UpIterator.State state = task.advance();
switch (state) {
case AVAILABLE:
processAvailableElement();
break;
case BLOCKED:
task.waitFor();
break;
case FINISHED:
return;
default:
// Treat unknown states as terminal, consistent with GlutenSourceFunction.
LOG.warn("Unexpected Velox task state in finishTask: {}", state);
return;
}
}
}
Comment on lines +270 to +274
// Allow filesystem sink to flush and commit partitions after job completion.
Thread.sleep(2000);
if ("q10_orc.sql".equals(queryFileName)) {
verifyQ10OrcOutput(queryStartMillis);
}
Copilot AI review requested due to automatic review settings July 10, 2026 08:47
@zhanglistar
zhanglistar force-pushed the codex/flink-orc-sink-format branch from 84ef067 to 7749d8a Compare July 10, 2026 08:49

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 14 out of 14 changed files in this pull request and generated 2 comments.

Comment thread gluten-flink/pom.xml Outdated
Comment on lines +31 to +35
<properties>
<flink.version>1.19.2</flink.version>
<velox4j.version>0.1.0-SNAPSHOT</velox4j.version>
<protobuf.version>3.25.5</protobuf.version>
<hadoop.version>2.7.2</hadoop.version>
Copilot AI review requested due to automatic review settings July 10, 2026 11:44

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 14 out of 14 changed files in this pull request and generated 3 comments.

Comment on lines 65 to 68
} else if (operatorFactory instanceof GlutenOneInputOperatorFactory) {
return Optional.of(((GlutenOneInputOperatorFactory) operatorFactory).getOperator());
} else if (operatorFactory instanceof GlutenTwoInputOperatorFactory) {
return Optional.of(((GlutenTwoInputOperatorFactory) operatorFactory).getOperator());
}
return Optional.empty();
Comment on lines +283 to +286
outputIdle = leftInputIdle && rightInputIdle;
if (wasIdle != outputIdle) {
output.emitWatermarkStatus(outputIdle ? WatermarkStatus.IDLE : WatermarkStatus.ACTIVE);
}
// ProcessingTimeService for Flink AbstractStreamOperator. GlutenTwoInputOperator uses
// GlutenAbstractStreamOperator so it needs the Gluten-specific factory.
offloadedOpConfig.setStreamOperatorFactory(new GlutenTwoInputOperatorFactory<>(newTwoInputOp));
offloadedOpConfig.setStreamOperator(newTwoInputOp);

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 18 out of 18 changed files in this pull request and generated 3 comments.

// ProcessingTimeService for Flink AbstractStreamOperator. GlutenTwoInputOperator uses
// GlutenAbstractStreamOperator so it needs the Gluten-specific factory.
offloadedOpConfig.setStreamOperatorFactory(new GlutenTwoInputOperatorFactory<>(newTwoInputOp));
offloadedOpConfig.setStreamOperator(newTwoInputOp);

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 18 out of 18 changed files in this pull request and generated 3 comments.

Comment on lines 65 to 68
} else if (operatorFactory instanceof GlutenOneInputOperatorFactory) {
return Optional.of(((GlutenOneInputOperatorFactory) operatorFactory).getOperator());
} else if (operatorFactory instanceof GlutenTwoInputOperatorFactory) {
return Optional.of(((GlutenTwoInputOperatorFactory) operatorFactory).getOperator());
}
return Optional.empty();
// ProcessingTimeService for Flink AbstractStreamOperator. GlutenTwoInputOperator uses
// GlutenAbstractStreamOperator so it needs the Gluten-specific factory.
offloadedOpConfig.setStreamOperatorFactory(new GlutenTwoInputOperatorFactory<>(newTwoInputOp));
offloadedOpConfig.setStreamOperator(newTwoInputOp);

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 18 out of 18 changed files in this pull request and generated 2 comments.

Comment on lines 65 to 67
} else if (operatorFactory instanceof GlutenOneInputOperatorFactory) {
return Optional.of(((GlutenOneInputOperatorFactory) operatorFactory).getOperator());
} else if (operatorFactory instanceof GlutenTwoInputOperatorFactory) {
return Optional.of(((GlutenTwoInputOperatorFactory) operatorFactory).getOperator());
}
Comment on lines 280 to 285
sourceOperator.getRightInputType(),
sourceOperator.getOutputTypes(),
inClass,
outClass);
// setStreamOperator would wrap this in Flink's SimpleOperatorFactory, which only initializes
// ProcessingTimeService for Flink AbstractStreamOperator. GlutenTwoInputOperator uses
// GlutenAbstractStreamOperator so it needs the Gluten-specific factory.
offloadedOpConfig.setStreamOperatorFactory(new GlutenTwoInputOperatorFactory<>(newTwoInputOp));
offloadedOpConfig.setStreamOperator(newTwoInputOp);
offloadedOpConfig.setStatePartitioner(0, new GlutenKeySelector());

@KevinyhZou KevinyhZou left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This seems to modify many other operators, which is unrelated to ORC format, we should split them into sperated prs

Comment thread gluten-flink/patches/fix-velox4j.patch Outdated
diff --git a/src/main/cpp/main/velox4j/query/StatefulQueryExecutor.cc b/src/main/cpp/main/velox4j/query/StatefulQueryExecutor.cc
index 2357cf2..159e3d6 100644
--- a/src/main/cpp/main/velox4j/query/StatefulQueryExecutor.cc
+++ b/src/main/cpp/main/velox4j/query/StatefulQueryExecutor.cc

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This patch should be applied to velox4j

leftTransform.getParallelism(),
false);
}
final TwoInputTransformation<RowData, RowData, RowData> transform =

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If operator is not GlutenTwoInputOperator, I think we should fallback to original flink operator

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 18 out of 18 changed files in this pull request and generated 4 comments.

leftTransform,
rightTransform,
createTransformationMeta(JOIN_TRANSFORMATION, config),
new GlutenTwoInputOperatorFactory<>(operator),
Comment on lines 65 to 68
} else if (operatorFactory instanceof GlutenOneInputOperatorFactory) {
return Optional.of(((GlutenOneInputOperatorFactory) operatorFactory).getOperator());
} else if (operatorFactory instanceof GlutenTwoInputOperatorFactory) {
return Optional.of(((GlutenTwoInputOperatorFactory) operatorFactory).getOperator());
}
return Optional.empty();
Comment on lines 451 to +455
public void initializeState(StateInitializationContext context) throws Exception {
// TODO: implement it
initializeNativeState();
super.initializeState(context);
}

private void snapshotNativeState(long checkpointId) {
if (task != null) {
task.snapshotState(checkpointId);
}
}

private void initializeNativeState() {
if (stateInitialized) {
return;
}
initSession();
// TODO: implement it
task.initializeState(0, null);
stateInitialized = true;
super.initializeState(context);

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 17 out of 17 changed files in this pull request and generated no new comments.

Comments suppressed due to low confidence (3)

gluten-flink/runtime/src/main/java/org/apache/gluten/util/Utils.java:68

  • Utils.getGlutenOperator(...) no longer recognizes GlutenTwoInputOperatorFactory. StreamExecJoin/WindowJoin use GlutenTwoInputOperatorFactory, and OffloadedJobGraphGenerator/OperatorChainSliceGraphGenerator call Utils.getGlutenOperator(...).get()/isPresent(); this will return empty for two-input operators and can trigger NoSuchElementException or incorrectly mark slices as unoffloadable.
  public static Optional<GlutenOperator> getGlutenOperator(
      StreamConfig streamConfig, ClassLoader userClassLoader) {
    StreamOperatorFactory operatorFactory = streamConfig.getStreamOperatorFactory(userClassLoader);
    if (operatorFactory instanceof SimpleOperatorFactory) {
      StreamOperator streamOperator = streamConfig.getStreamOperator(userClassLoader);
      if (streamOperator instanceof GlutenOperator) {
        return Optional.of((GlutenOperator) streamOperator);
      }
    } else if (operatorFactory instanceof GlutenOneInputOperatorFactory) {
      return Optional.of(((GlutenOneInputOperatorFactory) operatorFactory).getOperator());
    }
    return Optional.empty();

gluten-flink/runtime/src/main/java/org/apache/gluten/table/runtime/operators/GlutenTwoInputOperator.java:448

  • snapshotState(...) calls task.snapshotState(0), which ignores Flink's checkpoint id and differs from GlutenOneInputOperator/GlutenSourceFunction (they pass context.getCheckpointId()). This likely breaks checkpoint consistency and can cause state corruption or restore failures.
  @Override
  public void snapshotState(StateSnapshotContext context) throws Exception {
    // TODO: implement it
    task.snapshotState(0);
    super.snapshotState(context);
  }

gluten-flink/runtime/src/main/java/org/apache/gluten/client/OffloadedJobGraphGenerator.java:286

  • createOffloadedTwoInputOperator() uses StreamConfig#setStreamOperator(newTwoInputOp), which wraps the operator in Flink's SimpleOperatorFactory. That bypasses GlutenTwoInputOperatorFactory#createStreamOperator(), which currently performs required initialization for GlutenAbstractStreamOperator (setProcessingTimeService(...)) and binds the Gluten mailbox (GlutenMailboxOperatorHelper.bindAtTaskStartup(...)). Without that factory, GlutenAbstractStreamOperator may see a null ProcessingTimeService during setup/RuntimeContext initialization.
    GlutenTwoInputOperator<?, ?> newTwoInputOp =
        new GlutenTwoInputOperator<>(
            planNode,
            sourceOperator.getLeftId(),
            sourceOperator.getRightId(),
            sourceOperator.getLeftInputType(),
            sourceOperator.getRightInputType(),
            sourceOperator.getOutputTypes(),
            inClass,
            outClass);
    offloadedOpConfig.setStreamOperator(newTwoInputOp);
    offloadedOpConfig.setStatePartitioner(0, new GlutenKeySelector());
    offloadedOpConfig.setStatePartitioner(1, new GlutenKeySelector());

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 3 out of 3 changed files in this pull request and generated no new comments.

Comments suppressed due to low confidence (3)

gluten-flink/ut/src/test/java/org/apache/gluten/table/runtime/stream/custom/NexmarkTest.java:259

  • Using a fixed Thread.sleep(2000) to wait for filesystem/partition commits is flaky; on slower CI the ORC files (or success-file) may not be visible within 2s, causing intermittent failures. Poll with a timeout until output is committed (e.g., verify passes) instead of sleeping a fixed duration.
        if ("q10_orc.sql".equals(queryFileName)) {
          Thread.sleep(2000);
          verifyQ10OrcOutput();
        }

gluten-flink/ut/pom.xml:293

  • This module hardcodes protobuf-java 3.25.5, but the repo already defines a pinned protobuf version via the root pom's ${protobuf.version}. Hardcoding here risks dependency convergence/classpath conflicts across modules; prefer using the shared property so the build stays consistent.
      <version>3.25.5</version>

gluten-flink/ut/pom.xml:299

  • This module hardcodes hadoop-client 2.7.2, but the repo already centralizes Hadoop versioning via ${hadoop.version} in the root pom (and profiles may override it). Hardcoding a different version here can introduce conflicting transitive deps in the test classpath; prefer using ${hadoop.version}.
      <version>2.7.2</version>

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 3 out of 3 changed files in this pull request and generated no new comments.

Comments suppressed due to low confidence (3)

gluten-flink/ut/src/test/java/org/apache/gluten/table/runtime/stream/custom/NexmarkTest.java:260

  • This AssertJ statement doesn't assert anything because it doesn't include a terminal assertion (e.g., isTrue()). As written, the test will pass regardless of the condition's value.
      if (kafkaSource) {
        assertThat(checkJobRunningStatus(insertResult, 30000) == true);
      } else {

gluten-flink/ut/pom.xml:295

  • This introduces a hard-coded protobuf version that diverges from the parent POM's protobuf.version property (pom.xml:112). Using the shared property avoids version skew and dependency convergence issues across modules.
      <groupId>com.google.protobuf</groupId>
      <artifactId>protobuf-java</artifactId>
      <version>3.25.5</version>
      <scope>test</scope>
    </dependency>

gluten-flink/ut/pom.xml:301

  • This introduces a hard-coded Hadoop client version (2.7.2) that diverges from the build's managed hadoop.version (pom.xml:85). Aligning to the shared property avoids dependency conflicts and keeps test/classpath behavior consistent across modules.
      <groupId>org.apache.hadoop</groupId>
      <artifactId>hadoop-client</artifactId>
      <version>2.7.2</version>
      <scope>test</scope>
    </dependency>

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 3 out of 3 changed files in this pull request and generated 1 comment.

Comments suppressed due to low confidence (3)

gluten-flink/ut/src/test/java/org/apache/gluten/table/runtime/stream/custom/NexmarkTest.java:261

  • This AssertJ assertion is a no-op because it doesn't end with a terminal assertion (e.g., isTrue()). As written, the test will not fail even if the job never reaches RUNNING.
      if (kafkaSource) {
        assertThat(checkJobRunningStatus(insertResult, 30000) == true);
      } else {

gluten-flink/ut/src/test/java/org/apache/gluten/table/runtime/stream/custom/NexmarkTest.java:275

  • The Q10 ORC test hard-codes an absolute output path under /tmp. This makes the test environment-dependent and can cause collisions between parallel test runs or between local runs and CI. Prefer using a per-test temporary directory (e.g., JUnit temp dir) and wiring that path into the SQL (e.g., via a placeholder + replaceVariables) so the test is hermetic.
  private void cleanQ10OrcOutput() {
    Path outputDir = Paths.get("/tmp/data/output/bid_orc");
    if (!Files.exists(outputDir)) {
      return;
    }

gluten-flink/ut/src/test/resources/nexmark/q10_orc.sql:12

  • The sink path is hard-coded to file:///tmp/... which makes this test SQL non-portable and can conflict across runs. Consider parameterizing the output location (e.g., use a placeholder variable) and let the test harness substitute a per-run temp directory.
  'connector' = 'filesystem',
  'path' = 'file:///tmp/data/output/bid_orc/',
  'format' = 'orc',

Comment thread gluten-flink/ut/pom.xml
Comment on lines +296 to +301
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-client</artifactId>
<version>2.7.2</version>
<scope>test</scope>
</dependency>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants