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]