Hi Mate,
Hope you are doing well,

I wanted to address the comments you had with the FLIP,


Regarding your first point, I completely agree here, while implementing the 
function I also ran into this problem, and currently the implementation has it 
so that when I do not state partition by, the whole stream is sent to one 
operator (global partition).

Regarding your second comment, Gustavo raised the same point, It is a very good 
point I hadn't considered at the time, I will update the flip to match the idea 
that only +I and +U can be kept here, and -U/-D on a key not yet seen should be 
ignored without marking the key as seen, the implementation will be fixed to 
ignore -U and -D on unseen keys.
https://github.com/apache/flink/pull/29365#discussion_r4166714940

Regarding your third point here. Agreed, I will ensure that I explicitly state 
that if state_ttl is unset, the function uses table.exec.state.ttl as other 
PTFs do (0 = retained indefinitely), and an explicit state_ttl overrides it.

Will do these corrections today

Once these are in, I'll leave a few days for further comments before calling a 
vote        


Thanks
Vas
________________________________
From: Mate Czagany <[email protected]>
Sent: 05 October 2026 15:31
To: [email protected] <[email protected]>
Subject: Re: [DISCUSS] FLIP-614: Add built-in DEDUPLICATE_KEEP_FIRST PTF

Hi Vas,

Thank you for this FLIP. I have the following comments:

1. Unkeyed call: two conflicting definitions
Sections 4.1.1, 4.1.2.1 and 4.1.3.1 define the call without PARTITION BY as
whole-row deduplication. However, Section 8 says that exact-duplicate dedup
"is expressed today by listing all payload columns", which only makes sense
if the unkeyed call does not deduplicate on the whole row. In the PR, it
looks like the PTF will only ever emit a single row for the unkeyed
call. Could you please clarify which is intended?

2. Updating input
Section 4.1.3.7 says the function "keeps the first record observed per
key". In the example, the first record is +I, but the first record observed
can as well be -U or -D with a changelog source starting from a later
offset, or after this PTF's own state TTL has expired. I think it should be
stated explicitly here that only +I and +U can be kept here, and -U/-D on a
key not yet seen should be ignored without marking the key as seen.

3. State TTL default
Other PTFs already default to the value of table.exec.state.ttl if the TTL
is not set, thanks to method deriveStateTimeToLive [1]. It's contradictory
to what the FLIP says in Section 4.1.2.3. I think the fallback should be
kept to stay consistent with other PTFs and the FLIP should be reworded to
be specific about this fallback.

Best regards,
Mate Czagany

[1]
https://emea01.safelinks.protection.outlook.com/?url=https%3A%2F%2Fgithub.com%2Fapache%2Fflink%2Fblob%2Fe400c2c221368e98e80d5589cacbcb0b1585ba66%2Fflink-table%2Fflink-table-planner%2Fsrc%2Fmain%2Fjava%2Forg%2Fapache%2Fflink%2Ftable%2Fplanner%2Fplan%2Fnodes%2Fexec%2Fstream%2FStreamExecProcessTableFunction.java%23L404&data=05%7C02%7C%7C010f24b4ea434bbcb36f08df22ed74d0%7C84df9e7fe9f640afb435aaaaaaaaaaaa%7C1%7C0%7C639268075403195443%7CUnknown%7CTWFpbGZsb3d8eyJFbXB0eU1hcGkiOnRydWUsIlYiOiIwLjAuMDAwMCIsIlAiOiJXaW4zMiIsIkFOIjoiTWFpbCIsIldUIjoyfQ%3D%3D%7C0%7C%7C%7C&sdata=9SCoj6Lpam4NZdz19ctqxHbYvjbt3jTymKby3UWqkpw%3D&reserved=0<https://github.com/apache/flink/blob/e400c2c221368e98e80d5589cacbcb0b1585ba66/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecProcessTableFunction.java#L404>



On Fri, Oct 2, 2026 at 2:59 PM Gustavo de Morais <[email protected]>
wrote:

