twalthr commented on code in PR #29365: URL: https://github.com/apache/flink/pull/29365#discussion_r4207880641
########## docs/content.zh/docs/sql/reference/queries/deduplicate-keep-first.md: ########## @@ -0,0 +1,187 @@ +--- +title: "Deduplicate Keep First" +weight: 16 +type: docs +--- +<!-- +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +--> + +# Deduplicate Keep First + +{{< label Streaming >}} + +Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, append-in / append-out alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern. Review Comment: ```suggestion Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern with more configuration options. ``` ########## docs/content.zh/docs/sql/functions/built-in-functions.md: ########## @@ -147,6 +147,7 @@ Table functions can be used in two ways: as stand-alone inputs, where they are i | `FROM_CHANGELOG(input => TABLE t [, ...])` | Converts an append-only table with an explicit operation column into a dynamic table. See Changelog Conversion for the full list of arguments, semantics, and usage. | | `TO_CHANGELOG(input => TABLE t [, ...])` | Converts a dynamic table into an append-only table with an explicit operation column. See Changelog Conversion for the full list of arguments, semantics, and usage. | | `SNAPSHOT(input => TABLE t [, ...])` | Returns the current state of a dynamic table `t`. `SNAPSHOT` can only be used in a `LATERAL` context and not as a stand-alone table function. See [LATERAL SNAPSHOT join]({{< ref "docs/sql/reference/queries/joins" >}}#lateral-snapshot-join) for the full list of arguments, the join semantics, and usage. | +| `DEDUPLICATE_KEEP_FIRST(input => TABLE t [, ...])` | Keeps the first record per key as an insert-only stream (first arrival, or earliest event time with `on_time`). See [Deduplicate Keep First]({{< ref "docs/sql/reference/queries/deduplicate-keep-first" >}}#deduplicate_keep_first) for the full list of arguments, semantics, and usage. | Review Comment: ```suggestion | `DEDUPLICATE_KEEP_FIRST(input => TABLE t [, ...])` | Keeps the first record per key as an append-only table (first arrival, or earliest event time with `on_time`). See [Deduplicate Keep First]({{< ref "docs/sql/reference/queries/deduplicate-keep-first" >}}#deduplicate_keep_first) for the full list of arguments, semantics, and usage. | ``` ########## docs/content.zh/docs/sql/reference/queries/deduplicate-keep-first.md: ########## @@ -0,0 +1,187 @@ +--- +title: "Deduplicate Keep First" +weight: 16 +type: docs +--- +<!-- +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +--> + +# Deduplicate Keep First + +{{< label Streaming >}} + +Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, append-in / append-out alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern. + +| Function | Description | +|:---------|:------------| +| [DEDUPLICATE_KEEP_FIRST](#deduplicate_keep_first) | Keeps the first record per key as an insert-only stream (first arrival, or earliest event time with `on_time`) | + +## DEDUPLICATE_KEEP_FIRST + +The `DEDUPLICATE_KEEP_FIRST` PTF keeps the first record per key and drops every later record for that key. It supports two ordering modes: + +* **Processing time (default):** without `on_time`, the first record *observed* for a key is emitted; later records for that key are dropped while its state exists. This requires no watermark or event-time attribute. +* **Event time:** with `on_time`, the record with the *smallest event time* per key is kept and emitted once the watermark makes the choice final. A record arriving later but carrying an earlier event time replaces the candidate before finalization; records that can no longer become the earliest are dropped as late. + +The input may be insert-only or updating; the output is insert-only in every case. Review Comment: ```suggestion The input may be an append-only or updating table; the output is append-only in every case. ``` ########## docs/content.zh/docs/sql/reference/queries/deduplicate-keep-first.md: ########## @@ -0,0 +1,187 @@ +--- +title: "Deduplicate Keep First" +weight: 16 +type: docs +--- +<!-- +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +--> + +# Deduplicate Keep First + +{{< label Streaming >}} + +Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, append-in / append-out alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern. + +| Function | Description | +|:---------|:------------| +| [DEDUPLICATE_KEEP_FIRST](#deduplicate_keep_first) | Keeps the first record per key as an insert-only stream (first arrival, or earliest event time with `on_time`) | + +## DEDUPLICATE_KEEP_FIRST + +The `DEDUPLICATE_KEEP_FIRST` PTF keeps the first record per key and drops every later record for that key. It supports two ordering modes: + +* **Processing time (default):** without `on_time`, the first record *observed* for a key is emitted; later records for that key are dropped while its state exists. This requires no watermark or event-time attribute. +* **Event time:** with `on_time`, the record with the *smallest event time* per key is kept and emitted once the watermark makes the choice final. A record arriving later but carrying an earlier event time replaces the candidate before finalization; records that can no longer become the earliest are dropped as late. + +The input may be insert-only or updating; the output is insert-only in every case. + +### Syntax + +```sql +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE source_table [PARTITION BY key_col [, key_col2 ...]], + [on_time => DESCRIPTOR(rowtime_column),] + [state_ttl => <interval>,] + [reset_ttl_on_duplicate => <boolean>] +) +``` + +### Parameters + +| Parameter | Required | Description | +|:----------|:---------|:------------| +| `input` | Yes | The input table (insert-only or updating). Use `PARTITION BY` to deduplicate per key; all rows for a key are routed to the same parallel instance. Without `PARTITION BY`, the whole input is a single group processed at parallelism 1, so only the first row of the entire stream is emitted and every later row is dropped regardless of its content. | Review Comment: ```suggestion | `input` | Yes | The input table (append-only or updating). Use `PARTITION BY` to deduplicate per key; all rows for a key are routed to the same parallel instance. Without `PARTITION BY`, the whole input is a single group processed at parallelism 1, so only the first row of the entire stream is emitted and every later row is dropped regardless of its content. | ``` ########## docs/content.zh/docs/sql/reference/queries/deduplicate-keep-first.md: ########## @@ -0,0 +1,187 @@ +--- +title: "Deduplicate Keep First" +weight: 16 +type: docs +--- +<!-- +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +--> + +# Deduplicate Keep First + +{{< label Streaming >}} + +Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, append-in / append-out alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern. + +| Function | Description | +|:---------|:------------| +| [DEDUPLICATE_KEEP_FIRST](#deduplicate_keep_first) | Keeps the first record per key as an insert-only stream (first arrival, or earliest event time with `on_time`) | Review Comment: ```suggestion | [DEDUPLICATE_KEEP_FIRST](#deduplicate_keep_first) | Keeps the first record per key as an append-only table (first arrival, or earliest event time with `on_time`) | ``` ########## docs/content.zh/docs/sql/reference/queries/deduplicate-keep-first.md: ########## @@ -0,0 +1,187 @@ +--- +title: "Deduplicate Keep First" +weight: 16 +type: docs +--- +<!-- +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +--> + +# Deduplicate Keep First + +{{< label Streaming >}} + +Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, append-in / append-out alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern. + +| Function | Description | +|:---------|:------------| +| [DEDUPLICATE_KEEP_FIRST](#deduplicate_keep_first) | Keeps the first record per key as an insert-only stream (first arrival, or earliest event time with `on_time`) | + +## DEDUPLICATE_KEEP_FIRST + +The `DEDUPLICATE_KEEP_FIRST` PTF keeps the first record per key and drops every later record for that key. It supports two ordering modes: + +* **Processing time (default):** without `on_time`, the first record *observed* for a key is emitted; later records for that key are dropped while its state exists. This requires no watermark or event-time attribute. +* **Event time:** with `on_time`, the record with the *smallest event time* per key is kept and emitted once the watermark makes the choice final. A record arriving later but carrying an earlier event time replaces the candidate before finalization; records that can no longer become the earliest are dropped as late. + +The input may be insert-only or updating; the output is insert-only in every case. + +### Syntax + +```sql +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE source_table [PARTITION BY key_col [, key_col2 ...]], + [on_time => DESCRIPTOR(rowtime_column),] + [state_ttl => <interval>,] + [reset_ttl_on_duplicate => <boolean>] +) +``` + +### Parameters + +| Parameter | Required | Description | +|:----------|:---------|:------------| +| `input` | Yes | The input table (insert-only or updating). Use `PARTITION BY` to deduplicate per key; all rows for a key are routed to the same parallel instance. Without `PARTITION BY`, the whole input is a single group processed at parallelism 1, so only the first row of the entire stream is emitted and every later row is dropped regardless of its content. | +| `on_time` | No | A `DESCRIPTOR` naming a single rowtime attribute. When provided, the row with the smallest event time per key is kept and emitted once the watermark passes that timestamp (late rows are dropped). When omitted, the function keeps the first row observed for a key (arrival order). Requires insert-only input; combining `on_time` with an updating input is rejected at planning time. | +| `state_ttl` | No | An `INTERVAL` giving the processing-time retention for the per-key deduplication state. If omitted, falls back to `table.exec.state.ttl`, which retains state indefinitely at its default of 0. After a key's state expires, a later record for that key starts a new deduplication period and may be emitted again. `INTERVAL '0'` disables retention. | Review Comment: > `INTERVAL '0'` disables retention is this true? 0 in PTFs should mean we disable TTL and since the state layout does not even contain a timestamp, one can not enabled it as part of an evolution step. ########## docs/content.zh/docs/sql/reference/queries/deduplicate-keep-first.md: ########## @@ -0,0 +1,187 @@ +--- +title: "Deduplicate Keep First" +weight: 16 +type: docs +--- +<!-- +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +--> + +# Deduplicate Keep First + +{{< label Streaming >}} + +Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, append-in / append-out alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern. + +| Function | Description | +|:---------|:------------| +| [DEDUPLICATE_KEEP_FIRST](#deduplicate_keep_first) | Keeps the first record per key as an insert-only stream (first arrival, or earliest event time with `on_time`) | + +## DEDUPLICATE_KEEP_FIRST + +The `DEDUPLICATE_KEEP_FIRST` PTF keeps the first record per key and drops every later record for that key. It supports two ordering modes: + +* **Processing time (default):** without `on_time`, the first record *observed* for a key is emitted; later records for that key are dropped while its state exists. This requires no watermark or event-time attribute. +* **Event time:** with `on_time`, the record with the *smallest event time* per key is kept and emitted once the watermark makes the choice final. A record arriving later but carrying an earlier event time replaces the candidate before finalization; records that can no longer become the earliest are dropped as late. + +The input may be insert-only or updating; the output is insert-only in every case. + +### Syntax + +```sql +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE source_table [PARTITION BY key_col [, key_col2 ...]], + [on_time => DESCRIPTOR(rowtime_column),] + [state_ttl => <interval>,] + [reset_ttl_on_duplicate => <boolean>] +) +``` + +### Parameters + +| Parameter | Required | Description | +|:----------|:---------|:------------| +| `input` | Yes | The input table (insert-only or updating). Use `PARTITION BY` to deduplicate per key; all rows for a key are routed to the same parallel instance. Without `PARTITION BY`, the whole input is a single group processed at parallelism 1, so only the first row of the entire stream is emitted and every later row is dropped regardless of its content. | +| `on_time` | No | A `DESCRIPTOR` naming a single rowtime attribute. When provided, the row with the smallest event time per key is kept and emitted once the watermark passes that timestamp (late rows are dropped). When omitted, the function keeps the first row observed for a key (arrival order). Requires insert-only input; combining `on_time` with an updating input is rejected at planning time. | +| `state_ttl` | No | An `INTERVAL` giving the processing-time retention for the per-key deduplication state. If omitted, falls back to `table.exec.state.ttl`, which retains state indefinitely at its default of 0. After a key's state expires, a later record for that key starts a new deduplication period and may be emitted again. `INTERVAL '0'` disables retention. | +| `reset_ttl_on_duplicate` | No | Whether a later duplicate refreshes the key's `state_ttl`. Defaults to `TRUE`, so retention tracks the most recent occurrence of a key. Only meaningful when a TTL applies, whether set through `state_ttl` or inherited from `table.exec.state.ttl`. | + +### Output Schema + +When `PARTITION BY` is used, the partition key columns are prepended to the output, followed by the remaining (non-key) input columns. In event-time mode a rowtime column is appended. Review Comment: ```suggestion When `PARTITION BY` is used, the partition key columns are prepended to the output, followed by the remaining (non-key) input columns. In event-time mode a rowtime column is appended for (potentially following) time-based operations. ``` ########## docs/content.zh/docs/sql/reference/queries/deduplicate-keep-first.md: ########## @@ -0,0 +1,187 @@ +--- +title: "Deduplicate Keep First" +weight: 16 +type: docs +--- +<!-- +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +--> + +# Deduplicate Keep First + +{{< label Streaming >}} + +Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, append-in / append-out alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern. + +| Function | Description | +|:---------|:------------| +| [DEDUPLICATE_KEEP_FIRST](#deduplicate_keep_first) | Keeps the first record per key as an insert-only stream (first arrival, or earliest event time with `on_time`) | + +## DEDUPLICATE_KEEP_FIRST + +The `DEDUPLICATE_KEEP_FIRST` PTF keeps the first record per key and drops every later record for that key. It supports two ordering modes: + +* **Processing time (default):** without `on_time`, the first record *observed* for a key is emitted; later records for that key are dropped while its state exists. This requires no watermark or event-time attribute. +* **Event time:** with `on_time`, the record with the *smallest event time* per key is kept and emitted once the watermark makes the choice final. A record arriving later but carrying an earlier event time replaces the candidate before finalization; records that can no longer become the earliest are dropped as late. + +The input may be insert-only or updating; the output is insert-only in every case. + +### Syntax + +```sql +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE source_table [PARTITION BY key_col [, key_col2 ...]], + [on_time => DESCRIPTOR(rowtime_column),] Review Comment: we should also document "uid" argument that is required when more than one dedup is used. btw system args are at the end of the signature, this is important if positional args instead f named args are used. ########## docs/content.zh/docs/sql/reference/queries/deduplicate-keep-first.md: ########## @@ -0,0 +1,187 @@ +--- +title: "Deduplicate Keep First" +weight: 16 +type: docs +--- +<!-- +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +--> + +# Deduplicate Keep First + +{{< label Streaming >}} + +Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, append-in / append-out alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern. + +| Function | Description | +|:---------|:------------| +| [DEDUPLICATE_KEEP_FIRST](#deduplicate_keep_first) | Keeps the first record per key as an insert-only stream (first arrival, or earliest event time with `on_time`) | + +## DEDUPLICATE_KEEP_FIRST + +The `DEDUPLICATE_KEEP_FIRST` PTF keeps the first record per key and drops every later record for that key. It supports two ordering modes: + +* **Processing time (default):** without `on_time`, the first record *observed* for a key is emitted; later records for that key are dropped while its state exists. This requires no watermark or event-time attribute. +* **Event time:** with `on_time`, the record with the *smallest event time* per key is kept and emitted once the watermark makes the choice final. A record arriving later but carrying an earlier event time replaces the candidate before finalization; records that can no longer become the earliest are dropped as late. + +The input may be insert-only or updating; the output is insert-only in every case. + +### Syntax + +```sql +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE source_table [PARTITION BY key_col [, key_col2 ...]], + [on_time => DESCRIPTOR(rowtime_column),] + [state_ttl => <interval>,] + [reset_ttl_on_duplicate => <boolean>] +) +``` + +### Parameters + +| Parameter | Required | Description | +|:----------|:---------|:------------| +| `input` | Yes | The input table (insert-only or updating). Use `PARTITION BY` to deduplicate per key; all rows for a key are routed to the same parallel instance. Without `PARTITION BY`, the whole input is a single group processed at parallelism 1, so only the first row of the entire stream is emitted and every later row is dropped regardless of its content. | +| `on_time` | No | A `DESCRIPTOR` naming a single rowtime attribute. When provided, the row with the smallest event time per key is kept and emitted once the watermark passes that timestamp (late rows are dropped). When omitted, the function keeps the first row observed for a key (arrival order). Requires insert-only input; combining `on_time` with an updating input is rejected at planning time. | +| `state_ttl` | No | An `INTERVAL` giving the processing-time retention for the per-key deduplication state. If omitted, falls back to `table.exec.state.ttl`, which retains state indefinitely at its default of 0. After a key's state expires, a later record for that key starts a new deduplication period and may be emitted again. `INTERVAL '0'` disables retention. | +| `reset_ttl_on_duplicate` | No | Whether a later duplicate refreshes the key's `state_ttl`. Defaults to `TRUE`, so retention tracks the most recent occurrence of a key. Only meaningful when a TTL applies, whether set through `state_ttl` or inherited from `table.exec.state.ttl`. | + +### Output Schema + +When `PARTITION BY` is used, the partition key columns are prepended to the output, followed by the remaining (non-key) input columns. In event-time mode a rowtime column is appended. + +``` +[partition_key_columns] + [remaining_input_columns] +``` + +### Examples + +#### Keyed keep-first (processing time) + +```sql +-- Input (append-only): +-- +I[user_name:'Vas', action:'login'] +-- +I[user_name:'Vas', action:'click'] + +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE user_events PARTITION BY user_name +) + +-- Output (insert-only): +-- +I[user_name:'Vas', action:'login'] +``` + +The first row for `Vas` is emitted; the later `click` on the same key is dropped. + +#### Whole input, no PARTITION BY + +```sql +-- Input (append-only): +-- +I[user_name:'Vas', action:'login'] +-- +I[user_name:'Alice', action:'click'] +-- +I[user_name:'Bob', action:'view'] + +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE user_events Review Comment: we should rather partition by all columns in this case, don't show examples that are difficult to scale. ########## flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/strategies/SpecificInputTypeStrategies.java: ########## @@ -142,6 +142,10 @@ public static InputTypeStrategy plainJsonPath(final InputTypeStrategy signatures public static final InputTypeStrategy FROM_CHANGELOG_INPUT_TYPE_STRATEGY = FromChangelogTypeStrategy.INPUT_TYPE_STRATEGY; + /** Input strategy for {@link BuiltInFunctionDefinitions#DEDUPLICATE_KEEP_FIRST}. */ + public static final InputTypeStrategy DEDUPLICATE_KEEP_FIRST_INPUT_TYPE_STRATEGY = Review Comment: No need to add it here. ########## flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/strategies/DeduplicateKeepFirstTypeStrategy.java: ########## @@ -0,0 +1,191 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.types.inference.strategies; + +import org.apache.flink.annotation.Internal; +import org.apache.flink.table.api.DataTypes; +import org.apache.flink.table.api.DataTypes.Field; +import org.apache.flink.table.api.ValidationException; +import org.apache.flink.table.api.dataview.ValueView; +import org.apache.flink.table.functions.TableSemantics; +import org.apache.flink.table.types.DataType; +import org.apache.flink.table.types.inference.CallContext; +import org.apache.flink.table.types.inference.InputTypeStrategy; +import org.apache.flink.table.types.inference.StateTypeStrategy; +import org.apache.flink.table.types.inference.TypeStrategy; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Optional; +import java.util.Set; +import java.util.stream.Collectors; + +/** Type strategies for the {@code DEDUPLICATE_KEEP_FIRST} process table function. */ +@Internal +public final class DeduplicateKeepFirstTypeStrategy { + + public static final int ARG_INPUT = 0; + public static final int ARG_STATE_TTL = 1; + public static final int ARG_RESET_TTL_ON_DUPLICATE = 2; + private static final String SEEN_STATE_NAME = "seen"; + private static final String CANDIDATE_STATE_NAME = "candidate"; + + public static final InputTypeStrategy INPUT_TYPE_STRATEGY = + new ValidationOnlyInputTypeStrategy() { + @Override + public Optional<List<DataType>> inferInputTypes( + final CallContext callContext, final boolean throwOnFailure) { + final Optional<List<DataType>> stateTtlError = + validateStateTtl(callContext, throwOnFailure); + if (stateTtlError.isPresent()) { + return stateTtlError; + } + + final Optional<List<DataType>> resetTtlError = + validateResetTtlOnDuplicate(callContext, throwOnFailure); + if (resetTtlError.isPresent()) { + return resetTtlError; + } + + return Optional.of(callContext.getArgumentDataTypes()); + } + }; + + public static final TypeStrategy OUTPUT_TYPE_STRATEGY = + callContext -> { + final TableSemantics tableSemantics = + callContext + .getTableSemantics(ARG_INPUT) + .orElseThrow( + () -> + new ValidationException( + "First argument must be a table for DEDUPLICATE_KEEP_FIRST.")); + + final List<Field> inputFields = DataType.getFields(tableSemantics.dataType()); + final List<Field> outputFields = + Arrays.stream( + ChangelogTypeStrategyUtils.computeOutputIndices( + tableSemantics)) + .mapToObj(inputFields::get) + .collect(Collectors.toList()); + return Optional.of(DataTypes.ROW(outputFields).notNull()); + }; + + public static final StateTypeStrategy SEEN_STATE_TYPE_STRATEGY = + new StateTypeStrategy() { + @Override + public Optional<DataType> inferType(final CallContext callContext) { + return Optional.of( + ValueView.newValueViewDataType(DataTypes.BOOLEAN().notNull())); + } + + @Override + public Optional<Duration> getTimeToLive(final CallContext callContext) { + return callContext.getArgumentValue(ARG_STATE_TTL, Duration.class); + } + }; + + public static final StateTypeStrategy CANDIDATE_STATE_TYPE_STRATEGY = + new StateTypeStrategy() { + @Override + public Optional<DataType> inferType(final CallContext callContext) { + final TableSemantics tableSemantics = + callContext + .getTableSemantics(ARG_INPUT) + .orElseThrow( + () -> + new ValidationException( + "First argument must be a table for DEDUPLICATE_KEEP_FIRST.")); + final List<Field> inputFields = DataType.getFields(tableSemantics.dataType()); + final List<Field> candidateFields = + Arrays.stream( + ChangelogTypeStrategyUtils.computeOutputIndices( + tableSemantics)) + .mapToObj(inputFields::get) + .collect(Collectors.toCollection(ArrayList::new)); + // The trailing timestamp field is read positionally at runtime, so its name + // only needs to avoid clashing with a payload column of the same name. + final Set<String> takenNames = + candidateFields.stream() + .map(Field::getName) + .collect(Collectors.toSet()); + String timestampField = "event_time"; + for (int i = 0; takenNames.contains(timestampField); i++) { Review Comment: can we just solve duplicates by x$0, x$1, etc? ########## docs/content.zh/docs/sql/reference/queries/deduplicate-keep-first.md: ########## @@ -0,0 +1,187 @@ +--- +title: "Deduplicate Keep First" +weight: 16 +type: docs +--- +<!-- +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +--> + +# Deduplicate Keep First + +{{< label Streaming >}} + +Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, append-in / append-out alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern. + +| Function | Description | +|:---------|:------------| +| [DEDUPLICATE_KEEP_FIRST](#deduplicate_keep_first) | Keeps the first record per key as an insert-only stream (first arrival, or earliest event time with `on_time`) | + +## DEDUPLICATE_KEEP_FIRST + +The `DEDUPLICATE_KEEP_FIRST` PTF keeps the first record per key and drops every later record for that key. It supports two ordering modes: + +* **Processing time (default):** without `on_time`, the first record *observed* for a key is emitted; later records for that key are dropped while its state exists. This requires no watermark or event-time attribute. +* **Event time:** with `on_time`, the record with the *smallest event time* per key is kept and emitted once the watermark makes the choice final. A record arriving later but carrying an earlier event time replaces the candidate before finalization; records that can no longer become the earliest are dropped as late. + +The input may be insert-only or updating; the output is insert-only in every case. + +### Syntax + +```sql +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE source_table [PARTITION BY key_col [, key_col2 ...]], + [on_time => DESCRIPTOR(rowtime_column),] + [state_ttl => <interval>,] + [reset_ttl_on_duplicate => <boolean>] +) +``` + +### Parameters + +| Parameter | Required | Description | +|:----------|:---------|:------------| +| `input` | Yes | The input table (insert-only or updating). Use `PARTITION BY` to deduplicate per key; all rows for a key are routed to the same parallel instance. Without `PARTITION BY`, the whole input is a single group processed at parallelism 1, so only the first row of the entire stream is emitted and every later row is dropped regardless of its content. | +| `on_time` | No | A `DESCRIPTOR` naming a single rowtime attribute. When provided, the row with the smallest event time per key is kept and emitted once the watermark passes that timestamp (late rows are dropped). When omitted, the function keeps the first row observed for a key (arrival order). Requires insert-only input; combining `on_time` with an updating input is rejected at planning time. | +| `state_ttl` | No | An `INTERVAL` giving the processing-time retention for the per-key deduplication state. If omitted, falls back to `table.exec.state.ttl`, which retains state indefinitely at its default of 0. After a key's state expires, a later record for that key starts a new deduplication period and may be emitted again. `INTERVAL '0'` disables retention. | +| `reset_ttl_on_duplicate` | No | Whether a later duplicate refreshes the key's `state_ttl`. Defaults to `TRUE`, so retention tracks the most recent occurrence of a key. Only meaningful when a TTL applies, whether set through `state_ttl` or inherited from `table.exec.state.ttl`. | + +### Output Schema + +When `PARTITION BY` is used, the partition key columns are prepended to the output, followed by the remaining (non-key) input columns. In event-time mode a rowtime column is appended. + +``` +[partition_key_columns] + [remaining_input_columns] +``` + +### Examples + +#### Keyed keep-first (processing time) + +```sql +-- Input (append-only): +-- +I[user_name:'Vas', action:'login'] +-- +I[user_name:'Vas', action:'click'] + +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE user_events PARTITION BY user_name +) + +-- Output (insert-only): +-- +I[user_name:'Vas', action:'login'] +``` + +The first row for `Vas` is emitted; the later `click` on the same key is dropped. + +#### Whole input, no PARTITION BY + +```sql +-- Input (append-only): +-- +I[user_name:'Vas', action:'login'] +-- +I[user_name:'Alice', action:'click'] +-- +I[user_name:'Bob', action:'view'] + +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE user_events +) + +-- Output (insert-only): +-- +I[user_name:'Vas', action:'login'] +``` + +Without `PARTITION BY` the whole input is a single group at parallelism 1, so only the first row of the entire stream is emitted; every later row is dropped regardless of its content. + +#### Bounded state with state_ttl + +```sql +-- Input (append-only): +-- +I[user_name:'Vas', action:'login'] +-- +I[user_name:'Vas', action:'click'] +-- +I[user_name:'Vas', action:'view'] + +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE user_events PARTITION BY user_name, + state_ttl => INTERVAL '5' SECOND, + reset_ttl_on_duplicate => FALSE +) + +-- Output (insert-only): +-- +I[user_name:'Vas', action:'login'] +``` + +`state_ttl` bounds how long per-key state is retained (processing time). While the state exists, duplicates are dropped and only the first row per key is emitted; after the state expires, a later record for the key starts a new deduplication period. With the default `reset_ttl_on_duplicate => TRUE`, each dropped duplicate refreshes the retention window; with `FALSE`, retention is measured from the first record for the key. + +#### Earliest by event time + +```sql +-- Source declares: WATERMARK FOR ts AS ts - INTERVAL '10' SECOND +-- Input (ts = event time): +-- +I[user_name:'Vas', action:'login', ts:3.000] +-- +I[user_name:'Vas', action:'click', ts:1.000] +-- +I[user_name:'Vas', action:'view', ts:5.000] + +SELECT user_name, action FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE user_events PARTITION BY user_name, + on_time => DESCRIPTOR(ts) +) + +-- Output (insert-only, once the watermark finalizes the earliest): +-- +I[user_name:'Vas', action:'click'] +``` + +With `on_time`, the record with the smallest event time per key is kept (here `click` at `ts 1.000`, not `login` which arrived first) and emitted once the watermark passes that timestamp. + +#### Updating input + +```sql +-- Input (updating changelog): +-- -D[user_name:'Vas', action:'logout'] key not seen yet -> ignored +-- +I[user_name:'Vas', action:'login'] first record for the key -> emitted +-- +I[user_name:'Vas', action:'login'] duplicate -> swallowed +-- -U[user_name:'Vas', action:'login'] retraction -> swallowed +-- +U[user_name:'Vas', action:'click'] update -> swallowed +-- -D[user_name:'Vas', action:'click'] delete -> swallowed + +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE user_events PARTITION BY user_name +) + +-- Output (insert-only): +-- +I[user_name:'Vas', action:'login'] +``` + +The input may be an updating changelog. `DEDUPLICATE_KEEP_FIRST` keeps the first `+I` or `+U` observed per key and swallows every later change to that key (duplicate, `-U`, `+U`, `-D`), so the output stays insert-only. A `-U` or `-D` for a key that has not been seen yet is ignored and does not mark the key as seen, so the next `+I` or `+U` for that key is still emitted. Event-time mode (`on_time`) is not supported with updating input. Review Comment: people might not know +I or -U ########## flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/strategies/DeduplicateKeepFirstTypeStrategy.java: ########## @@ -0,0 +1,191 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.types.inference.strategies; + +import org.apache.flink.annotation.Internal; +import org.apache.flink.table.api.DataTypes; +import org.apache.flink.table.api.DataTypes.Field; +import org.apache.flink.table.api.ValidationException; +import org.apache.flink.table.api.dataview.ValueView; +import org.apache.flink.table.functions.TableSemantics; +import org.apache.flink.table.types.DataType; +import org.apache.flink.table.types.inference.CallContext; +import org.apache.flink.table.types.inference.InputTypeStrategy; +import org.apache.flink.table.types.inference.StateTypeStrategy; +import org.apache.flink.table.types.inference.TypeStrategy; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Optional; +import java.util.Set; +import java.util.stream.Collectors; + +/** Type strategies for the {@code DEDUPLICATE_KEEP_FIRST} process table function. */ +@Internal +public final class DeduplicateKeepFirstTypeStrategy { + + public static final int ARG_INPUT = 0; + public static final int ARG_STATE_TTL = 1; + public static final int ARG_RESET_TTL_ON_DUPLICATE = 2; + private static final String SEEN_STATE_NAME = "seen"; + private static final String CANDIDATE_STATE_NAME = "candidate"; + + public static final InputTypeStrategy INPUT_TYPE_STRATEGY = + new ValidationOnlyInputTypeStrategy() { + @Override + public Optional<List<DataType>> inferInputTypes( + final CallContext callContext, final boolean throwOnFailure) { + final Optional<List<DataType>> stateTtlError = + validateStateTtl(callContext, throwOnFailure); + if (stateTtlError.isPresent()) { + return stateTtlError; + } + + final Optional<List<DataType>> resetTtlError = + validateResetTtlOnDuplicate(callContext, throwOnFailure); + if (resetTtlError.isPresent()) { + return resetTtlError; + } + + return Optional.of(callContext.getArgumentDataTypes()); + } + }; + + public static final TypeStrategy OUTPUT_TYPE_STRATEGY = + callContext -> { + final TableSemantics tableSemantics = + callContext + .getTableSemantics(ARG_INPUT) + .orElseThrow( + () -> + new ValidationException( + "First argument must be a table for DEDUPLICATE_KEEP_FIRST.")); + + final List<Field> inputFields = DataType.getFields(tableSemantics.dataType()); + final List<Field> outputFields = + Arrays.stream( + ChangelogTypeStrategyUtils.computeOutputIndices( + tableSemantics)) + .mapToObj(inputFields::get) + .collect(Collectors.toList()); + return Optional.of(DataTypes.ROW(outputFields).notNull()); + }; + + public static final StateTypeStrategy SEEN_STATE_TYPE_STRATEGY = + new StateTypeStrategy() { + @Override + public Optional<DataType> inferType(final CallContext callContext) { + return Optional.of( + ValueView.newValueViewDataType(DataTypes.BOOLEAN().notNull())); + } + + @Override + public Optional<Duration> getTimeToLive(final CallContext callContext) { + return callContext.getArgumentValue(ARG_STATE_TTL, Duration.class); + } + }; + + public static final StateTypeStrategy CANDIDATE_STATE_TYPE_STRATEGY = + new StateTypeStrategy() { + @Override + public Optional<DataType> inferType(final CallContext callContext) { + final TableSemantics tableSemantics = + callContext + .getTableSemantics(ARG_INPUT) + .orElseThrow( + () -> + new ValidationException( + "First argument must be a table for DEDUPLICATE_KEEP_FIRST.")); + final List<Field> inputFields = DataType.getFields(tableSemantics.dataType()); + final List<Field> candidateFields = + Arrays.stream( + ChangelogTypeStrategyUtils.computeOutputIndices( + tableSemantics)) + .mapToObj(inputFields::get) + .collect(Collectors.toCollection(ArrayList::new)); + // The trailing timestamp field is read positionally at runtime, so its name + // only needs to avoid clashing with a payload column of the same name. + final Set<String> takenNames = + candidateFields.stream() + .map(Field::getName) + .collect(Collectors.toSet()); + String timestampField = "event_time"; + for (int i = 0; takenNames.contains(timestampField); i++) { Review Comment: or even better: just nest one more time ########## docs/content.zh/docs/sql/reference/queries/deduplicate-keep-first.md: ########## @@ -0,0 +1,187 @@ +--- +title: "Deduplicate Keep First" +weight: 16 +type: docs +--- +<!-- +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +--> + +# Deduplicate Keep First + +{{< label Streaming >}} + +Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, append-in / append-out alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern. + +| Function | Description | +|:---------|:------------| +| [DEDUPLICATE_KEEP_FIRST](#deduplicate_keep_first) | Keeps the first record per key as an insert-only stream (first arrival, or earliest event time with `on_time`) | + +## DEDUPLICATE_KEEP_FIRST + +The `DEDUPLICATE_KEEP_FIRST` PTF keeps the first record per key and drops every later record for that key. It supports two ordering modes: + +* **Processing time (default):** without `on_time`, the first record *observed* for a key is emitted; later records for that key are dropped while its state exists. This requires no watermark or event-time attribute. +* **Event time:** with `on_time`, the record with the *smallest event time* per key is kept and emitted once the watermark makes the choice final. A record arriving later but carrying an earlier event time replaces the candidate before finalization; records that can no longer become the earliest are dropped as late. + +The input may be insert-only or updating; the output is insert-only in every case. + +### Syntax + +```sql +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE source_table [PARTITION BY key_col [, key_col2 ...]], + [on_time => DESCRIPTOR(rowtime_column),] + [state_ttl => <interval>,] + [reset_ttl_on_duplicate => <boolean>] +) +``` + +### Parameters + +| Parameter | Required | Description | +|:----------|:---------|:------------| +| `input` | Yes | The input table (insert-only or updating). Use `PARTITION BY` to deduplicate per key; all rows for a key are routed to the same parallel instance. Without `PARTITION BY`, the whole input is a single group processed at parallelism 1, so only the first row of the entire stream is emitted and every later row is dropped regardless of its content. | +| `on_time` | No | A `DESCRIPTOR` naming a single rowtime attribute. When provided, the row with the smallest event time per key is kept and emitted once the watermark passes that timestamp (late rows are dropped). When omitted, the function keeps the first row observed for a key (arrival order). Requires insert-only input; combining `on_time` with an updating input is rejected at planning time. | Review Comment: ```suggestion | `on_time` | No | A `DESCRIPTOR` naming a single watermarked column. When provided, the row with the smallest event time per key is kept and emitted once the watermark passes that timestamp (late rows are dropped). When omitted, the function keeps the first row observed for a key (arrival order). Requires append-only input; combining `on_time` with an updating input is rejected at planning time. | ``` ########## flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/strategies/DeduplicateKeepFirstTypeStrategy.java: ########## @@ -0,0 +1,191 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.types.inference.strategies; + +import org.apache.flink.annotation.Internal; +import org.apache.flink.table.api.DataTypes; +import org.apache.flink.table.api.DataTypes.Field; +import org.apache.flink.table.api.ValidationException; +import org.apache.flink.table.api.dataview.ValueView; +import org.apache.flink.table.functions.TableSemantics; +import org.apache.flink.table.types.DataType; +import org.apache.flink.table.types.inference.CallContext; +import org.apache.flink.table.types.inference.InputTypeStrategy; +import org.apache.flink.table.types.inference.StateTypeStrategy; +import org.apache.flink.table.types.inference.TypeStrategy; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Optional; +import java.util.Set; +import java.util.stream.Collectors; + +/** Type strategies for the {@code DEDUPLICATE_KEEP_FIRST} process table function. */ +@Internal +public final class DeduplicateKeepFirstTypeStrategy { + + public static final int ARG_INPUT = 0; + public static final int ARG_STATE_TTL = 1; + public static final int ARG_RESET_TTL_ON_DUPLICATE = 2; + private static final String SEEN_STATE_NAME = "seen"; + private static final String CANDIDATE_STATE_NAME = "candidate"; + + public static final InputTypeStrategy INPUT_TYPE_STRATEGY = + new ValidationOnlyInputTypeStrategy() { + @Override + public Optional<List<DataType>> inferInputTypes( + final CallContext callContext, final boolean throwOnFailure) { + final Optional<List<DataType>> stateTtlError = + validateStateTtl(callContext, throwOnFailure); + if (stateTtlError.isPresent()) { + return stateTtlError; + } + + final Optional<List<DataType>> resetTtlError = + validateResetTtlOnDuplicate(callContext, throwOnFailure); + if (resetTtlError.isPresent()) { + return resetTtlError; + } + + return Optional.of(callContext.getArgumentDataTypes()); + } + }; + + public static final TypeStrategy OUTPUT_TYPE_STRATEGY = + callContext -> { + final TableSemantics tableSemantics = + callContext + .getTableSemantics(ARG_INPUT) + .orElseThrow( + () -> + new ValidationException( + "First argument must be a table for DEDUPLICATE_KEEP_FIRST.")); + + final List<Field> inputFields = DataType.getFields(tableSemantics.dataType()); + final List<Field> outputFields = + Arrays.stream( + ChangelogTypeStrategyUtils.computeOutputIndices( + tableSemantics)) + .mapToObj(inputFields::get) + .collect(Collectors.toList()); + return Optional.of(DataTypes.ROW(outputFields).notNull()); + }; + + public static final StateTypeStrategy SEEN_STATE_TYPE_STRATEGY = + new StateTypeStrategy() { + @Override + public Optional<DataType> inferType(final CallContext callContext) { + return Optional.of( + ValueView.newValueViewDataType(DataTypes.BOOLEAN().notNull())); + } + + @Override + public Optional<Duration> getTimeToLive(final CallContext callContext) { + return callContext.getArgumentValue(ARG_STATE_TTL, Duration.class); + } + }; + + public static final StateTypeStrategy CANDIDATE_STATE_TYPE_STRATEGY = + new StateTypeStrategy() { + @Override + public Optional<DataType> inferType(final CallContext callContext) { + final TableSemantics tableSemantics = + callContext + .getTableSemantics(ARG_INPUT) + .orElseThrow( + () -> + new ValidationException( + "First argument must be a table for DEDUPLICATE_KEEP_FIRST.")); + final List<Field> inputFields = DataType.getFields(tableSemantics.dataType()); + final List<Field> candidateFields = + Arrays.stream( + ChangelogTypeStrategyUtils.computeOutputIndices( + tableSemantics)) + .mapToObj(inputFields::get) + .collect(Collectors.toCollection(ArrayList::new)); + // The trailing timestamp field is read positionally at runtime, so its name + // only needs to avoid clashing with a payload column of the same name. + final Set<String> takenNames = + candidateFields.stream() + .map(Field::getName) + .collect(Collectors.toSet()); + String timestampField = "event_time"; + for (int i = 0; takenNames.contains(timestampField); i++) { + timestampField = "event_time_" + i; + } + candidateFields.add(DataTypes.FIELD(timestampField, DataTypes.BIGINT())); + return Optional.of( Review Comment: We don't always need this state right? ########## docs/content.zh/docs/sql/reference/queries/deduplicate-keep-first.md: ########## @@ -0,0 +1,187 @@ +--- +title: "Deduplicate Keep First" +weight: 16 +type: docs +--- +<!-- +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +--> + +# Deduplicate Keep First + +{{< label Streaming >}} + +Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, append-in / append-out alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern. + +| Function | Description | +|:---------|:------------| +| [DEDUPLICATE_KEEP_FIRST](#deduplicate_keep_first) | Keeps the first record per key as an insert-only stream (first arrival, or earliest event time with `on_time`) | + +## DEDUPLICATE_KEEP_FIRST + +The `DEDUPLICATE_KEEP_FIRST` PTF keeps the first record per key and drops every later record for that key. It supports two ordering modes: + +* **Processing time (default):** without `on_time`, the first record *observed* for a key is emitted; later records for that key are dropped while its state exists. This requires no watermark or event-time attribute. +* **Event time:** with `on_time`, the record with the *smallest event time* per key is kept and emitted once the watermark makes the choice final. A record arriving later but carrying an earlier event time replaces the candidate before finalization; records that can no longer become the earliest are dropped as late. + +The input may be insert-only or updating; the output is insert-only in every case. + +### Syntax + +```sql +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE source_table [PARTITION BY key_col [, key_col2 ...]], + [on_time => DESCRIPTOR(rowtime_column),] + [state_ttl => <interval>,] + [reset_ttl_on_duplicate => <boolean>] +) +``` + +### Parameters + +| Parameter | Required | Description | +|:----------|:---------|:------------| +| `input` | Yes | The input table (insert-only or updating). Use `PARTITION BY` to deduplicate per key; all rows for a key are routed to the same parallel instance. Without `PARTITION BY`, the whole input is a single group processed at parallelism 1, so only the first row of the entire stream is emitted and every later row is dropped regardless of its content. | +| `on_time` | No | A `DESCRIPTOR` naming a single rowtime attribute. When provided, the row with the smallest event time per key is kept and emitted once the watermark passes that timestamp (late rows are dropped). When omitted, the function keeps the first row observed for a key (arrival order). Requires insert-only input; combining `on_time` with an updating input is rejected at planning time. | +| `state_ttl` | No | An `INTERVAL` giving the processing-time retention for the per-key deduplication state. If omitted, falls back to `table.exec.state.ttl`, which retains state indefinitely at its default of 0. After a key's state expires, a later record for that key starts a new deduplication period and may be emitted again. `INTERVAL '0'` disables retention. | Review Comment: so there is a difference between 0 and optional. Optional meaning one can still change it later. 0 disabled for ever. ########## docs/content.zh/docs/sql/reference/queries/deduplicate-keep-first.md: ########## @@ -0,0 +1,187 @@ +--- +title: "Deduplicate Keep First" +weight: 16 +type: docs +--- +<!-- +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +--> + +# Deduplicate Keep First + +{{< label Streaming >}} + +Flink SQL provides the `DEDUPLICATE_KEEP_FIRST` process table function (PTF) for removing duplicate rows while keeping only the **first** record per key. Its output is always **insert-only**, which makes it a concise, append-in / append-out alternative to the `ROW_NUMBER()` [deduplication]({{< ref "docs/sql/reference/queries/deduplication" >}}) pattern. + +| Function | Description | +|:---------|:------------| +| [DEDUPLICATE_KEEP_FIRST](#deduplicate_keep_first) | Keeps the first record per key as an insert-only stream (first arrival, or earliest event time with `on_time`) | + +## DEDUPLICATE_KEEP_FIRST + +The `DEDUPLICATE_KEEP_FIRST` PTF keeps the first record per key and drops every later record for that key. It supports two ordering modes: + +* **Processing time (default):** without `on_time`, the first record *observed* for a key is emitted; later records for that key are dropped while its state exists. This requires no watermark or event-time attribute. +* **Event time:** with `on_time`, the record with the *smallest event time* per key is kept and emitted once the watermark makes the choice final. A record arriving later but carrying an earlier event time replaces the candidate before finalization; records that can no longer become the earliest are dropped as late. + +The input may be insert-only or updating; the output is insert-only in every case. + +### Syntax + +```sql +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE source_table [PARTITION BY key_col [, key_col2 ...]], + [on_time => DESCRIPTOR(rowtime_column),] + [state_ttl => <interval>,] + [reset_ttl_on_duplicate => <boolean>] +) +``` + +### Parameters + +| Parameter | Required | Description | +|:----------|:---------|:------------| +| `input` | Yes | The input table (insert-only or updating). Use `PARTITION BY` to deduplicate per key; all rows for a key are routed to the same parallel instance. Without `PARTITION BY`, the whole input is a single group processed at parallelism 1, so only the first row of the entire stream is emitted and every later row is dropped regardless of its content. | +| `on_time` | No | A `DESCRIPTOR` naming a single rowtime attribute. When provided, the row with the smallest event time per key is kept and emitted once the watermark passes that timestamp (late rows are dropped). When omitted, the function keeps the first row observed for a key (arrival order). Requires insert-only input; combining `on_time` with an updating input is rejected at planning time. | +| `state_ttl` | No | An `INTERVAL` giving the processing-time retention for the per-key deduplication state. If omitted, falls back to `table.exec.state.ttl`, which retains state indefinitely at its default of 0. After a key's state expires, a later record for that key starts a new deduplication period and may be emitted again. `INTERVAL '0'` disables retention. | +| `reset_ttl_on_duplicate` | No | Whether a later duplicate refreshes the key's `state_ttl`. Defaults to `TRUE`, so retention tracks the most recent occurrence of a key. Only meaningful when a TTL applies, whether set through `state_ttl` or inherited from `table.exec.state.ttl`. | + +### Output Schema + +When `PARTITION BY` is used, the partition key columns are prepended to the output, followed by the remaining (non-key) input columns. In event-time mode a rowtime column is appended. + +``` +[partition_key_columns] + [remaining_input_columns] +``` + +### Examples + +#### Keyed keep-first (processing time) + +```sql +-- Input (append-only): +-- +I[user_name:'Vas', action:'login'] +-- +I[user_name:'Vas', action:'click'] + +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE user_events PARTITION BY user_name +) + +-- Output (insert-only): +-- +I[user_name:'Vas', action:'login'] +``` + +The first row for `Vas` is emitted; the later `click` on the same key is dropped. + +#### Whole input, no PARTITION BY + +```sql +-- Input (append-only): +-- +I[user_name:'Vas', action:'login'] +-- +I[user_name:'Alice', action:'click'] +-- +I[user_name:'Bob', action:'view'] + +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE user_events +) + +-- Output (insert-only): +-- +I[user_name:'Vas', action:'login'] +``` + +Without `PARTITION BY` the whole input is a single group at parallelism 1, so only the first row of the entire stream is emitted; every later row is dropped regardless of its content. + +#### Bounded state with state_ttl + +```sql +-- Input (append-only): +-- +I[user_name:'Vas', action:'login'] +-- +I[user_name:'Vas', action:'click'] +-- +I[user_name:'Vas', action:'view'] + +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE user_events PARTITION BY user_name, + state_ttl => INTERVAL '5' SECOND, + reset_ttl_on_duplicate => FALSE +) + +-- Output (insert-only): +-- +I[user_name:'Vas', action:'login'] +``` + +`state_ttl` bounds how long per-key state is retained (processing time). While the state exists, duplicates are dropped and only the first row per key is emitted; after the state expires, a later record for the key starts a new deduplication period. With the default `reset_ttl_on_duplicate => TRUE`, each dropped duplicate refreshes the retention window; with `FALSE`, retention is measured from the first record for the key. + +#### Earliest by event time + +```sql +-- Source declares: WATERMARK FOR ts AS ts - INTERVAL '10' SECOND +-- Input (ts = event time): +-- +I[user_name:'Vas', action:'login', ts:3.000] +-- +I[user_name:'Vas', action:'click', ts:1.000] +-- +I[user_name:'Vas', action:'view', ts:5.000] + +SELECT user_name, action FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE user_events PARTITION BY user_name, + on_time => DESCRIPTOR(ts) +) + +-- Output (insert-only, once the watermark finalizes the earliest): +-- +I[user_name:'Vas', action:'click'] +``` + +With `on_time`, the record with the smallest event time per key is kept (here `click` at `ts 1.000`, not `login` which arrived first) and emitted once the watermark passes that timestamp. + +#### Updating input + +```sql +-- Input (updating changelog): +-- -D[user_name:'Vas', action:'logout'] key not seen yet -> ignored +-- +I[user_name:'Vas', action:'login'] first record for the key -> emitted +-- +I[user_name:'Vas', action:'login'] duplicate -> swallowed +-- -U[user_name:'Vas', action:'login'] retraction -> swallowed +-- +U[user_name:'Vas', action:'click'] update -> swallowed +-- -D[user_name:'Vas', action:'click'] delete -> swallowed + +SELECT * FROM DEDUPLICATE_KEEP_FIRST( + input => TABLE user_events PARTITION BY user_name +) + +-- Output (insert-only): +-- +I[user_name:'Vas', action:'login'] +``` + +The input may be an updating changelog. `DEDUPLICATE_KEEP_FIRST` keeps the first `+I` or `+U` observed per key and swallows every later change to that key (duplicate, `-U`, `+U`, `-D`), so the output stays insert-only. A `-U` or `-D` for a key that has not been seen yet is ignored and does not mark the key as seen, so the next `+I` or `+U` for that key is still emitted. Event-time mode (`on_time`) is not supported with updating input. Review Comment: link "updating changelog" to a page, currently PTF page until Gustavos PR is in -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
