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

DongShengHe updated FLINK-40888:
--------------------------------
    Description: 
h3. Motivation:

For MATCH_RECOGNIZE queries with DEFINE conditions, the generated 
IterativeCondition eagerly performs the following work for every incoming event:

1. All matched pattern events are extracted. For each pattern variable 
referenced in the condition, all events matched so far for that variable are 
extracted from the NFA shared state into a java.util.List, even when the 
condition only needs the current (last) event. For quantified patterns 
accumulating many events (e.g. A+ over hundreds of rows), this is {{O(N)}} work 
per event, i.e. {{O(N^2)}} overall.

{{2. All aggregations are computed unconditionally. Every aggregate referenced 
in the condition is computed per record before the condition body is evaluated. 
For a condition like DEFINE C AS C.flag > 0 AND SUM(C.amount) > 100, 
SUM(C.amount) (which itself iterates over all events matched for C) is computed 
even when C.flag > 0 is false and the AND is already decided.}}

In production jobs with long matched sequences this leads to significant CPU 
overhead and backpressure. We observed repeated checkpoint failures and job 
restarts for patterns matching more than 500 events on an internal Flink 1.18 
based deployment.
h3. Proposed changes:

1. {*}Current-event shortcut{*}: when a DEFINE condition references the last 
event (offset 0, LAST) of the pattern variable it defines, generate a 
single-element list wrapping the incoming event instead of extracting all 
matched events from the NFA state. This extends the existing fast path that is 
already used for the star pattern variable *.

2. {*}Lazy aggregation via memoized suppliers{*}: the per-record construction 
of pattern event lists and the aggregate computation are wrapped in 
java.util.function.Supplier instances that are created per record but evaluated 
on first access only. The && short-circuit in the condition body therefore 
skips the work entirely when it is not needed. One memoized aggregate-row 
supplier is shared by all aggregates of a pattern variable, so the aggregate 
function still runs at most once per record even when the condition reads 
multiple aggregates of the same variable.

  was:
h3. Motivation:

For MATCH_RECOGNIZE queries with DEFINE conditions, the generated 
IterativeCondition eagerly performs the following work for every incoming event:

1. All matched pattern events are extracted. For each pattern variable 
referenced in the condition, all events matched so far for that variable are 
extracted from the NFA shared state into a java.util.List, even when the 
condition only needs the current (last) event. For quantified patterns 
accumulating many events (e.g. A+ over hundreds of rows), this is {{O(N)}} work 
per event, i.e. {{O(N^2)}} overall.

{{2. All aggregations are computed unconditionally. Every aggregate referenced 
in the condition is computed per record before the condition body is evaluated. 
For a condition like DEFINE C AS C.flag > 0 AND SUM(C.amount) > 100, 
SUM(C.amount) (which itself iterates over all events matched for C) is computed 
even when C.flag > 0 is false and the AND is already decided.}}

In production jobs with long matched sequences this leads to significant CPU 
overhead and backpressure. We observed repeated checkpoint failures and job 
restarts for patterns matching more than 500 events on an internal Flink 1.18 
based deployment.
h3. Proposed changes:

1. {*}Current-event shortcut{*}: when a DEFINE condition references the last 
event (offset 0, LAST) of the pattern variable it defines, generate a 
single-element list wrapping the incoming event instead of extracting all 
matched events from the NFA state. This extends the existing fast path that is 
already used for the star pattern variable (*)

2. {*}Lazy aggregation via memoized suppliers{*}: the per-record construction 
of pattern event lists and the aggregate computation are wrapped in 
java.util.function.Supplier instances that are created per record but evaluated 
on first access only. The && short-circuit in the condition body therefore 
skips the work entirely when it is not needed. One memoized aggregate-row 
supplier is shared by all aggregates of a pattern variable, so the aggregate 
function still runs at most once per record even when the condition reads 
multiple aggregates of the same variable.


> Improve MATCH_RECOGNIZE performance: lazily evaluate DEFINE aggregations and 
> avoid extracting all matched pattern events
> ------------------------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40888
>                 URL: https://issues.apache.org/jira/browse/FLINK-40888
>             Project: Flink
>          Issue Type: Improvement
>          Components: Table SQL / Planner
>            Reporter: DongShengHe
>            Priority: Minor
>
> h3. Motivation:
> For MATCH_RECOGNIZE queries with DEFINE conditions, the generated 
> IterativeCondition eagerly performs the following work for every incoming 
> event:
> 1. All matched pattern events are extracted. For each pattern variable 
> referenced in the condition, all events matched so far for that variable are 
> extracted from the NFA shared state into a java.util.List, even when the 
> condition only needs the current (last) event. For quantified patterns 
> accumulating many events (e.g. A+ over hundreds of rows), this is {{O(N)}} 
> work per event, i.e. {{O(N^2)}} overall.
> {{2. All aggregations are computed unconditionally. Every aggregate 
> referenced in the condition is computed per record before the condition body 
> is evaluated. For a condition like DEFINE C AS C.flag > 0 AND SUM(C.amount) > 
> 100, SUM(C.amount) (which itself iterates over all events matched for C) is 
> computed even when C.flag > 0 is false and the AND is already decided.}}
> In production jobs with long matched sequences this leads to significant CPU 
> overhead and backpressure. We observed repeated checkpoint failures and job 
> restarts for patterns matching more than 500 events on an internal Flink 1.18 
> based deployment.
> h3. Proposed changes:
> 1. {*}Current-event shortcut{*}: when a DEFINE condition references the last 
> event (offset 0, LAST) of the pattern variable it defines, generate a 
> single-element list wrapping the incoming event instead of extracting all 
> matched events from the NFA state. This extends the existing fast path that 
> is already used for the star pattern variable *.
> 2. {*}Lazy aggregation via memoized suppliers{*}: the per-record construction 
> of pattern event lists and the aggregate computation are wrapped in 
> java.util.function.Supplier instances that are created per record but 
> evaluated on first access only. The && short-circuit in the condition body 
> therefore skips the work entirely when it is not needed. One memoized 
> aggregate-row supplier is shared by all aggregates of a pattern variable, so 
> the aggregate function still runs at most once per record even when the 
> condition reads multiple aggregates of the same variable.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to