This is an automated email from the ASF dual-hosted git repository.
lgbo-ustc pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git
The following commit(s) were added to refs/heads/main by this push:
new b939e826b2 [GLUTEN-12468][FLINK] Handle WatermarkStatus elements in
GlutenSourceFunction (#12461)
b939e826b2 is described below
commit b939e826b2d751552ef657197b92fb769986b4ca
Author: lgbo <[email protected]>
AuthorDate: Thu Jul 16 11:16:14 2026 +0800
[GLUTEN-12468][FLINK] Handle WatermarkStatus elements in
GlutenSourceFunction (#12461)
* fix: Handle WatermarkStatus elements in GlutenSourceFunction
Add processing for WatermarkStatus elements from native idle detection.
When IDLE is received, call sourceContext.markAsTemporarilyIdle() to notify
Flink that this source is temporarily idle, allowing watermark progress to
continue from other sources.
* test: Add unit tests for GlutenSourceFunction WatermarkStatus handling
Adds GlutenSourceFunctionWatermarkStatusTest with 5 test cases covering:
- IDLE status triggers markAsTemporarilyIdle()
- ACTIVE status is a no-op on SourceContext
- IDLE→ACTIVE transition does not call markAsTemporarilyIdle()
- Repeated IDLE calls are idempotent
- ACTIVE status produces no invocations
Uses reflection to invoke private processWatermarkStatus() and a custom
TrackingSourceContext spy — no Mockito or native session required.
* test: Add integration test for idle watermark status in WatermarkAssigner
* fix: Update velox4j reference to feature/idle-source-handling branch
* fix: Add shouldCallNoMoreSplits option to GlutenSourceFunction for
unbounded test scenarios
* feat: E2E test for idle WatermarkStatus detection with Kafka + MiniCluster
- Add GlutenStreamSource.isShouldCallNoMoreSplits() delegating to source
- Add GlutenSourceFunction.isShouldCallNoMoreSplits() getter
- Extend OffloadedJobGraphGenerator to preserve shouldCallNoMoreSplits
when creating a new GlutenSourceFunction during offloading
- Rewrite GlutenSourceFunctionWatermarkStatusE2ETest as a real E2E test
using embedded Kafka broker + Flink MiniCluster, verifying that
WatermarkStatus.IDLE is emitted after idle timeout
- Fix EmptyNode output type in WatermarkPushDownSpec project to match
table scan schema (avoids FieldNotFound error during plan init)
* chore: remove obsolete GlutenSourceFunctionWatermarkStatusTest
Covered by GlutenSourceFunctionWatermarkStatusE2ETest which tests
the same behavior end-to-end with Kafka + MiniCluster.
* test: verify idle inputs are excluded from combined min-watermark
Add testIdleInputExcludedFromMinWatermark to
GlutenStreamTwoInputWatermarkStatusTest: when one input is marked
IDLE, its watermark is excluded from min-watermark calculation so
the other active input can advance freely.
* fix: upgrade surefire in gluten-flink-ut from 3.0.0-M5 to 3.3.0
* fix: revert gluten-flink surefire upgrade
* fix: clean up idle source E2E test
* Update velox4j reference for idle source handling
* fix: Remove noMoreSplits test toggle
* Update velox4j reference for watermark destructor fix
* Update velox4j reference for idle timer tests
* Update velox4j reference for idle source handling
* Update velox4j reference for callback bridge fix
* Update velox4j reference to gluten branch
---
.github/workflows/flink.yml | 2 +-
.../gluten/client/OffloadedJobGraphGenerator.java | 16 +-
.../runtime/operators/GlutenSourceFunction.java | 15 +-
.../GlutenOneInputWatermarkAssignerIdleTest.java | 241 +++++++++++++++++++
.../GlutenStreamTwoInputWatermarkStatusTest.java | 26 +++
...GlutenSourceFunctionWatermarkStatusE2ETest.java | 254 +++++++++++++++++++++
6 files changed, 544 insertions(+), 10 deletions(-)
diff --git a/.github/workflows/flink.yml b/.github/workflows/flink.yml
index 2342f4c75e..0fbaa01897 100644
--- a/.github/workflows/flink.yml
+++ b/.github/workflows/flink.yml
@@ -88,7 +88,7 @@ jobs:
export fmt_SOURCE=BUNDLED
export folly_SOURCE=BUNDLED
git clone -b gluten-0530 https://github.com/bigo-sg/velox4j.git
- cd velox4j && git reset --hard
b3987293d08c8b4b13a37162ba7238b45df7feb0
+ cd velox4j && git reset --hard
edffdc6404e942e1eb7b848c6517fa763bb91c7e
git apply $GITHUB_WORKSPACE/gluten-flink/patches/fix-velox4j.patch
$GITHUB_WORKSPACE/build/mvn clean install -DskipTests -Dgpg.skip
-Dspotless.skip=true
cd ..
diff --git
a/gluten-flink/runtime/src/main/java/org/apache/gluten/client/OffloadedJobGraphGenerator.java
b/gluten-flink/runtime/src/main/java/org/apache/gluten/client/OffloadedJobGraphGenerator.java
index a9fa29557d..69e97ecbd8 100644
---
a/gluten-flink/runtime/src/main/java/org/apache/gluten/client/OffloadedJobGraphGenerator.java
+++
b/gluten-flink/runtime/src/main/java/org/apache/gluten/client/OffloadedJobGraphGenerator.java
@@ -187,14 +187,14 @@ public class OffloadedJobGraphGenerator {
boolean supportsVectorOutput =
supportsVectorOutput(sourceChainSlice, chainSliceGraph, jobVertex);
Class<?> outClass = supportsVectorOutput ? StatefulRecord.class :
RowData.class;
- GlutenStreamSource newSourceOp =
- new GlutenStreamSource(
- new GlutenSourceFunction<>(
- planNode,
- sourceOperator.getOutputTypes(),
- sourceOperator.getId(),
- ((GlutenStreamSource) sourceOperator).getConnectorSplit(),
- outClass));
+ GlutenSourceFunction<?> newFn =
+ new GlutenSourceFunction<>(
+ planNode,
+ sourceOperator.getOutputTypes(),
+ sourceOperator.getId(),
+ ((GlutenStreamSource) sourceOperator).getConnectorSplit(),
+ outClass);
+ GlutenStreamSource newSourceOp = new GlutenStreamSource(newFn);
offloadedOpConfig.setStreamOperator(newSourceOp);
if (supportsVectorOutput) {
setOffloadedOutputSerializer(offloadedOpConfig, sourceOperator);
diff --git
a/gluten-flink/runtime/src/main/java/org/apache/gluten/table/runtime/operators/GlutenSourceFunction.java
b/gluten-flink/runtime/src/main/java/org/apache/gluten/table/runtime/operators/GlutenSourceFunction.java
index a47b6c7c17..e7dfab6f09 100644
---
a/gluten-flink/runtime/src/main/java/org/apache/gluten/table/runtime/operators/GlutenSourceFunction.java
+++
b/gluten-flink/runtime/src/main/java/org/apache/gluten/table/runtime/operators/GlutenSourceFunction.java
@@ -31,6 +31,7 @@ import io.github.zhztheplayer.velox4j.session.Session;
import io.github.zhztheplayer.velox4j.stateful.StatefulElement;
import io.github.zhztheplayer.velox4j.stateful.StatefulRecord;
import io.github.zhztheplayer.velox4j.stateful.StatefulWatermark;
+import io.github.zhztheplayer.velox4j.stateful.StatefulWatermarkStatus;
import io.github.zhztheplayer.velox4j.type.RowType;
import org.apache.flink.api.common.state.ListState;
@@ -132,8 +133,10 @@ public class GlutenSourceFunction<OUT> extends
RichParallelSourceFunction<OUT>
processRecord(sourceContext, element.asRecord());
} else if (element.isWatermark()) {
processWatermark(sourceContext, element.asWatermark());
+ } else if (element.isWatermarkStatus()) {
+ processWatermarkStatus(sourceContext, element.asWatermarkStatus());
} else {
- LOG.debug("Ignoring element that is neither record nor watermark");
+ LOG.debug("Ignoring element that is neither record, watermark, nor
watermark status");
}
} finally {
element.close();
@@ -169,6 +172,16 @@ public class GlutenSourceFunction<OUT> extends
RichParallelSourceFunction<OUT>
sourceContext.emitWatermark(new Watermark(watermark.getTimestamp()));
}
+ /** Processes a watermark status and notifies the source context about
idleness. */
+ private void processWatermarkStatus(
+ SourceContext<OUT> sourceContext, StatefulWatermarkStatus status) {
+ if (status.isIdle()) {
+ sourceContext.markAsTemporarilyIdle();
+ }
+ // ACTIVE: no explicit action needed; the source context will resume
+ // activity tracking when the next record or watermark is emitted.
+ }
+
/** Collects a StatefulRecord as RowData by converting the RowVector. */
private void collectAsRowData(SourceContext<OUT> sourceContext,
StatefulRecord record) {
List<RowData> rows =
diff --git
a/gluten-flink/ut/src/test/java/org/apache/gluten/streaming/api/operators/GlutenOneInputWatermarkAssignerIdleTest.java
b/gluten-flink/ut/src/test/java/org/apache/gluten/streaming/api/operators/GlutenOneInputWatermarkAssignerIdleTest.java
new file mode 100644
index 0000000000..00f28d52c8
--- /dev/null
+++
b/gluten-flink/ut/src/test/java/org/apache/gluten/streaming/api/operators/GlutenOneInputWatermarkAssignerIdleTest.java
@@ -0,0 +1,241 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.gluten.streaming.api.operators;
+
+import org.apache.gluten.rexnode.RexConversionContext;
+import org.apache.gluten.rexnode.RexNodeConverter;
+import org.apache.gluten.rexnode.Utils;
+import org.apache.gluten.table.runtime.operators.GlutenOneInputOperator;
+import org.apache.gluten.table.runtime.stream.common.Velox4jEnvironment;
+import org.apache.gluten.util.LogicalTypeConverter;
+import org.apache.gluten.util.PlanNodeIdGenerator;
+
+import io.github.zhztheplayer.velox4j.expression.TypedExpr;
+import io.github.zhztheplayer.velox4j.plan.EmptyNode;
+import io.github.zhztheplayer.velox4j.plan.ProjectNode;
+import io.github.zhztheplayer.velox4j.plan.StatefulPlanNode;
+import io.github.zhztheplayer.velox4j.plan.WatermarkAssignerNode;
+
+import org.apache.flink.api.common.serialization.SerializerConfigImpl;
+import org.apache.flink.api.common.typeinfo.TypeInformation;
+import org.apache.flink.streaming.api.watermark.Watermark;
+import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
+import org.apache.flink.streaming.runtime.watermarkstatus.WatermarkStatus;
+import org.apache.flink.streaming.util.OneInputStreamOperatorTestHarness;
+import org.apache.flink.table.data.GenericRowData;
+import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.planner.calcite.FlinkRexBuilder;
+import org.apache.flink.table.planner.calcite.FlinkTypeFactory;
+import org.apache.flink.table.planner.calcite.FlinkTypeSystem;
+import org.apache.flink.table.runtime.typeutils.InternalTypeInfo;
+import org.apache.flink.table.types.logical.BigIntType;
+import org.apache.flink.table.types.logical.IntType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.RowType;
+
+import org.apache.calcite.rex.RexNode;
+import org.apache.calcite.sql.type.SqlTypeName;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+
+import java.util.List;
+import java.util.Map;
+import java.util.Queue;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * Integration test for WatermarkAssigner idle detection.
+ *
+ * <p>Native {@code checkWatermarkStatus} is driven by the {@code next()} →
{@code advance()} loop
+ * inside the {@code GlutenOneInputOperator}'s drain pipeline. Since {@code
addInput()} resets the
+ * idle baseline on every record, the idle check must happen on an {@code
advance()} call that
+ * processes no new input. We achieve this by calling {@code
processWatermark()} (which triggers a
+ * drain cycle without new data) after waiting past the idle timeout.
+ */
+public class GlutenOneInputWatermarkAssignerIdleTest {
+
+ private static FlinkTypeFactory typeFactory;
+ private static FlinkRexBuilder rexBuilder;
+ private static RowType inputFlinkRowType;
+ private static io.github.zhztheplayer.velox4j.type.RowType inputVeloxType;
+ private static TypeInformation<RowData> typeInfo;
+
+ @BeforeAll
+ static void setUpClass() {
+ Velox4jEnvironment.initializeOnce();
+
+ typeFactory =
+ new FlinkTypeFactory(
+ Thread.currentThread().getContextClassLoader(),
FlinkTypeSystem.INSTANCE);
+ rexBuilder = new FlinkRexBuilder(typeFactory);
+
+ inputFlinkRowType =
+ RowType.of(new LogicalType[] {new IntType(), new BigIntType()}, new
String[] {"id", "ts"});
+ inputVeloxType =
+ (io.github.zhztheplayer.velox4j.type.RowType)
+ LogicalTypeConverter.toVLType(inputFlinkRowType);
+ typeInfo = InternalTypeInfo.of(inputFlinkRowType);
+ }
+
+ @Test
+ void testIdleDetectionWithRealTimePassage() throws Exception {
+ long idleTimeout = 100L; // 100 ms
+ long watermarkInterval = 50L;
+
+ // ── Watermark expression: reference the ts field (index 1) ──
+ List<String> fieldNames = Utils.getNamesFromRowType(inputFlinkRowType);
+ RexNode tsRef =
rexBuilder.makeInputRef(typeFactory.createSqlType(SqlTypeName.BIGINT), 1);
+ TypedExpr watermarkExpr =
+ RexNodeConverter.toTypedExpr(tsRef, new
RexConversionContext(fieldNames));
+
+ ProjectNode watermarkProject =
+ new ProjectNode(
+ PlanNodeIdGenerator.newId(),
+ List.of(new EmptyNode(inputVeloxType)),
+ List.of("TIMESTAMP"),
+ List.of(watermarkExpr));
+
+ // ── WatermarkAssignerNode ──
+ WatermarkAssignerNode assignerNode =
+ new WatermarkAssignerNode(
+ PlanNodeIdGenerator.newId(),
+ null,
+ watermarkProject,
+ idleTimeout,
+ 1, // rowtimeFieldIndex (ts at index 1)
+ watermarkInterval);
+
+ // ── GlutenOneInputOperator ──
+ GlutenOneInputOperator<RowData, RowData> operator =
+ new GlutenOneInputOperator<>(
+ new StatefulPlanNode(assignerNode.getId(), assignerNode),
+ PlanNodeIdGenerator.newId(),
+ inputVeloxType,
+ Map.of(assignerNode.getId(), inputVeloxType),
+ RowData.class,
+ RowData.class,
+ "IdleDetectionTest");
+
+ // ── Test harness ──
+ OneInputStreamOperatorTestHarness<RowData, RowData> harness =
+ new OneInputStreamOperatorTestHarness<>(
+ operator, typeInfo.createSerializer(new SerializerConfigImpl()));
+ harness.setup(typeInfo.createSerializer(new SerializerConfigImpl()));
+ harness.open();
+
+ try {
+ // ── Phase 1: feed one record ──
+ GenericRowData record1 = GenericRowData.of(1, 1000L);
+ harness.processElement(new StreamRecord<>(record1, 1000L));
+ // output: StreamRecord(record1), Watermark(1000)
+ // After drain, checkWatermarkStatus(now1) scheduled timer at now1+100ms.
+
+ // ── Phase 2: wait past idleTimeout, then trigger a drain WITHOUT new
input ──
+ Thread.sleep(idleTimeout * 2); // 200 ms > 100 ms
+ harness.processWatermark(new Watermark(0));
+ // Inside drain: advance → next → advanceWithFuture → blocked
+ // → checkWatermarkStatus(now2) → idle detected (now2 - lastRecordTime
> 100ms)
+ // → push WatermarkStatus.IDLE to pendings_
+
+ // ── Phase 3: feed a second record — the drain first pops pending IDLE
──
+ GenericRowData record2 = GenericRowData.of(2, 2000L);
+ harness.processElement(new StreamRecord<>(record2, 2000L));
+ // drain: advance → next → pendings_ non-empty (IDLE) → pop IDLE → emit
IDLE
+ // → advance → next → process record2 → addInput → onRecord → idle was
true
+ // → emit ACTIVE → push ACTIVE → advance → push record2 + watermark
+ // → pop ACTIVE → emit ACTIVE → pop record2 → collect → pop watermark
→ emit
+
+ // ── Assertions ──
+ Queue<Object> output = harness.getOutput();
+ assertThat(output)
+ .as("Output must contain WatermarkStatus.IDLE after idle timeout")
+ .anyMatch(e -> e instanceof WatermarkStatus && ((WatermarkStatus)
e).isIdle());
+ assertThat(output)
+ .as("Output must contain WatermarkStatus.ACTIVE after idle→active
transition")
+ .anyMatch(e -> e instanceof WatermarkStatus && !((WatermarkStatus)
e).isIdle());
+ assertThat(output)
+ .as("All input records must be preserved")
+ .anyMatch(
+ e ->
+ e instanceof StreamRecord
+ && ((StreamRecord<RowData>) e).getValue().getInt(0) == 1)
+ .anyMatch(
+ e ->
+ e instanceof StreamRecord
+ && ((StreamRecord<RowData>) e).getValue().getInt(0) ==
2);
+ } finally {
+ harness.close();
+ }
+ }
+
+ @Test
+ void testNoIdleWithContinuousRecords() throws Exception {
+ long idleTimeout = 100L;
+ long watermarkInterval = 50L;
+
+ List<String> fieldNames = Utils.getNamesFromRowType(inputFlinkRowType);
+ RexNode tsRef =
rexBuilder.makeInputRef(typeFactory.createSqlType(SqlTypeName.BIGINT), 1);
+ TypedExpr watermarkExpr =
+ RexNodeConverter.toTypedExpr(tsRef, new
RexConversionContext(fieldNames));
+
+ ProjectNode watermarkProject =
+ new ProjectNode(
+ PlanNodeIdGenerator.newId(),
+ List.of(new EmptyNode(inputVeloxType)),
+ List.of("TIMESTAMP"),
+ List.of(watermarkExpr));
+
+ WatermarkAssignerNode assignerNode =
+ new WatermarkAssignerNode(
+ PlanNodeIdGenerator.newId(), null, watermarkProject, idleTimeout,
1, watermarkInterval);
+
+ GlutenOneInputOperator<RowData, RowData> operator =
+ new GlutenOneInputOperator<>(
+ new StatefulPlanNode(assignerNode.getId(), assignerNode),
+ PlanNodeIdGenerator.newId(),
+ inputVeloxType,
+ Map.of(assignerNode.getId(), inputVeloxType),
+ RowData.class,
+ RowData.class,
+ "NoIdleTest");
+
+ OneInputStreamOperatorTestHarness<RowData, RowData> harness =
+ new OneInputStreamOperatorTestHarness<>(
+ operator, typeInfo.createSerializer(new SerializerConfigImpl()));
+ harness.setup(typeInfo.createSerializer(new SerializerConfigImpl()));
+ harness.open();
+
+ try {
+ GenericRowData record1 = GenericRowData.of(1, 1000L);
+ GenericRowData record2 = GenericRowData.of(2, 2000L);
+
+ harness.processElement(new StreamRecord<>(record1, 1000L));
+ // Feed second record immediately (no idle gap) → addInput resets
baseline
+ harness.processElement(new StreamRecord<>(record2, 2000L));
+ // Then drain without new data — idle timeout has NOT elapsed since
record2
+ harness.processWatermark(new Watermark(0));
+
+ Queue<Object> output = harness.getOutput();
+ assertThat(output)
+ .as("No WatermarkStatus.IDLE should appear when records arrive
continuously")
+ .noneMatch(e -> e instanceof WatermarkStatus && ((WatermarkStatus)
e).isIdle());
+ } finally {
+ harness.close();
+ }
+ }
+}
diff --git
a/gluten-flink/ut/src/test/java/org/apache/gluten/streaming/api/operators/GlutenStreamTwoInputWatermarkStatusTest.java
b/gluten-flink/ut/src/test/java/org/apache/gluten/streaming/api/operators/GlutenStreamTwoInputWatermarkStatusTest.java
index 8306c9d20d..a35b21716d 100644
---
a/gluten-flink/ut/src/test/java/org/apache/gluten/streaming/api/operators/GlutenStreamTwoInputWatermarkStatusTest.java
+++
b/gluten-flink/ut/src/test/java/org/apache/gluten/streaming/api/operators/GlutenStreamTwoInputWatermarkStatusTest.java
@@ -79,6 +79,32 @@ public class GlutenStreamTwoInputWatermarkStatusTest extends
GlutenStreamJoinOpe
}
}
+ @Test
+ public void testIdleInputExcludedFromMinWatermark() throws Exception {
+ // When one input is idle, its watermark is excluded from the combined
min-watermark
+ // calculation. The other active input's watermark can advance freely.
+ GlutenTwoInputOperator operator =
createGlutenJoinOperator(FlinkJoinType.INNER);
+
+ try (TwoInputStreamOperatorTestHarness<StatefulRecord, StatefulRecord,
StatefulRecord> harness =
+ new TwoInputStreamOperatorTestHarness<>(operator)) {
+ harness.setup();
+ harness.open();
+
+ harness.processWatermark1(new Watermark(100L));
+ harness.processWatermark2(new Watermark(90L));
+ assertThat(harness.getOutput()).containsExactly(new Watermark(90L));
+
+ harness.processWatermarkStatus1(WatermarkStatus.IDLE);
+ // Input 1 (watermark=100) is idle and excluded. Combined = input 2
(90). No change.
+ assertThat(harness.getOutput()).containsExactly(new Watermark(90L));
+
+ // Input 2 advances to 120. Since input 1 is idle and excluded, combined
= 120.
+ // If input 1 were still active, combined would be min(100, 120) = 100.
+ harness.processWatermark2(new Watermark(120L));
+ assertThat(harness.getOutput()).containsExactly(new Watermark(90L), new
Watermark(120L));
+ }
+ }
+
@Test
public void testWatermarksUseNativeTwoInputMinimum() throws Exception {
// While both inputs are active, native execution should combine indexed
input watermarks by
diff --git
a/gluten-flink/ut/src/test/java/org/apache/gluten/table/runtime/operators/GlutenSourceFunctionWatermarkStatusE2ETest.java
b/gluten-flink/ut/src/test/java/org/apache/gluten/table/runtime/operators/GlutenSourceFunctionWatermarkStatusE2ETest.java
new file mode 100644
index 0000000000..afcc605a2a
--- /dev/null
+++
b/gluten-flink/ut/src/test/java/org/apache/gluten/table/runtime/operators/GlutenSourceFunctionWatermarkStatusE2ETest.java
@@ -0,0 +1,254 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.gluten.table.runtime.operators;
+
+import org.apache.gluten.streaming.api.operators.GlutenStreamSource;
+import org.apache.gluten.table.runtime.stream.common.Velox4jEnvironment;
+
+import io.github.zhztheplayer.velox4j.connector.KafkaConnectorSplit;
+import io.github.zhztheplayer.velox4j.connector.KafkaTableHandle;
+import io.github.zhztheplayer.velox4j.expression.FieldAccessTypedExpr;
+import io.github.zhztheplayer.velox4j.plan.EmptyNode;
+import io.github.zhztheplayer.velox4j.plan.ProjectNode;
+import io.github.zhztheplayer.velox4j.plan.StatefulPlanNode;
+import io.github.zhztheplayer.velox4j.plan.TableScanWithWatermarkNode;
+import io.github.zhztheplayer.velox4j.plan.WatermarkPushDownSpec;
+import io.github.zhztheplayer.velox4j.stateful.StatefulRecord;
+import io.github.zhztheplayer.velox4j.type.BigIntType;
+import io.github.zhztheplayer.velox4j.type.RowType;
+import io.github.zhztheplayer.velox4j.type.VarCharType;
+
+import org.apache.flink.api.common.typeinfo.TypeInformation;
+import org.apache.flink.core.execution.JobClient;
+import org.apache.flink.streaming.api.datastream.DataStreamSource;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.operators.AbstractStreamOperator;
+import org.apache.flink.streaming.api.operators.OneInputStreamOperator;
+import org.apache.flink.streaming.api.watermark.Watermark;
+import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
+import org.apache.flink.streaming.runtime.watermarkstatus.WatermarkStatus;
+
+import com.salesforce.kafka.test.junit5.SharedKafkaTestResource;
+import com.salesforce.kafka.test.listeners.PlainListener;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import java.nio.charset.StandardCharsets;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.UUID;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.TimeUnit;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * End-to-end test that verifies WatermarkStatus.IDLE is emitted through
GlutenSourceFunction when
+ * the native Kafka source detects idleness.
+ *
+ * <p>The test produces a few records to an embedded Kafka broker, starts a
Flink MiniCluster job
+ * with the Gluten native source pipeline, waits for the idle timeout to
expire, and checks that
+ * WatermarkStatus.IDLE is captured by a downstream operator.
+ */
+class GlutenSourceFunctionWatermarkStatusE2ETest {
+
+ private static final int KAFKA_PORT = 19093;
+ private static final long IDLE_TIMEOUT_MS = 5000;
+ private static final long WATERMARK_INTERVAL_MS = 500;
+
+ @RegisterExtension
+ static final SharedKafkaTestResource KAFKA =
+ new SharedKafkaTestResource()
+ .withBrokerProperty("host.name", "127.0.0.1")
+ .withBrokers(1)
+ .registerListener(new PlainListener().onPorts(KAFKA_PORT));
+
+ private static final CopyOnWriteArrayList<WatermarkStatus> capturedStatuses =
+ new CopyOnWriteArrayList<>();
+
+ @BeforeAll
+ static void setupGluten() {
+ Velox4jEnvironment.initializeOnce();
+ }
+
+ @BeforeEach
+ void clearCaptured() {
+ capturedStatuses.clear();
+ }
+
+ @AfterEach
+ void ensureJobCancelled() throws Exception {
+ cancelJob();
+ }
+
+ private JobClient jobClient;
+ private GlutenSourceFunction<StatefulRecord> sourceFunction;
+
+ @Test
+ void testIdleDetectionAfterStopWritingToKafka() throws Exception {
+ String topic = "idle_e2e_" + UUID.randomUUID().toString().replace("-", "");
+ KAFKA.getKafkaTestUtils().createTopic(topic, 1, (short) 1);
+ KAFKA
+ .getKafkaTestUtils()
+ .produceRecords(
+ List.of(
+ jsonRecord(topic, "{\"id\":1000,\"name\":\"r0\"}"),
+ jsonRecord(topic, "{\"id\":2000,\"name\":\"r1\"}"),
+ jsonRecord(topic, "{\"id\":3000,\"name\":\"r2\"}")));
+
+ // Build pipeline: GlutenStreamSource -> WatermarkStatusCaptureOperator
+ StreamExecutionEnvironment env =
StreamExecutionEnvironment.createLocalEnvironment(1);
+ env.getConfig().disableClosureCleaner();
+ env.setParallelism(1);
+
+ DataStreamSource<StatefulRecord> source = addSourceToEnv(env, topic);
+ source
+ .transform(
+ "capture",
+ TypeInformation.of(Object.class),
+ (OneInputStreamOperator) new StatusCaptureOp())
+ .setParallelism(1);
+
+ jobClient = env.executeAsync("IdleDetectionE2ETest");
+ try {
+ waitForIdleStatus();
+ assertThat(capturedStatuses)
+ .as("Should have received WatermarkStatus.IDLE after idle timeout")
+ .contains(WatermarkStatus.IDLE);
+ } finally {
+ cancelJob();
+ }
+ }
+
+ // -- helpers --
+
+ private void waitForIdleStatus() throws InterruptedException {
+ long deadline = System.currentTimeMillis() + IDLE_TIMEOUT_MS + 10000;
+ while (System.currentTimeMillis() < deadline) {
+ if (capturedStatuses.contains(WatermarkStatus.IDLE)) {
+ return;
+ }
+ Thread.sleep(WATERMARK_INTERVAL_MS);
+ }
+ }
+
+ private void cancelJob() throws Exception {
+ try {
+ if (jobClient != null) {
+ JobClient client = jobClient;
+ jobClient = null;
+ try {
+ client.cancel().get(30, TimeUnit.SECONDS);
+ client.getJobExecutionResult().get(30, TimeUnit.SECONDS);
+ } catch (Exception e) {
+ }
+ }
+ } finally {
+ if (sourceFunction != null) {
+ sourceFunction.close();
+ sourceFunction = null;
+ }
+ }
+ }
+
+ private DataStreamSource<StatefulRecord> addSourceToEnv(
+ StreamExecutionEnvironment env, String topic) {
+ RowType veloxRowType =
+ new RowType(List.of("id", "name"), List.of(new BigIntType(), new
VarCharType()));
+
+ ProjectNode watermarkProject =
+ new ProjectNode(
+ "watermark_project",
+ List.of(new EmptyNode(veloxRowType)),
+ List.of("watermark"),
+ List.of(FieldAccessTypedExpr.create(new BigIntType(), "id")));
+ WatermarkPushDownSpec watermarkSpec =
+ new WatermarkPushDownSpec(watermarkProject, IDLE_TIMEOUT_MS,
WATERMARK_INTERVAL_MS, 0);
+
+ Map<String, String> tableParams = new HashMap<>();
+ tableParams.put("bootstrap.servers", "127.0.0.1:" + KAFKA_PORT);
+ tableParams.put("client.id", "test-client-e2e-" + UUID.randomUUID());
+ tableParams.put("group.id", "test-group-e2e");
+ tableParams.put("topic", topic);
+ tableParams.put("format", "json");
+ tableParams.put("scan.startup.mode", "earliest-offsets");
+ tableParams.put("enable.auto.commit", "false");
+
+ KafkaTableHandle tableHandle =
+ new KafkaTableHandle("connector-kafka", topic, veloxRowType,
tableParams);
+ String planId = "plan_" + UUID.randomUUID().toString().replace("-", "");
+ TableScanWithWatermarkNode scanNode =
+ new TableScanWithWatermarkNode(planId, veloxRowType, tableHandle,
List.of(), watermarkSpec);
+ KafkaConnectorSplit connectorSplit =
+ new KafkaConnectorSplit(
+ "connector-kafka",
+ 0,
+ false,
+ "127.0.0.1:" + KAFKA_PORT,
+ "test-group-e2e",
+ "json",
+ false,
+ "earliest-offset",
+ List.of(new KafkaConnectorSplit.TopicPartitionOffset(topic, 0,
-1L)));
+
+ sourceFunction =
+ new GlutenSourceFunction<>(
+ new StatefulPlanNode(scanNode.getId(), scanNode),
+ Map.of(scanNode.getId(), veloxRowType),
+ scanNode.getId(),
+ connectorSplit,
+ StatefulRecord.class);
+
+ GlutenStreamSource sourceOp = new GlutenStreamSource(sourceFunction,
"KafkaSource");
+ return new DataStreamSource<StatefulRecord>(
+ env, TypeInformation.of(StatefulRecord.class), sourceOp, false,
"KafkaSource");
+ }
+
+ private static ProducerRecord<byte[], byte[]> jsonRecord(String topic,
String value) {
+ return new ProducerRecord<>(topic, value.getBytes(StandardCharsets.UTF_8));
+ }
+
+ // -- Capture operator --
+
+ private static class StatusCaptureOp extends AbstractStreamOperator<Object>
+ implements OneInputStreamOperator<Object, Object> {
+
+ @Override
+ public void processElement(StreamRecord<Object> element) throws Exception {
+ Object value = element.getValue();
+ if (value instanceof StatefulRecord) {
+ ((StatefulRecord) value).close();
+ }
+ }
+
+ @Override
+ public void processWatermark(Watermark mark) throws Exception {
+ output.emitWatermark(mark);
+ }
+
+ @Override
+ public void processWatermarkStatus(WatermarkStatus status) throws
Exception {
+ capturedStatuses.add(status);
+ super.processWatermarkStatus(status);
+ }
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]