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]