Hi Roman,

let me try to answer your questions:

1. PTF parameter keep=last

The work on TO_CHANGELOG/FROM_CHANGELOG PTFs has shown that functions with too many parameters and overloaded functionality quickly lead to a sparse feature matrix (i.e. some parameters don't work with other parameters). This leads to a complex user experience.

Instead, we should apply divide-and-conquer here and rather introduce a larger set of functions where the name clearly indicates what it does and each function does exactly one thing right. So DEDUPLICATE_KEEP_LAST should be a separate function in the future. Btw we even discussed to split into DEDUPLICATE_KEEP_FIRST_APPEND and DEDUPLICATE_KEEP_FIRST_UPDATE but focused on splitting on logical level only.

2. ordered until "finalized"

Not sure if I understand this question. Section 4.1.3.7 mentions that "input runs in watermarkless mode only". So there is no changelog stream ordering taking place.

Cheers,
Timo


On 26.09.26 10:53, Roman Khachatryan wrote:
Hi Vas,

Thanks for the proposal.
I have a couple of questions:

1. Would it make sense to change the syntax so that KEEP_LAST semantics
could be added later? (e.g. via PTF parameter keep=last)

2. Could you clarify how changelog stream events (section 4.1.3.7) are
ordered until "finalized" by watermark?

Regards,
Roman


On Fri, Sep 25, 2026 at 5:35 PM Vas Shabu <[email protected]> wrote:

Hi all,

I’d like to propose FLIP-614: Add built-in DEDUPLICATE_KEEP_FIRST PTF [1]
for discussion.

Deduplication is one of the most common transformations in Flink SQL, but
today keep-first deduplication can only be expressed through a ROW_NUMBER()
over-window filtered to the first row. That pattern is verbose and easy to
get wrong: it must be written exactly for the planner to recognise it as
deduplication. This FLIP introduces a built-in Process Table Function (PTF)
that replaces it with a single, self-describing call:

SELECT * FROM DEDUPLICATE_KEEP_FIRST(
   input => TABLE(user_events) PARTITION BY user_id
);

DEDUPLICATE_KEEP_FIRST keeps the first record per key and always produces
an insert-only output. It supports two ordering modes:
- Watermarkless (default): keeps the first record observed for a key, with
no watermark or event-time attribute required.
- Event-time: deterministically keeps the record with the earliest event
time and emits it once the watermark makes the choice final.

The function also accepts updating input in watermarkless mode. It keeps
the first record per key and swallows all later changes, so the result
stays insert-only. Additional state management configuration parameters are
detailed in the FLIP, and we are happy to receive community feedback on two
open design choices: whether reset_ttl_on_duplicate is useful (and its
ideal default), and whether defaulting state_ttl to no TTL, following other
operators, aligns with expectations. Both of these points are under open
design points in the FLIP.

With this proposal, we aim to improve user experience by introducing a
user-friendly feature to perform one of the most common tasks in Flink.
Looking forward to your feedback and thoughts.

Kind regards,
Vas Shabu

[1]
https://cwiki.apache.org/confluence/spaces/FLINK/pages/451975182/FLIP-614+Add+built-in+DEDUPLICATE_KEEP_FIRST+PTF
<http://Hi
%20all,%20%20I’d%20like%20to%20propose%20FLIP-614:%20Add%20built-in%20DEDUPLICATE_KEEP_FIRST%20PTF%20[1]%20for%20discussion.%20%20Deduplication%20is%20one%20of%20the%20most%20common%20transformations%20in%20Flink%20SQL,%20but%20today%20keep-first%20deduplication%20can%20only%20be%20expressed%20through%20a%20ROW_NUMBER()%20over-window%20filtered%20to%20the%20first%20row.%20That%20pattern%20is%20verbose%20and%20easy%20to%20get%20wrong:%20it%20must%20be%20written%20exactly%20for%20the%20planner%20to%20recognise%20it%20as%20deduplication.%20This%20FLIP%20introduces%20a%20built-in%20Process%20Table%20Function%20(PTF)%20that%20replaces%20it%20with%20a%20single,%20self-describing%20call:%20%20SELECT%20*%20FROM%20DEDUPLICATE_KEEP_FIRST(%20
  
input%20=>%20TABLE(user_events)%20PARTITION%20BY%20user_id%20)%20%20DEDUPLICATE_KEEP_FIRST%20keeps%20the%20first%20record%20per%20key%20and%20always%20produces%20an%20insert-only%20output.%20It%20supports%20two%20ordering%20modes:%20-%20Watermarkless%20(default):%20keeps%20the%20first%20record%20observed%20for%20a%20key,%20with%20no%20watermark%20or%20event-time%20attribute%20required.%20-%20Event-time:%20deterministically%20keeps%20the%20record%20with%20the%20earliest%20event%20time%20and%20emits%20it%20once%20the%20watermark%20makes%20the%20choice%20final.%20%20The%20function%20also%20accepts%20updating%20input%20in%20watermarkless%20mode.%20It%20keeps%20the%20first%20record%20per%20key%20and%20swallows%20all%20later%20changes,%20so%20the%20result%20stays%20insert-only.%20Further%20configuration%20parameters%20are%20included%20in%20the%20FLIP%20for%20state%20management.%20%20With%20this%20proposal,%20we%20aim%20to%20improve%20user%20experience%20by%20introducing%20a%20user-friendly%20feature%20to%20perform%20one%20of%20the%20most%20common%20tasks%20in%20Flink.%20Looking%20forward%20to%20your%20feedback%20and%20thoughts.%20%20Kind%20regards,%20Vas%20Shabu%20%20[1]%20
https://cwiki.apache.org/confluence/spaces/FLINK/pages/451975182/FLIP-614+Add+built-in+DEDUPLICATE_KEEP_FIRST+PTF





Reply via email to