Copilot commented on code in PR #29007:
URL: https://github.com/apache/flink/pull/29007#discussion_r4206526636


##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/runtime/stream/sql/UdfMetricsITCase.java:
##########
@@ -281,6 +372,93 @@ void testDistinctFunctionsGetSeparateHandles() throws 
Exception {
                 .isEqualTo(SOURCE_ROWS.size());
     }
 
+    @Test
+    void testAsyncScalarMetricsRecorded() throws Exception {
+        StreamTableEnvironment tEnv = createTableEnv(true);
+        tEnv.createTemporarySystemFunction("asyncudf", AsyncIntDoubler.class);
+        createSource(tEnv, "src", "id INT");
+        createBlackHoleSink(tEnv, "sink", "v INT");
+
+        JobID jobId = execute(tEnv, "INSERT INTO sink SELECT asyncudf(id) FROM 
src");
+
+        // The processing time spans dispatch to off-thread completion; one 
sample per input row.
+        assertThat(histogram(jobId, 
processingTimePattern("asyncudf")).getCount())
+                .isEqualTo(SOURCE_ROWS.size());
+        assertThat(counter(jobId, 
exceptionCountPattern("asyncudf")).getCount()).isZero();

Review Comment:
   This only checks that samples exist, so it would still pass if timing 
started at completion and recorded near-zero values. Since the UDF schedules 
completion after 5 ms, also assert that the recorded minimum is at least that 
delay to verify the promised dispatch-to-completion span.



##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/correlate/async/DelegatingAsyncTableResultFuture.java:
##########
@@ -41,21 +44,54 @@ public class DelegatingAsyncTableResultFuture implements 
BiConsumer<Collection<O
 
     private final CompletableFuture<Collection<Object>> completableFuture;
 
+    // Null unless UDF metrics are enabled. The sample decision and start-time 
are taken on the task
+    // thread in the constructor (invoked at dispatch, before eval); the 
histogram/counter are
+    // updated at completion (accept, callback thread). The two updated 
metrics are internally
+    // thread-safe; sample/startNanos are published to the callback via the 
future's completion.
+    @Nullable private final UdfMetrics udfMetrics;
+    private boolean sample;
+    private long startNanos;
+
     public DelegatingAsyncTableResultFuture(
             ResultFuture<Object> delegatedResultFuture,
             boolean needsWrapping,
             boolean isInternalResultType) {
+        this(delegatedResultFuture, needsWrapping, isInternalResultType, null);
+    }
+
+    public DelegatingAsyncTableResultFuture(
+            ResultFuture<Object> delegatedResultFuture,
+            boolean needsWrapping,
+            boolean isInternalResultType,
+            @Nullable UdfMetrics udfMetrics) {
         this.delegatedResultFuture = delegatedResultFuture;
         this.wrapFunction =
                 needsWrapping
                         ? (isInternalResultType ? this::wrapInternal : 
this::wrapExternal)
                         : outs -> outs;
         this.completableFuture = new CompletableFuture<>();
+        this.udfMetrics = udfMetrics;
+        // Sample decision taken on the task thread; the sampler counter is 
never touched
+        // off-thread. These writes must precede whenComplete below: the 
callback registration
+        // performs the volatile completion-stack push that establishes the 
happens-before edge
+        // carrying sample/startNanos to the completing thread's accept().
+        if (udfMetrics != null) {
+            sample = udfMetrics.shouldSample();
+            startNanos = sample ? System.nanoTime() : 0L;
+        }
         this.completableFuture.whenComplete(this);
     }
 
     @Override
     public void accept(Collection<Object> outs, Throwable throwable) {
+        if (udfMetrics != null) {
+            if (throwable != null) {
+                udfMetrics.markException();
+            }

Review Comment:
   The new async-table exception-counting branch has no test: the table IT only 
completes successfully, while the exceptional-completion test exercises the 
separate scalar delegate. Add a unit or integration case that completes an 
async table future exceptionally and asserts this counter increments.



##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/runtime/stream/sql/UdfMetricsITCase.java:
##########
@@ -281,6 +372,93 @@ void testDistinctFunctionsGetSeparateHandles() throws 
Exception {
                 .isEqualTo(SOURCE_ROWS.size());
     }
 
+    @Test
+    void testAsyncScalarMetricsRecorded() throws Exception {
+        StreamTableEnvironment tEnv = createTableEnv(true);
+        tEnv.createTemporarySystemFunction("asyncudf", AsyncIntDoubler.class);
+        createSource(tEnv, "src", "id INT");
+        createBlackHoleSink(tEnv, "sink", "v INT");
+
+        JobID jobId = execute(tEnv, "INSERT INTO sink SELECT asyncudf(id) FROM 
src");
+
+        // The processing time spans dispatch to off-thread completion; one 
sample per input row.
+        assertThat(histogram(jobId, 
processingTimePattern("asyncudf")).getCount())
+                .isEqualTo(SOURCE_ROWS.size());
+        assertThat(counter(jobId, 
exceptionCountPattern("asyncudf")).getCount()).isZero();
+    }
+
+    @Test
+    void testAsyncScalarNullArgumentNotTimed() throws Exception {
+        StreamTableEnvironment tEnv = createTableEnv(true);
+        tEnv.createTemporarySystemFunction("asyncnulludf", 
AsyncIntDoubler.class);
+        String dataId =
+                TestValuesTableFactory.registerData(
+                        Arrays.asList(Row.of((Integer) null), Row.of((Integer) 
null)));
+        tEnv.executeSql(
+                "CREATE TABLE nullsrc (id INT) WITH ("
+                        + "'connector' = 'values', 'bounded' = 'true', 
'data-id' = '"
+                        + dataId
+                        + "')");
+        createBlackHoleSink(tEnv, "sink", "v INT");
+
+        JobID jobId = execute(tEnv, "INSERT INTO sink SELECT asyncnulludf(id) 
FROM nullsrc");
+
+        // A null argument short-circuits to a null result without invoking 
eval, so there is no
+        // call to time. The handle is still registered, but the histogram 
stays empty rather than
+        // collecting near-zero observations that would pull the reported 
latency down.
+        assertThat(histogram(jobId, 
processingTimePattern("asyncnulludf")).getCount()).isZero();
+        assertThat(counter(jobId, 
exceptionCountPattern("asyncnulludf")).getCount()).isZero();
+    }
+
+    @Test
+    void testAsyncCompletionExceptionSurvivesJob() throws Exception {
+        StreamTableEnvironment tEnv = createTableEnv(true);
+        // Serialize invocations so the shared failure counter drives a 
deterministic retry.
+        tEnv.getConfig()
+                
.set(ExecutionConfigOptions.TABLE_EXEC_ASYNC_SCALAR_MAX_CONCURRENT_OPERATIONS, 
1);
+        tEnv.createTemporarySystemFunction("flakyudf", new AsyncFlaky(2));
+        createSource(tEnv, "src", "id INT");
+        createBlackHoleSink(tEnv, "sink", "v INT");
+
+        // Two exceptional completions are counted, then the async retry 
succeeds and the job
+        // finishes normally: an exceptional completion is a soft error, not a 
job failure.
+        JobID jobId = execute(tEnv, "INSERT INTO sink SELECT flakyudf(id) FROM 
src");
+
+        assertThat(counter(jobId, 
exceptionCountPattern("flakyudf")).getCount())
+                .isGreaterThanOrEqualTo(1);

Review Comment:
   `AsyncFlaky(2)` deterministically produces two exceptional completions 
before succeeding because invocations are serialized and the default retry 
budget permits the third attempt. Using `>= 1` lets an implementation that 
misses one retry exception pass, so assert the exact count promised by the test 
comment.



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