Weiqing Yang created FLINK-40294:
------------------------------------
Summary: Instrument async scalar and table UDF calls with metrics
Key: FLINK-40294
URL: https://issues.apache.org/jira/browse/FLINK-40294
Project: Flink
Issue Type: Sub-task
Components: Table SQL / Planner, Table SQL / Runtime
Reporter: Weiqing Yang
This is the third of four sub-tasks implementing FLIP-485.
FLIP:
https://cwiki.apache.org/confluence/spaces/FLINK/pages/373885706/FLIP-485+Add+UDF+Metrics
Vote thread: https://lists.apache.org/thread/symqpswsohl2s5wmtkcw0jjp1w5dot0n
Extend the UDF metrics to asynchronous scalar and table user-defined functions,
so
{\{udfProcessingTime}} and \{{udfExceptionCount}} cover the async paths as well
as the
synchronous ones.
Unlike a synchronous call, an async invocation does not begin and end on the
same thread,
so the instrumentation is split:
* At dispatch, on the task thread, the sampling decision is taken and the start
timestamp
is captured.
* At completion, on the callback thread, the elapsed time is recorded and an
exceptional
completion increments the exception counter.
The completion side is handled in the per-invocation delegating futures,
DelegatingAsyncResultFuture for async scalar and
DelegatingAsyncTableResultFuture for
async table. Both the sample decision and the start timestamp are written
before the
completion callback is registered, so there is a happens-before edge to the
callback
thread. The histogram is synchronized and the exception counter is thread-safe.
{\{udfProcessingTime}} for an async function therefore measures the full
dispatch to
completion span, not just the synchronous hand-off.
Main changes:
* AsyncCodeGenerator and BridgingFunctionGenUtil emit the dispatch-side
instrumentation
for the async scalar and async table call sites.
* CodeGeneratorContext gains a single-handle accessor so the generated fetcher
resolves
the shared UdfMetrics handle.
* DelegatingAsyncResultFuture and DelegatingAsyncTableResultFuture record the
completion.
As on the synchronous path, the instrumentation is emitted at code generation
only when
{\{table.exec.udf-metric-enabled}} is true, so the generated code is unchanged
when the
feature is disabled.
Known limitation: an async scalar call whose arguments are all null
short-circuits but
still flows through the future, so with metrics enabled it records one near-zero
histogram entry. This is a bounded skew in the timing distribution, not a
correctness
issue. Avoiding it cleanly needs a separate dispatch path for the short-circuit
case.
Covered by new async cases in UdfMetricsITCase: async scalar and async table
metrics,
naming and scope, and an exceptional completion that increments
udfExceptionCount while
the job still finishes. The existing AsyncCalcITCase and AsyncCorrelateITCase
suites
guard against regressions on the async paths.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)