Copilot commented on code in PR #2415:
URL: https://github.com/apache/auron/pull/2415#discussion_r3627656861


##########
auron-flink-extension/auron-flink-planner/src/test/java/org/apache/auron/flink/table/kafka/AuronKafkaSourceMergeITCase.java:
##########
@@ -140,4 +140,42 @@ public void 
testSharedSourceUnionAllDoesNotFuseUnderDefaultReuse() {
         rows.sort(Comparator.comparingInt(o -> (int) o.getField(0)));
         assertThat(rows).isEqualTo(Arrays.asList(Row.of(21), Row.of(21), 
Row.of(22), Row.of(22), Row.of(23)));
     }
+
+    /**
+     * The same shared-source {@code UNION ALL} as {@link
+     * #testSharedSourceUnionAllDoesNotFuseUnderDefaultReuse}, but with object 
reuse enabled. Under
+     * object reuse Flink hands the <em>same</em> {@code AuronColumnarRowData} 
reference to both
+     * standalone Calc consumers with no defensive copy in between (the 
sibling test runs with reuse
+     * off, where Flink deep-copies the row to heap per consumer). The row set 
is the empirical
+     * guarantee here: the projections and filters convert, so each consumer 
runs as a standalone
+     * native Auron Calc that eagerly copies every column out of the shared 
columnar view into its
+     * own Arrow batch before returning, and neither consumer retains a 
reference to the reused row.
+     * If a native Calc instead aliased the shared row id or read the batch 
after it was recycled,
+     * the row set would collapse or corrupt (all rows folding onto the last 
row of a batch) or fault
+     * on freed off-heap buffers, so the row-set assertion below is what pins 
that native path safe.
+     *
+     * <p>The gate assertion still holds under object reuse: a multi-consumer 
source is never fused,
+     * so both Calcs remain standalone operators and {@code calcOperatorCount} 
is 2.
+     */
+    @Test
+    public void testSharedSourceUnionAllFanOutSafeWithObjectReuse() {
+        environment.setParallelism(1);
+        environment.getConfig().enableObjectReuse();
+        try {
+            String sql =
+                    "SELECT `age` FROM T5 WHERE `age` > 20 " + "UNION ALL 
SELECT `age` + 1 FROM T5 WHERE `age` > 10";
+
+            assertThat(calcOperatorCount(sql))
+                    .as(
+                            "a shared (multi-consumer) source must not fuse 
under object reuse; both Calcs must remain operators")
+                    .isEqualTo(2);
+
+            List<Row> rows = CollectionUtil.iteratorToList(
+                    tableEnvironment.executeSql(sql).collect());
+            rows.sort(Comparator.comparingInt(o -> (int) o.getField(0)));
+            assertThat(rows).isEqualTo(Arrays.asList(Row.of(21), Row.of(21), 
Row.of(22), Row.of(22), Row.of(23)));
+        } finally {
+            environment.getConfig().disableObjectReuse();
+        }

Review Comment:
   This test toggles ExecutionConfig object reuse on a 
StreamExecutionEnvironment that is shared across tests (PER_CLASS lifecycle). 
The finally block unconditionally disables object reuse, which can silently 
change the pre-test setting if it was already enabled by setup or another test. 
Capture the previous setting and restore it to keep the test isolated and 
order-independent.



-- 
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