[ 
https://issues.apache.org/jira/browse/FLINK-40343?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18102295#comment-18102295
 ] 

colagy wang commented on FLINK-40343:
-------------------------------------

Even if only *one single row* inside a partition is updated, the current FULL 
refresh logic still triggers recomputation for the entire partition. This 
behavior frequently causes a large number of untouched partitions to be 
refreshed unnecessarily, bringing significant extra compute overhead, 
especially for scenarios with sparse, scattered row‑level modifications across 
many partitions. We hope the new precise batch‑incremental mode can handle this 
case: only process the actually‑changed rows within affected partitions instead 
of re‑computing the whole partition.

> Support row‑level precise batch incremental update for Materialized Table 
> with Left Join support
> ------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40343
>                 URL: https://issues.apache.org/jira/browse/FLINK-40343
>             Project: Flink
>          Issue Type: Improvement
>          Components: Table SQL / API
>    Affects Versions: 2.3.0
>         Environment: Apache Flink 2.3, Apache Paimon
>            Reporter: colagy wang
>            Priority: Major
>
> h3. Background
> Currently Flink Materialized Table provides two refresh modes:
>  # {{{}REFRESH_MODE = CONTINUOUS{}}}: long‑running streaming job. It 
> maintains full operator state for a long time to handle changelog updates. 
> For scenarios where only sparse updates arrive occasionally (e.g. only a few 
> records changed within 180 days), keeping large state for months brings 
> excessive storage and CPU overhead.
>  # {{{}REFRESH_MODE = FULL{}}}: periodic batch refresh. Existing batch 
> incremental capability only works for simple append‑only projection/filter. 
> Once the query contains {{{}JOIN{}}}, {{{}GROUP BY{}}}, batch incremental is 
> silently degraded to full recomputation. Even if only a tiny fraction of 
> rows/partitions are modified, the whole dataset and massive partitions have 
> to be recomputed, resulting in extremely high compute cost.
> h3. Business Pain Point
> Take {{LEFT JOIN(order, product)}} as a typical business case:
>  The business runs on a 15‑minute scheduling cycle. Only sparse partial 
> updates come in each cycle.
>  * Streaming mode: We have to maintain huge join state continuously, just for 
> very few changed records, state cost is unacceptable for long‑term running.
>  * Existing full batch mode: Every refresh triggers recomputation of massive 
> unchanged partitions, compute cost is expensive and wasteful.
> We expect a *row‑level precise batch incremental refresh mode* for 
> Materialized Table, based on Paimon source table:
>  # It does NOT require long‑lived streaming operator state; there is no large 
> state to persist for months.
>  # Only process rows/partitions which have real changes within the scheduling 
> interval (for example every 15 minutes). Avoid recomputing massive untouched 
> partitions.
>  # Support incremental computation for {{{}LEFT JOIN{}}}. This feature can be 
> disabled by default for safety, but should be available when users explicitly 
> enable it.
> h3. Industry Reference
> Aliyun Ververica platform has already implemented similar precise 
> batch‑incremental materialized table capability in production. It proves this 
> technical solution is feasible and valuable for real‑world production 
> workloads.
> h3. Expected Benefit
>  # Reduce compute resource cost: avoid full partition recomputation when only 
> sparse updates occur.
>  # Eliminate heavy long‑term streaming state maintenance cost for 
> low‑change‑frequency business.
>  # Expand batch incremental capability beyond simple filter/projection, 
> support common join workloads such as {{{}LEFT JOIN{}}}.
> h3. Suggested acceptance criteria
>  * New refresh mode for Materialized Table: precise row‑level batch 
> incremental update against Paimon source.
>  * Support incremental {{LEFT JOIN}} calculation, guarded by an explicit 
> config switch (disabled by default).
>  * Only process actually‑modified data set in each scheduled batch run; skip 
> untouched partitions/rows.
>  * No requirement for persistent long‑running streaming operator state.



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

Reply via email to