Hi, I'm new to Daft. It appears to be fast, but I'...
# general
d
Hi, I'm new to Daft. It appears to be fast, but I'm running into a blocker with the datatype FixedSizeList[float64] . I have a series of parquet files that have an
image_id
column and a
predictions
column in s3. The
image_id
column is of type
Utf8
and the
predictions
column is of type
FixedSizeList[Float64; 1000]
I want to group by
image_id
and take the mean value of item
150
in
predictions
column.
Copy code
╭────────────────────────────────┬────────────────────────────────╮                                                                                                                                         
│ image_id                       ┆ predictions                    │                                                                                                                                         
│ ---                            ┆ ---                            │
│ Utf8                           ┆ FixedSizeList[Float64; 1000]   │
╞════════════════════════════════╪════════════════════════════════╡
I'd like to explicitly access a single element in the
FixedSizeList
during aggregation, but I cannot find any documentation of the methods of a column full of datatypes
FixedSizeList
. The only way that I have found to select an element is to first convert the
FixedSizedList
to a Python
list
and then to a
list
datatype followed by using the
get
method. This blows up the memory of the pipeline. Ideally, I would like to call a method directly on the column object that selects a single element from
FixedSizeList.
Does anyone have any tips on how to do this? Working attempt, which blows up memory
Copy code
df3 = (
    df2.select(col("image_id"), col("predictions"))
    .with_column(
        "cool_channel",
        col("predictions").cast(daft.DataType.python()).cast(daft.DataType.list(daft.DataType.float64())).list.get(150)
    )
    .groupby('image_id').agg(col("cool_channel").mean().alias("mean_channel"))
)
Ideal solution
Copy code
df3 = (
    df2.select(col("image_id"), col("predictions"))
    .with_column(
        "cool_channel",
        col("predictions").pluck(150)  # method does not exist
    )
    .groupby('image_id').agg(col("cool_channel").mean().alias("mean_channel"))
)
Can someone point me to some more documentation of
FixedSizeLists
or any solution to this issue?
j
Welcome @Dexter Antonio! Have you tried the
.list.get
expression?
👀 1
col("predictions").list.get(150)
image.png
d
Copy code
>>> df3 = (
...     df2.select(col("image_id"), col("predictions"))
...     .with_column(
...         "cool_channel",
...         col("predictions").list.get(150)
...     )
...     .groupby('image_id').agg(col("cool_channel").mean().alias("mean_channel"))
... ).show(1)
Traceback (most recent call last):
  File "<stdin>", line 3, in <module>
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/api_annotations.py", line 26, in _wrap
    return timed_method(*args, **kwargs)
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/analytics.py", line 199, in tracked_method
    result = method(*args, **kwargs)
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/dataframe/dataframe.py", line 1473, in with_column
    return self.with_columns({column_name: expr})
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/api_annotations.py", line 26, in _wrap
    return timed_method(*args, **kwargs)
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/analytics.py", line 199, in tracked_method
    result = method(*args, **kwargs)
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/dataframe/dataframe.py", line 1520, in with_columns
    builder = self._builder.with_columns(new_columns)
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/logical/builder.py", line 159, in with_columns
    builder = self._builder.with_columns(column_pyexprs)
daft.exceptions.DaftCoreException: DaftError::External Unable to create logical plan node.
Due to: DaftError::ValueError Column "predictions" with dtype FixedSizeList[Float64; 1000] cannot be exploded, must be a List or FixedSizeList column.
This is the error that I get when I attempt that approach
Copy code
Due to: DaftError::ValueError Column "predictions" with dtype FixedSizeList[Float64; 1000] cannot be exploded, must be a List or FixedSizeList column.
I guess the issue isn't with the
.list.get
methods but with something else 🤔.
j
Very odd… What version of Daft are you using?
d
Copy code
$ uv pip freeze | grep getdaft
getdaft==0.3.14
j
Super weird, it’s working fine for me and I’m looking at the code and see no possible way why this is failing 😓 the error message suggests the column is indeed a
FixedSizeList[Float64; 1000]
type, but the code should never match into that error statement. Could you try running this code for me? It runs perfectly on my machine, but I wonder if your installation of Daft would error out:
Copy code
import daft