> Hi Vas,
>
> Thanks driving this FLIP! Deduplication is one of the most common uses
> cases that we have in Flink. Currently, there are multiple ways of doing
> deduplication with Flink SQL: some not really intuitive and others,
> inefficient. With this and other incoming PTFs, we'll hopefully provide an
> easier user experience for doing this common task.
>
> Regarding your second point, Roman: in case the user wants to do event-time
> deduplication, we can keep earliest record in state and output it when the
> watermark surpasses it. We don't have to do buffering/sorting in this case
> since it's just one record. That's also what makes it more efficient since
> we don't have to keep multiple records in state.
>
> Friendly bump if others want to review. It'd be nice to start voting on
> this early next week.
>
> Kind regards,
> Gustavo
>
>
> On Thu, 1 Oct 2026 at 12:26, Vas Shabu <[email protected]> wrote:
>
> > Hi all.
> > Just wondering if anyone has any feedback regarding this FLIP, I intend
> to
> > start voting on Monday, October 5th.
> >
> > Thanks
> > Vas
> >
> > ________________________________
> > From: Timo Walther <[email protected]>
> > Sent: 28 September 2026 09:03
> > To: [email protected] <[email protected]>
> > Subject: Re: [DISCUSS] FLIP-614: Add built-in DEDUPLICATE_KEEP_FIRST PTF
> >
> > 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://emea01.safelinks.protection.outlook.com/?url=https%3A%2F%2Fcwiki.apache.org%2Fconfluence%2Fspaces%2FFLINK%2Fpages%2F451975182%2FFLIP-614%2BAdd%2Bbuilt-in%2BDEDUPLICATE_KEEP_FIRST%2BPTF&data=05%7C02%7C%7C010f24b4ea434bbcb36f08df22ed74d0%7C84df9e7fe9f640afb435aaaaaaaaaaaa%7C1%7C0%7C639268075403232395%7CUnknown%7CTWFpbGZsb3d8eyJFbXB0eU1hcGkiOnRydWUsIlYiOiIwLjAuMDAwMCIsIlAiOiJXaW4zMiIsIkFOIjoiTWFpbCIsIldUIjoyfQ%3D%3D%7C0%7C%7C%7C&sdata=Eyx%2BtBXFJTUXnhXMD4B8p%2FHQoZPEHfTkLupI1ky5nV8%3D&reserved=0<https://cwiki.apache.org/confluence/spaces/FLINK/pages/451975182/FLIP-614+Add+built-in+DEDUPLICATE_KEEP_FIRST+PTF>
> > <
> >
> https://emea01.safelinks.protection.outlook.com/?url=https%3A%2F%2Fcwiki.apache.org%2Fconfluence%2Fspaces%2FFLINK%2Fpages%2F451975182%2FFLIP-614%2BAdd%2Bbuilt-in%2BDEDUPLICATE_KEEP_FIRST%2BPTF&data=05%7C02%7C%7C010f24b4ea434bbcb36f08df22ed74d0%7C84df9e7fe9f640afb435aaaaaaaaaaaa%7C1%7C0%7C639268075403260854%7CUnknown%7CTWFpbGZsb3d8eyJFbXB0eU1hcGkiOnRydWUsIlYiOiIwLjAuMDAwMCIsIlAiOiJXaW4zMiIsIkFOIjoiTWFpbCIsIldUIjoyfQ%3D%3D%7C0%7C%7C%7C&sdata=2qb4ccGmgs3UyaDXYhm0woNdlVFUXv5er6QX5lCPjOM%3D&reserved=0<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://emea01.safelinks.protection.outlook.com/?url=https%3A%2F%2Fcwiki.apache.org%2Fconfluence%2Fspaces%2FFLINK%2Fpages%2F451975182%2FFLIP-614%2BAdd%2Bbuilt-in%2BDEDUPLICATE_KEEP_FIRST%2BPTF&data=05%7C02%7C%7C010f24b4ea434bbcb36f08df22ed74d0%7C84df9e7fe9f640afb435aaaaaaaaaaaa%7C1%7C0%7C639268075403287462%7CUnknown%7CTWFpbGZsb3d8eyJFbXB0eU1hcGkiOnRydWUsIlYiOiIwLjAuMDAwMCIsIlAiOiJXaW4zMiIsIkFOIjoiTWFpbCIsIldUIjoyfQ%3D%3D%7C0%7C%7C%7C&sdata=kFlnl2kdNnkrugzNaPfIBRXVV3e6AaWc1N6jedLlu30%3D&reserved=0<https://cwiki.apache.org/confluence/spaces/FLINK/pages/451975182/FLIP-614+Add+built-in+DEDUPLICATE_KEEP_FIRST+PTF>
> > <
> >
> https://emea01.safelinks.protection.outlook.com/?url=https%3A%2F%2Fcwiki.apache.org%2Fconfluence%2Fspaces%2FFLINK%2Fpages%2F451975182%2FFLIP-614%2BAdd%2Bbuilt-in%2BDEDUPLICATE_KEEP_FIRST%2BPTF&data=05%7C02%7C%7C010f24b4ea434bbcb36f08df22ed74d0%7C84df9e7fe9f640afb435aaaaaaaaaaaa%7C1%7C0%7C639268075403313275%7CUnknown%7CTWFpbGZsb3d8eyJFbXB0eU1hcGkiOnRydWUsIlYiOiIwLjAuMDAwMCIsIlAiOiJXaW4zMiIsIkFOIjoiTWFpbCIsIldUIjoyfQ%3D%3D%7C0%7C%7C%7C&sdata=dGJt7z4JjddXlHFp%2FIW370FIUmx%2BwoGcVapP4lHt1iQ%3D&reserved=0<https://cwiki.apache.org/confluence/spaces/FLINK/pages/451975182/FLIP-614+Add+built-in+DEDUPLICATE_KEEP_FIRST+PTF>
> > >
> > >>>
> > >>
> > >>
> > >
> >
> >
>

Reply via email to