Hey y'all. Found some interesting behavior involvi...
# general
s
Hey y'all. Found some interesting behavior involving structs. For context, I have a pyiceberg table that I'm trying to read it. When I use
daft.read_iceberg
and try to execute on it I get the following error.
daft.exceptions.DaftCoreException: DaftError::ArrowError validity mask length must match the number of values
. Here is a code example:
Copy code
daft_df = daft.read_iceberg(table)
daft_df.to_pandas()
However if I use pyiceberg to read the table into pandas and then go to daft and back to pandas, it works without errors
Copy code
# iceberg table to pandas to daft and back to pandas 
p_pd = table.scan().to_pandas()
p_df = daft.from_pandas(p_pd)
p_pd2 = p_df.to_pandas()
I think I've isolated the issue to a column of dtype struct that has a
null
in it. In both paths, when I get into daft the schema for the issue column is the same
Struct[decimal_key: Decimal(precision=10, scale=2), string_key: Utf8]
. Also, with the working path. I can filter nulls and execute on that. Which makes me think it has something to do with
read_iceberg
Copy code
p_df.where(p_df["null_struct_column"].is_null()).to_pandas()
Any guidance would be appreciated!
Looks like
collect()
seems to work however
to_pandas()
and in our case
to_ray_dataset()
are the issues
a
Thanks for letting us know. @Kevin Wang any thoughts?
@Sam Gottlieb do you have a minimum reproducible example we can look at?
s
I need to find a way to share. It only works if I go from iceberg -> daft -> pandas/ray. Tried saving same DF as local parquet file and couldn't reproduce
t
doesn't work:
Copy code
from pyiceberg.catalog import load_catalog

catalog = load_catalog("glue", type="glue")
table = catalog.load_table(table_name)
df = daft.read_iceberg(table)

df = df.select(df["sale_price"].is_null())
df.to_pandas()
works:
Copy code
fields = [
    ("amount", pa.decimal128(10, 2)),
    ("currency", pa.string())
]
struct_type =  pa.struct(fields)
data = [
    {'amount': Decimal(99.99).quantize(Decimal('0.01')), 'currency': 'USD'},
    {'amount': Decimal(149.99).quantize(Decimal('0.01')), 'currency': 'EUR'},
    {'amount': Decimal(199.99).quantize(Decimal('0.01')), 'currency': 'GBP'},
    None
]
sale_price = pa.array(data,  type=struct_type)

table = pa.table([sale_price], names=['sale_price'])
df = daft.from_arrow(table)

df = df.select(df["sale_price"].is_null())
df.to_pandas()
schemas are the same for both:
Copy code
╞═════════════╪════════════════════════════════════════════════════════════════╡
│ sale_price  ┆ Struct[amount: Decimal(precision=10, scale=2), currency: Utf8] │
╰─────────────┴────────────────────────────────────────────────────────────────╯
is there something special about how data is read from iceberg that could be mangling structs?
k
Hi @Sam Gottlieb @Tabrez Mohammed would it be possible to get us an example of a file that causes this error?
I don't think it has to do with Iceberg but I may be wrong
t
What we're seeing from our latest iceberg snapshot is the column containing structs that has mixed nulls and non-null will fail to materialize. Here's a loop we did to test:
Copy code
feed_ids = df.select("feed_id").distinct().to_pydict()["feed_id"]
for feed_id in feed_ids:
    print(feed_id)
    feed_df = df.where(df["feed_id"] == feed_id)
    feed_df = feed_df.to_pandas()
    print(f"sale price null count: {feed_df['sale_price'].isnull().sum()}")
    print(f"sale price not null count: {feed_df['sale_price'].notnull().sum()}")
