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


##########
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:
   Fixed. The count assertion alone would have held if the span were measured 
from completion rather than dispatch. `testAsyncScalarMetricsRecorded` now also 
asserts the histogram minimum is at least the scheduled completion delay, 
pulled out as `ASYNC_DELAY_MILLIS` so the UDF and the assertion can't drift 
apart.
   



##########
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:
   Fixed. `testAsyncCompletionExceptionSurvivesJob` uses an 
`AsyncScalarFunction`, so it only ever exercised the scalar delegate. Added 
`testAsyncTableCompletionExceptionCounted`, which runs an `AsyncTableFunction` 
that completes exceptionally twice before succeeding and asserts the async 
table counter reaches 2.
   



##########
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:
   Fixed. I checked the determinism claim before tightening the bound: 
`table.exec.async-scalar.max-attempts` defaults to 3 and the test pins 
max-concurrent-operations to 1, so the first row fails twice and succeeds on 
its third attempt, while the later rows see a counter that has already reached 
the failure budget. Exactly 2 is right, now with the reasoning in a 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