[
https://issues.apache.org/jira/browse/FLINK-40856?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18124791#comment-18124791
]
Dian Fu edited comment on FLINK-40856 at 10/8/26 2:58 AM:
----------------------------------------------------------
[~beetleyibo] I have investigated engines such as Pandas, Polars, Daft and
their behavior are all throwing exceptions when the inputs have non-equal
lengths. I tend to align the behavior with them.
was (Author: dianfu):
[~beetleyibo] I have investigated the other engines such as Pandas, Polars,
Daft and their behavior are all throw exceptions when the inputs have non-equal
lengths.
> 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.
> - 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.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)