feed_id
is a partition key for the table. This prints successfully where a partition is all null or all non null. When it fails, it raises this:
Copy code
DaftCoreException: DaftError::External Unable to create arrow chunk from streaming file readers3://[...]/data/advertiser_id=1101l859/feed_id=1101l2560/00008-191-e1ac2758-b781-4dfd-bf31-0b4e570b12e7-00004.parquet: validity mask length must match the number of values
@Sam Gottlieb can you share this file? ^
s
Upon digging we were able to get a similar error. This file should reproduce
Copy code
import daft
# pd.read_parquet("/Users/sgottlieb/Downloads/00008-191-e1ac2758-b781-4dfd-bf31-0b4e570b12e7-00004.parquet")
daft_df = daft.read_parquet("00008-191-e1ac2758-b781-4dfd-bf31-0b4e570b12e7-00004.parquet")
daft_df.to_pandas()
k
Taking a look
@Desmond Cheong would you know what's going on? It seems to only be failing on reading the UTF8
currency
field of the
sale_price
column, where arrow2 sees a validity bitmap of length 2797 instead of the array length 3571. 2797 also happens to be the number of nulls in the
sale_price
column, not sure why that would be related but it'd be a weird coincidence
👀 2
Wait Desmond is out of office, @Sammy Sidhu would you mind looking into this?
j
He’s on a flight, I’m poking around at it rn
😄 1
🙏 1
🙏🏽 1
So I have a fix, but I haven’t managed to really build out the root cause analysis. It has something to do with the way required fields interact with nested parquet decoding 😬
d
Haha took a quick look too
not sure why that would be related but it'd be a weird coincidence
It's never a coincidence, something smells 👃 And to me it smells over here https://github.com/Eventual-Inc/Daft/blob/95a61d26bfb198c590570229a81f7e3bec0049a0/src/arrow2/src/io/parquet/read/deserialize/binary/nested.rs#L[…]2 If the size of validity matches the number of nulls, then we're pushing to the validity mask only when there's a null. In the
push_valid
case we're somehow in the
RequiredDictionary
state (no nulls, dict-encoded) even though having nulls means that we should be in the
OptionalDictionary
state (nulls, dict-encoded). Slightly further up in the same file https://github.com/Eventual-Inc/Daft/blob/95a61d26bfb198c590570229a81f7e3bec0049a0/src/arrow2/src/io/parquet/read/deserialize/binary/nested.rs#L[…]7 we're indeed building a
RequiredDictionary
. So the bug either lives further upstream or there's an encoding issue with
page.descriptor.primitive_type.field_info.repetition == Repetition::Optional
j
Yes… the fix is very simple: https://github.com/Eventual-Inc/Daft/pull/3572
d
yeah but that's odd to me because
RequiredDictionary
should never need a validity mask...
👍 1
j
I think somewhere in the code it’s not handling the
RequiredDictionary
case correctly and thus using the “wrongly built” validity
Oh I see what you mean — it should be
OptionalDictionary
but somehow it’s
RequiredDictionary
👍 1
d
yeah, here's a smaller repro if it helps. Easier to see the metadata lol nope doesn't actually repro, just looks nice
j
My guess is that it’s something to do with
optional group
containing a
required
column:
😮 1
👀 1
Which is a pretty interesting choice by whatever application wrote this Parquet file 😬
t
spark likes to assume that if a column is filled it must be required
This was written using the dataframe writer v2
d
I think that bit of metadata is saying that
sale_price
is optional, but if it's defined then
currency
is required. But I guess since
currency
is in an optional struct, it should also be optional
👍 1
t
If I give that same parquet file from another snapshot, would be able to see if that metadata differs? I'm not sure how you dumped that datatype info
j
I used the parquet CLI tool: https://formulae.brew.sh/formula/parquet-cli
parquet meta <path_to_file>
is how it works
👀 1
t
Interesting, older file for same partition:
Copy code
optional group sale_price = 10 {
    optional int64 amount (DECIMAL(10,0)) = 37;
    optional binary currency (STRING) = 38;
  }
😓 2
the change from 10,0 > 10,2 is on purpose, but I didn't touch anything for how currency is set
j
Ok I think we’ve found the culprit haha. Likely some internal application logic in Spark that’s writing these weird optional::required structs
I think that bit of metadata is saying that
sale_price
is optional, but if it’s defined then
currency
is required. But I guess since
currency
is in an optional struct, it should also be optional
This makes sense, I’m guessing then that this meaning of
required
is getting picked up by us also when parsing the actual currency array. I wonder if there’s an easy way to make
currency
optional as well by detecting that it is the descendant of an optional struct.
One really dumb solution I can think of is to default to Optional parsing behavior if we every see that it is a nested field 🤣
t
Is there a bug on spark where it should have coalesced to null/optional in the schema? I guess the files we can read have that field as optional
d
I’m guessing then that this meaning of
required
is getting picked up by us also when parsing the actual currency array.
Yeah looks like it. This is what we deserialize as the nested field's metadata
field_info: FieldInfo { name: "currency", repetition: Required, id: Some(38)
.
Is there a bug on spark where it should have coalesced to null/optional in the schema?
I think this is per-row group metadata, so technically it's valid. The overall schema for the column should correctly mark the column as nullable
Ok here's my proposal. In short: nested children should be nullable if any of their parents are nullable. Otherwise we simply check if the current field is nullable. Thoughts @jay?
j
Sounds reasonable I think…
d
back from vacation, here's the fix https://github.com/Eventual-Inc/Daft/pull/3598
👀 3
s
Thanks for the quick support y'all. Tested with the 0.4 release last night and it no issues
🙌 1
d
Glad to hear it!