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

lincoln lee edited comment on FLINK-27849 at 10/5/26 2:22 PM:
--------------------------------------------------------------

Thanks for taking stock of the remaining issues [~martijnvisser]. I would 
welcome you taking the lead on this and re-scoping the ticket into concrete 
follow-ups.

For FLINK-33183, I think there is a direction worth exploring: an explicit, 
user-provided constraint for metadata columns, similar to the idea from 
[FLIP-566|https://cwiki.apache.org/confluence/spaces/FLINK/pages/406621001/FLIP-566+Introduce+a+new+IMMUTABLE+columns+constraint].
 My concern with using {{VIRTUAL}} versus non-{{{}VIRTUAL{}}} remains that it 
describes participation in the sink schema, without guaranteeing that 
retractions preserve the original metadata values.

A possible constraint could be called {*}{{{}RETRACTION CONSISTENT (or 
IMMUTABLE ON RETRACTION{}}}{{{}){}}}{*}, with illustrative syntax such as:
{code:java}
CONSTRAINT metadata_retractions

      COLUMNS (origin_time)

      RETRACTION CONSISTENT NOT ENFORCED

{code}
The contract would be that, for every {{UPDATE_BEFORE}} or {{{}DELETE{}}}, the 
constrained columns carry the same values as in the corresponding row version 
previously emitted as {{INSERT}} or {{{}UPDATE_AFTER{}}}. Values may change 
between versions; retracting an old version must preserve that version’s 
values. This differs from {{{}IMMUTABLE{}}}, which requires values to remain 
unchanged across updates for the same primary key.

The data owner would be responsible for this guarantee. The planner could then 
accept the declared columns when checking source-side NDU requirements, while 
retaining checks for other columns and downstream operations. The constraint 
itself would neither validate nor repair the incoming changelog.

This would give FLINK-33183 a more explicit contract and a focused scope. The 
syntax and semantics would need community discussion, but I think it is a 
useful starting point. The other NDU gaps you identified, and any 
reconsideration of the {{IGNORE}} default, can be evaluated separately without 
waiting for this proposal.


was (Author: lincoln.86xy):
Thanks for taking stock of the remaining issues [~martijnvisser]. I would 
welcome you taking the lead on this and re-scoping the ticket into concrete 
follow-ups.

For FLINK-33183, I think there is a direction worth exploring: an explicit, 
user-provided constraint for metadata columns, similar to the idea from 
[FLIP-566|https://cwiki.apache.org/confluence/spaces/FLINK/pages/406621001/FLIP-566+Introduce+a+new+IMMUTABLE+columns+constraint].
 My concern with using {{VIRTUAL}} versus non-{{{}VIRTUAL{}}} remains that it 
describes participation in the sink schema, without guaranteeing that 
retractions preserve the original metadata values.

A possible constraint could be called {*}{{RETRACTION CONSISTENT (or 
}}{{{}IMMUTABLE ON RETRACTION{}}}{{{}){}}}{*}, with illustrative syntax such as:

{code}

{{CONSTRAINT metadata_retractions}}

  {{    COLUMNS (origin_time)}}

  {{    RETRACTION CONSISTENT NOT ENFORCED}}

{code}{{{}{}}}

The contract would be that, for every {{UPDATE_BEFORE}} or {{{}DELETE{}}}, the 
constrained columns carry the same values as in the corresponding row version 
previously emitted as {{INSERT}} or {{{}UPDATE_AFTER{}}}. Values may change 
between versions; retracting an old version must preserve that version’s 
values. This differs from {{{}IMMUTABLE{}}}, which requires values to remain 
unchanged across updates for the same primary key.

The data owner would be responsible for this guarantee. The planner could then 
accept the declared columns when checking source-side NDU requirements, while 
retaining checks for other columns and downstream operations. The constraint 
itself would neither validate nor repair the incoming changelog.

This would give FLINK-33183 a more explicit contract and a focused scope. The 
syntax and semantics would need community discussion, but I think it is a 
useful starting point. The other NDU gaps you identified, and any 
reconsideration of the {{IGNORE}} default, can be evaluated separately without 
waiting for this proposal.

> Harden correctness for non-deterministic updates present in the changelog 
> pipeline
> ----------------------------------------------------------------------------------
>
>                 Key: FLINK-27849
>                 URL: https://issues.apache.org/jira/browse/FLINK-27849
>             Project: Flink
>          Issue Type: Bug
>          Components: Table SQL / Runtime
>    Affects Versions: 2.1.0
>            Reporter: lincoln lee
>            Assignee: lincoln lee
>            Priority: Major
>             Fix For: 2.4.0
>
>
> There commonly exists updates(which means not only RowKind.INSERT messages) 
> in a streaming pipeline, then wrong results or error may occurs when use some 
> non-deterministic functions or operations.
> It is a long lived issue since the first day that flink sql was available in 
> streaming, but it still not totally be eliminated though some efforts have 
> been taken.
> We should detect all the non-deterministic operations in the changelog 
> pipelines, raise an error to tell users the risk and also add an mechanism 
> that can process such a issue if a user is willing to pay some cost(probably 
> introduce the state).
> All non-deterministic operations include builtin temporal functions(now, 
> current_timestamp...), UUID, RAND... 
> or user defined non-deterministic functions (override isDeterministic return 
> false)
> or a lookup join on a lookup source which data may change over time
> or a cdc-source with meta data field (described in FLINK-28242)
>  
>  
> ====== Solution  ======
> Will introduce a physical plan checker to validate if there's any 
> non-deterministic updates which may cause wrong result, and also a physical 
> plan rewriter to eliminate the non determinism generated by lookup join node 
> (which we think is commonly used in sql, and hard to solve by users 
> themselves).
> For implementation steps, the main changes may include 4 parts:
>  # [preparing work] Adds an internal postOptimize method for physical dag 
> processing
>  # Introduces a `StreamNonDeterministicPlanResolver` to validate if there's 
> any non-deterministic updates which may cause wrong result and rewrite lookup 
> join node with materialization (to eliminate the non determinism generated by 
> lookup join node)
>  # Implements a new lookup join operator (sync mode only) with state to 
> eliminate the non determinism
>  # [optimization] SinkUpsertMaterializer should be aware of the input 
> upsertKey if it is not empty
>  
>  
>  
>  
>  



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

Reply via email to