df = daft.from_pydict({"foo": [[0.1, 0.2], [0.1, 0.2]]})
df = df.with_column(
    "foo",
    df["foo"].cast(daft.DataType.fixed_size_list(daft.DataType.float64(), 2)),
)
df = df.with_column("x", df["foo"].list.get(0))
df.collect()
👀 1
d
Your code works great
Copy code
>>> import daft
>>> 
>>> df = daft.from_pydict({"foo": [[0.1, 0.2], [0.1, 0.2]]})
",
    df["foo"].cast(daft.DataType.fixed_size_list(daft.DataType.float64(), 2)),
)
df = df.with_column("x", df["foo"].list.get(0))
df.collect()
>>> df = df.with_column(
...     "foo",
...     df["foo"].cast(daft.DataType.fixed_size_list(daft.DataType.float64(), 2)),
... )
>>> df = df.with_column("x", df["foo"].list.get(0))
>>> df.collect()
╭───────────────────────────┬─────────╮                                                                                                                                                                     
│ foo                       ┆ x       │
│ ---                       ┆ ---     │
│ FixedSizeList[Float64; 2] ┆ Float64 │
╞═══════════════════════════╪═════════╡
│ [0.1, 0.2]                ┆ 0.1     │
├╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌┤
│ [0.1, 0.2]                ┆ 0.1     │
╰───────────────────────────┴─────────╯

(Showing first 2 of 2 rows)
>>>
My guess is that my columns are not actually ``FixedSizeList[Float64; 1000]` or a more accurate representation of them is ``FixedSizeList[Float64 | None; 1000]` . What I can try doing is to pick a small subset of my data and see if the aggregation works on there. Maybe there are some strange values that are screwing things up.
j
Interesting… any chance you could share the Parquet files so we can attempt to reproduce this on our end?
My hypothesis is that these Parquet files might be producing a special type that looks like a fixed size list but is actually some kind of extension type that wraps it
d
Currently I'm scanning through hundreds of gigabytes of data. I will send you a smaller subset of my data that reproduces the error and masks some of the proprietary information.
j
That would be incredibly helpful, thank you! And yes, we probably just need access to a small subset with just the
predictions
column
d
great, I will get back to you
j
Yes ok I’m pretty sure we’re getting some extension types from your Parquet file and aren’t handling it properly. Made a PR to fix a display bug for our extension types which should at least get it to display correctly in error messages: https://github.com/Eventual-Inc/Daft/pull/3456 We haven’t really focused on handling external extension types in Daft, but can probably add some functionality around it once we can verify that this is indeed the case. I’m guessing that some external system produced these Parquet files and gave the field this special type!
d
Here is a minimal example
Copy code
import daft
from daft import col
df2 = daft.read_parquet("test3.parquet")
df3 = (
    df2.select(col("image_id"), col("predictions"))
    .with_column(
        "cool_channel",
        col("predictions").list.get(150)
    )
    .groupby('image_id').agg(col("cool_channel").mean().alias("mean_channel"))
).show(1)
Copy code
Traceback (most recent call last):
  File "<stdin>", line 3, in <module>
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/api_annotations.py", line 26, in _wrap
    return timed_method(*args, **kwargs)
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/analytics.py", line 199, in tracked_method
    result = method(*args, **kwargs)
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/dataframe/dataframe.py", line 1473, in with_column
    return self.with_columns({column_name: expr})
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/api_annotations.py", line 26, in _wrap
    return timed_method(*args, **kwargs)
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/analytics.py", line 199, in tracked_method
    result = method(*args, **kwargs)
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/dataframe/dataframe.py", line 1520, in with_columns
    builder = self._builder.with_columns(new_columns)
  File "/root/workspaces/noetik-analysis/.venv/lib/python3.10/site-packages/daft/logical/builder.py", line 159, in with_columns
    builder = self._builder.with_columns(column_pyexprs)
daft.exceptions.DaftCoreException: DaftError::External Unable to create logical plan node.
Due to: DaftError::ValueError Column "predictions" with dtype FixedSizeList[Float64; 1000] cannot be exploded, must be a List or FixedSizeList column.
If I create
test.parquet
with pandas, then the error is resolved. The problem could lie in how daft is writing parquet column metadata.
Thank you for creating a PR
j
Was test3.parquet written by Daft?
d
yes
👍 1
j
Yeah so after reading your data with my changes, I see that it’s actually an extension type:
We’ll take a look at what’s happening here… Thanks for raising this! cc @Colin Ho as well
d
great, thank you for addressing this so quickly. Is extension type an extension to the parquet file format and somehow daft is writing out the wrong extension type into the metadata of the parquet?
j
Yup! I think there’s some bad interactions between Daft and PyArrow that makes it write bad metadata into parquet. How did you write this data by the way? Was the datatype of the column a fixed shape tensor?
👍 1
Also curious about your version of PyArrow and Ray (if applicable)
d
Another member of my team created these files. Let me ask him and get back to you. I believe he was using Ray in part of the pipeline.