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]

Reply via email to