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

Josh Rosen updated SPARK-58487:
-------------------------------
    Summary: Pandas UDFs and Arrow-enabled toPandas() silently corrupt nullable 
BIGINT values of magnitude 2^53 or more by upcasting the column to float64  
(was: "Pandas UDFs and Arrow-enabled toPandas() silently corrupt nullable 
BIGINT values of magnitude 2^53 or more by upcasting the column to float64)

> Pandas UDFs and Arrow-enabled toPandas() silently corrupt nullable BIGINT 
> values of magnitude 2^53 or more by upcasting the column to float64
> ---------------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: SPARK-58487
>                 URL: https://issues.apache.org/jira/browse/SPARK-58487
>             Project: Spark
>          Issue Type: Bug
>          Components: PySpark
>    Affects Versions: 4.2.0
>            Reporter: Josh Rosen
>            Priority: Major
>              Labels: correctness
>
> A nullable {{LongType}} column holding a value of magnitude at least {{2^53}} 
> is silently changed to the nearest representable {{float64}} as soon as a 
> NULL appears in the same Arrow batch.
> {code:python}
> from pyspark.sql.types import StructType, StructField, LongType
> from pyspark.sql.functions import pandas_udf, col
> import pandas as pd
> big = 9007199254740993          # 2**53 + 1, not exactly representable in 
> float64
> schema = StructType([StructField("v", LongType())])
> df = spark.createDataFrame([(big,), (None,)], schema=schema)
> df.toPandas()['v'].iloc[0]
> # 9007199254740992.0   -- off by one, and now a float
> @pandas_udf(LongType())
> def identity(s: pd.Series) -> pd.Series:
>     return s
> df.coalesce(1).withColumn("v2", identity(col("v"))).collect()
> # [Row(v=9007199254740993, v2=9007199254740992), ...]   -- v2 != v
> {code}
> Remove the NULL row and both paths are exact. The NULL is what triggers the 
> upcast.
> *This is a Spark self-inconsistency rather than a pandas limitation.* On the 
> same DataFrame, {{collect()}} and {{toPandas()}} with Arrow disabled both 
> return the exact integer; only the Arrow path differs. Two supported ways of 
> reading the same column disagree.
> *Mechanism.* pyarrow's default {{Array.to_pandas()}} cannot represent an 
> int64 column containing a null, because numpy has no nullable integer dtype, 
> so it upcasts the whole column to {{float64}} with NaN as the null sentinel. 
> {{float64}} has a 53-bit mantissa, so integers at or above {{2^53}} are not 
> all representable and the upcast rounds.
> Neither {{pyspark.sql.pandas.conversion._convert_arrow_table_to_pandas}} nor 
> the Arrow-column-to-Series conversion in {{pyspark.sql.pandas.serializers}} 
> passes {{types_mapper=}} to opt into pandas' nullable extension dtypes, which 
> would keep both the nulls and the precision.
> On the way back, {{{}create_arrow_array_from_pandas{}}} casts the 
> already-corrupted float64 Series to the declared integer type. The corrupted 
> value is a whole number, so the cast raises nothing and the wrong value is 
> written out as if correct.
> *The trigger is batch-scoped, which matters for anyone trying to reproduce 
> it.* The NULL and the large value must land in the same Arrow batch. A 
> two-row, two-partition example does not reproduce, because each batch is 
> either all-non-null or trivially the null row. 
> *Measured scope:*
> ||path||result||
> |{{collect()}}|exact|
> |{{{}toPandas(){}}}, Arrow disabled|exact|
> |{{{}toPandas(){}}}, Arrow enabled|corrupted|
> |scalar and scalar-iter {{pandas_udf}}|corrupted|
> |grouped-agg {{{}pandas_udf{}}}, NULL and value in one batch|corrupted|
> |{{struct<a: bigint>}} through the Arrow paths above|corrupted|
> |{{{}mapInPandas{}}}, {{{}applyInPandas{}}}, {{applyInArrow}}|exact|
> |Arrow-optimized plain Python UDF|exact, by construction (per-row Python 
> objects, no dtype decision)|
> |{{DecimalType(38, 18)}}|exact (pyarrow yields {{decimal.Decimal}} objects, 
> no float64 fallback)|
> *Potential fixes.* Have the Arrow-to-pandas conversion sites opt in to
> {{{}integer_object_nulls=True{}}}, or pass a {{types_mapper}} that maps 
> nullable Arrow integers to pandas nullable extension dtypes. 
> {{ArrowArrayToPandasConversion.convert_legacy}} already uses the
> former for array, map and struct columns, so the mechanism is present and 
> simply not applied to the scalar integral branches. Failing that, the 
> return-path cast should not silently accept a float that lost precision.
> {*}Related bug{*}: {{createDataFrame()}} on a pandas {{Int64}} 
> extension-dtype column containing a null infers DOUBLE and corrupts the same 
> magnitudes. 



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

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to