Dian Fu created FLINK-40856:
-------------------------------
Summary: 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
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)