[ 
https://issues.apache.org/jira/browse/FLINK-40294?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated FLINK-40294:
-----------------------------------
    Labels: pull-request-available  (was: )

> 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
>            Priority: Major
>              Labels: pull-request-available
>
> 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)

Reply via email to