Sam Gottlieb
12/13/2024, 4:46 PMdaft.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:
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
# 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
p_df.where(p_df["null_struct_column"].is_null()).to_pandas()
Any guidance would be appreciated!Sam Gottlieb
12/13/2024, 5:36 PMcollect() seems to work however to_pandas() and in our case to_ray_dataset() are the issuesAndrew Gazelka
12/13/2024, 5:40 PMAndrew Gazelka
12/13/2024, 5:45 PMSam Gottlieb
12/13/2024, 6:04 PMTabrez Mohammed
12/13/2024, 6:24 PMfrom 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:
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:
╞═════════════╪════════════════════════════════════════════════════════════════╡
│ sale_price ┆ Struct[amount: Decimal(precision=10, scale=2), currency: Utf8] │
╰─────────────┴────────────────────────────────────────────────────────────────╯Tabrez Mohammed
12/13/2024, 6:29 PMKevin Wang
12/13/2024, 7:17 PMKevin Wang
12/13/2024, 7:17 PMTabrez Mohammed
12/13/2024, 7:42 PMfeed_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:
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? ^Sam Gottlieb
12/13/2024, 7:46 PMSam Gottlieb
12/13/2024, 7:46 PMimport 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()Kevin Wang
12/13/2024, 9:10 PMKevin Wang
12/13/2024, 9:42 PMcurrency 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 coincidenceKevin Wang
12/13/2024, 11:16 PMjay
12/14/2024, 12:21 AMjay
12/14/2024, 1:38 AMDesmond Cheong
12/14/2024, 1:39 AMnot sure why that would be related but it'd be a weird coincidenceIt'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::Optionaljay
12/14/2024, 1:39 AMDesmond Cheong
12/14/2024, 1:40 AMRequiredDictionary should never need a validity mask...jay
12/14/2024, 1:40 AMRequiredDictionary case correctly and thus using the “wrongly built” validityOptionalDictionary but somehow it’s RequiredDictionaryDesmond Cheong
12/14/2024, 1:47 AMjay
12/14/2024, 1:52 AMoptional group containing a required column:jay
12/14/2024, 1:53 AMTabrez Mohammed
12/14/2024, 1:53 AMTabrez Mohammed
12/14/2024, 1:54 AMDesmond Cheong
12/14/2024, 1:55 AMsale_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 optionalTabrez Mohammed
12/14/2024, 1:56 AMjay
12/14/2024, 1:58 AMparquet meta <path_to_file> is how it worksTabrez Mohammed
12/14/2024, 2:00 AMoptional group sale_price = 10 {
optional int64 amount (DECIMAL(10,0)) = 37;
optional binary currency (STRING) = 38;
}Tabrez Mohammed
12/14/2024, 2:00 AMjay
12/14/2024, 2:01 AMI think that bit of metadata is saying thatThis makes sense, I’m guessing then that this meaning ofis optional, but if it’s defined thensale_priceis required. But I guess sincecurrencyis in an optional struct, it should also be optionalcurrency
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.jay
12/14/2024, 2:02 AMTabrez Mohammed
12/14/2024, 2:04 AMDesmond Cheong
12/14/2024, 2:29 AMI’m guessing then that this meaning ofYeah looks like it. This is what we deserialize as the nested field's metadatais getting picked up by us also when parsing the actual currency array.required
field_info: FieldInfo { name: "currency", repetition: Required, id: Some(38) .Desmond Cheong
12/14/2024, 2:30 AMIs 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
Desmond Cheong
12/14/2024, 2:51 AMjay
12/14/2024, 3:32 AMDesmond Cheong
12/18/2024, 7:05 AMSam Gottlieb
12/20/2024, 3:55 PMDesmond Cheong
12/20/2024, 3:55 PM