[
https://issues.apache.org/jira/browse/FLINK-40856?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Dian Fu updated FLINK-40856:
----------------------------
Description:
DataFrame.explode() currently accepts only one collection column or expression.
Exploding multiple columns by chaining explode() calls produces a Cartesian
product, so users cannot expand logically paired arrays such as item names and
quantities.
This issue proposes extending DataFrame.explode() to accept multiple ARRAY
columns or expressions and explode them with positional zip semantics.
Proposed API:
def explode(
self,
column: Union[
str,
Expression,
List[Union[str, Expression]],
],
*,
output_column: Optional[Union[str, List[str]]] = None,
ignore_empty_and_null: bool = False,
) -> "DataFrame"
Proposed behavior:
- A list containing multiple inputs enables zip-style explode.
- All inputs in a multi-column explode must have ARRAY types. Multi-column MAP
and MULTISET inputs are out of scope.
- Elements are paired by their ordinal position.
- A directly referenced input column is replaced at its original position.
- For each input row, all input ARRAYs must have the same effective length.
Empty and NULL ARRAYs have an effective length of zero. If the effective
lengths differ, the operation fails at execution time. Otherwise, elements are
paired by their ordinal position.
- The output of a computed ARRAY expression is appended after the original
columns.
- output_column, when provided, must contain exactly one unique name per input
ARRAY.
- When output_column is omitted, each output reuses the resolved input or
expression name.
- ARRAY<ROW<...>> elements remain a single ROW column. Users can explicitly
call Expression.flatten in a subsequent select when needed.
- Length validation is performed before applying ignore_empty_and_null.
Therefore, ignore_empty_and_null does not suppress or pad unequal-length inputs.
Example:
Input:
+----+-------------------+-------+------------+
| id | items | label | quantities |
+----+-------------------+-------+------------+
| 1 | [apple, pear] | A | [2, 3] |
+----+-------------------+-------+------------+
df.explode(
["items", "quantities"],
output_column=["item", "quantity"],
)
Result:
+----+-------+-------+----------+
| id | item | label | quantity |
+----+-------+-------+----------+
| 1 | apple | A | 2 |
| 1 | pear | A | 3 |
+----+-------+-------+----------+
Empty and null handling:
- For each input row, all input ARRAYs must have the same effective length.
Empty and NULL ARRAYs have an effective length of zero.
- If the effective lengths differ, the operation fails at execution time,
regardless of ignore_empty_and_null.
- Non-empty ARRAYs of equal length are exploded in positional order.
- If all effective lengths are zero, the input row produces no result when
ignore_empty_and_null=True.
- If all effective lengths are zero, the input row produces one result row with
NULL for every exploded output when
ignore_empty_and_null=False.
- NULL elements inside a non-empty ARRAY are preserved at their corresponding
positions, subject to the existing FLINK-40658
limitation for actual NULL ROW elements.
Implementation note:
The implementation can use one ARRAY as an ordinality driver:
UNNEST(driver_array) WITH ORDINALITY
Before expansion, it must validate that all input ARRAYs have the same
effective length, for example by comparing
COALESCE(CARDINALITY(array), 0). A mismatch must fail the operation at
execution time.
After successful validation, values from the remaining ARRAYs can be retrieved
using the driver's ordinal position. The validation must also be evaluated when
the driver ARRAY is empty or NULL, so that an unequal-length row cannot be
silently discarded before the mismatch is detected.
This uses a single UNNEST and avoids both Cartesian products and shortest-input
truncation. Native multi-input UNNEST should not be assumed to provide the
required strict positional semantics.
Acceptance criteria:
- Support two or more ARRAY columns and expressions.
- Verify that equal-length ARRAYs are exploded positionally and that
unequal-length ARRAYs fail at execution time for both ignore_empty_and_null
modes.
- Cover empty and NULL ARRAYs with both ignore_empty_and_null modes.
- Preserve the original positions of directly referenced columns.
- Append outputs for computed expressions.
- Preserve ARRAY<ROW> elements without implicit flattening.
- Validate non-ARRAY inputs, duplicate direct inputs, output-name counts and
name conflicts.
- Add batch and streaming coverage without materially increasing the number of
IT jobs.
- Document the multi-column behavior and provide an API example.
- Update FLIP-591 to include the public API extension and its semantics.
- Verify length validation for mixed non-empty, empty, and NULL ARRAY inputs,
including mismatches where the ordinality driver is empty or NULL.
was:
DataFrame.explode() currently accepts only one collection column or expression.
Exploding multiple columns by chaining explode() calls produces a Cartesian
product, so users cannot expand logically paired arrays such as item names and
quantities.
This issue proposes extending DataFrame.explode() to accept multiple ARRAY
columns or expressions and explode them with positional zip semantics.
Proposed API:
def explode(
self,
column: Union[
str,
Expression,
List[Union[str, Expression]],
],
*,
output_column: Optional[Union[str, List[str]]] = None,
ignore_empty_and_null: bool = False,
) -> "DataFrame"
Proposed behavior:
- A list containing multiple inputs enables zip-style explode.
- All inputs in a multi-column explode must have ARRAY types. Multi-column MAP
and MULTISET inputs are out of scope.
- Elements are paired by their ordinal position.
- The result is truncated to the shortest ARRAY, matching Python zip semantics.
- A directly referenced input column is replaced at its original position.
- The output of a computed ARRAY expression is appended after the original
columns.
- output_column, when provided, must contain exactly one unique name per input
ARRAY.
- When output_column is omitted, each output reuses the resolved input or
expression name.
- ARRAY<ROW<...>> elements remain a single ROW column. Users can explicitly
call Expression.flatten in a subsequent select when needed.
Example:
Input:
+----+-------------------+-------+------------+
| id | items | label | quantities |
+----+-------------------+-------+------------+
| 1 | [apple, pear] | A | [2, 3, 4] |
+----+-------------------+-------+------------+
df.explode(
["items", "quantities"],
output_column=["item", "quantity"],
)
Result:
+----+-------+-------+----------+
| id | item | label | quantity |
+----+-------+-------+----------+
| 1 | apple | A | 2 |
| 1 | pear | A | 3 |
+----+-------+-------+----------+
Empty and null handling:
- For each input row, all input ARRAYs must have the same effective length;
empty and NULL ARRAYs have an effective length of zero. If the lengths differ,
the operation fails with a clear error.
- Non-empty ARRAYs are exploded in zip order, producing one output row for each
position.
- If all input ARRAYs are empty or NULL, the input row produces no result when
ignore_empty_and_null=True.
- If all input ARRAYs are empty or NULL, the input row produces one result row
with NULL for every exploded output when ignore_empty_and_null=False.
- NULL elements inside a non-empty ARRAY are preserved at their corresponding
positions, subject to the existing FLINK-40658 limitation for actual NULL ROW
elements.
Implementation note:
The implementation can use one ARRAY as an ordinality driver:
UNNEST(driver_array) WITH ORDINALITY
The remaining ARRAY values can be retrieved using their ordinal index, with the
result filtered by the minimum input cardinality. This avoids the
Cartesian-product behavior of independent lateral joins. Native multi-input
UNNEST should not be assumed to provide the required zip semantics.
Acceptance criteria:
- Support two or more ARRAY columns and expressions.
- Cover equal and unequal ARRAY lengths.
- Cover empty and NULL ARRAYs with both ignore_empty_and_null modes.
- Preserve the original positions of directly referenced columns.
- Append outputs for computed expressions.
- Preserve ARRAY<ROW> elements without implicit flattening.
- Validate non-ARRAY inputs, duplicate direct inputs, output-name counts and
name conflicts.
- Add batch and streaming coverage without materially increasing the number of
IT jobs.
- Document the multi-column behavior and provide an API example.
- Update FLIP-591 to include the public API extension and its semantics.
> Support zip-style multi-column explode in DataFrame API
> -------------------------------------------------------
>
> Key: FLINK-40856
> URL: https://issues.apache.org/jira/browse/FLINK-40856
> Project: Flink
> Issue Type: Sub-task
> Components: API / Python
> Reporter: Dian Fu
> Priority: Major
>
> DataFrame.explode() currently accepts only one collection column or
> expression. Exploding multiple columns by chaining explode() calls produces a
> Cartesian product, so users cannot expand logically paired arrays such as
> item names and quantities.
> This issue proposes extending DataFrame.explode() to accept multiple ARRAY
> columns or expressions and explode them with positional zip semantics.
> Proposed API:
> def explode(
> self,
> column: Union[
> str,
> Expression,
> List[Union[str, Expression]],
> ],
> *,
> output_column: Optional[Union[str, List[str]]] = None,
> ignore_empty_and_null: bool = False,
> ) -> "DataFrame"
> Proposed behavior:
> - A list containing multiple inputs enables zip-style explode.
> - All inputs in a multi-column explode must have ARRAY types. Multi-column
> MAP and MULTISET inputs are out of scope.
> - Elements are paired by their ordinal position.
> - A directly referenced input column is replaced at its original position.
> - For each input row, all input ARRAYs must have the same effective length.
> Empty and NULL ARRAYs have an effective length of zero. If the effective
> lengths differ, the operation fails at execution time. Otherwise, elements
> are paired by their ordinal position.
> - The output of a computed ARRAY expression is appended after the original
> columns.
> - output_column, when provided, must contain exactly one unique name per
> input ARRAY.
> - When output_column is omitted, each output reuses the resolved input or
> expression name.
> - ARRAY<ROW<...>> elements remain a single ROW column. Users can explicitly
> call Expression.flatten in a subsequent select when needed.
> - Length validation is performed before applying ignore_empty_and_null.
> Therefore, ignore_empty_and_null does not suppress or pad unequal-length
> inputs.
> Example:
> Input:
> +----+-------------------+-------+------------+
> | id | items | label | quantities |
> +----+-------------------+-------+------------+
> | 1 | [apple, pear] | A | [2, 3] |
> +----+-------------------+-------+------------+
> df.explode(
> ["items", "quantities"],
> output_column=["item", "quantity"],
> )
> Result:
> +----+-------+-------+----------+
> | id | item | label | quantity |
> +----+-------+-------+----------+
> | 1 | apple | A | 2 |
> | 1 | pear | A | 3 |
> +----+-------+-------+----------+
> Empty and null handling:
> - For each input row, all input ARRAYs must have the same effective length.
> Empty and NULL ARRAYs have an effective length of zero.
> - If the effective lengths differ, the operation fails at execution time,
> regardless of ignore_empty_and_null.
> - Non-empty ARRAYs of equal length are exploded in positional order.
> - If all effective lengths are zero, the input row produces no result when
> ignore_empty_and_null=True.
> - If all effective lengths are zero, the input row produces one result row
> with NULL for every exploded output when
> ignore_empty_and_null=False.
> - NULL elements inside a non-empty ARRAY are preserved at their corresponding
> positions, subject to the existing FLINK-40658
> limitation for actual NULL ROW elements.
> Implementation note:
> The implementation can use one ARRAY as an ordinality driver:
> UNNEST(driver_array) WITH ORDINALITY
> Before expansion, it must validate that all input ARRAYs have the same
> effective length, for example by comparing
> COALESCE(CARDINALITY(array), 0). A mismatch must fail the operation at
> execution time.
> After successful validation, values from the remaining ARRAYs can be
> retrieved using the driver's ordinal position. The validation must also be
> evaluated when the driver ARRAY is empty or NULL, so that an unequal-length
> row cannot be silently discarded before the mismatch is detected.
> This uses a single UNNEST and avoids both Cartesian products and
> shortest-input truncation. Native multi-input UNNEST should not be assumed to
> provide the required strict positional semantics.
> Acceptance criteria:
> - Support two or more ARRAY columns and expressions.
> - Verify that equal-length ARRAYs are exploded positionally and that
> unequal-length ARRAYs fail at execution time for both ignore_empty_and_null
> modes.
> - Cover empty and NULL ARRAYs with both ignore_empty_and_null modes.
> - Preserve the original positions of directly referenced columns.
> - Append outputs for computed expressions.
> - Preserve ARRAY<ROW> elements without implicit flattening.
> - Validate non-ARRAY inputs, duplicate direct inputs, output-name counts and
> name conflicts.
> - Add batch and streaming coverage without materially increasing the number
> of IT jobs.
> - Document the multi-column behavior and provide an API example.
> - Update FLIP-591 to include the public API extension and its semantics.
> - Verify length validation for mixed non-empty, empty, and NULL ARRAY inputs,
> including mismatches where the ordinality driver is empty or NULL.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)