[ 
https://issues.apache.org/jira/browse/FLINK-40705?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Dian Fu closed FLINK-40705.
---------------------------
    Fix Version/s: 1.20.6
                   2.3.1
                   2.4.0
       Resolution: Fixed

Fixed in:
- master via 9a84b5399f164ef46b0bd7bf6df25692e94bcada and 
3c78d5ff191d436da258e1646995b17845f87c84
- release-2.3 via ec5e1b708219ec5bffb726cbf246769ba9801d2a
- release-1.20 via f8b86430c4da559ab6061bfd063b27b47ad26e35

> PyFlink pandas UDFs fail on MAP inputs across Arrow batches
> -----------------------------------------------------------
>
>                 Key: FLINK-40705
>                 URL: https://issues.apache.org/jira/browse/FLINK-40705
>             Project: Flink
>          Issue Type: Bug
>          Components: API / Python
>    Affects Versions: 1.17.0
>            Reporter: Liu Liu
>            Assignee: Liu Liu
>            Priority: Major
>              Labels: pull-request-available
>             Fix For: 1.20.6, 2.3.1, 2.4.0
>
>
> Pandas UDFs consuming MAP columns can fail when input records span multiple 
> Arrow batches.
> {{MapWriter}} does not override {{reset()}} to reset its key and value 
> writers. Consequently, the map vector restarts at the beginning of each 
> batch, while its child writers retain their previous counters. Subsequent 
> batches write map entries at incorrect positions.
> This can produce incorrect decoded values or fail with:
> {{Map array child array should have no nulls}}
> To reproduce, pass two records containing \{1: 11} and \{2: 22} in a 
> {{MAP<INT NOT NULL, INT>}} column to a pandas UDF, with:
> {code:java}
> t_env.get_config().set("python.fn-execution.arrow.batch.size", "1")
> t_env.get_config().set("python.fn-execution.bundle.size", "100") {code}
> The second single-record batch fails; setting the Arrow batch size to {{2}} 
> allows both records to succeed in one batch.
> The defect was introduced with MAP support in FLINK-30607. 



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to