Is it possible to do transform data in Daft? Like ...
# general
k
Is it possible to do transform data in Daft? Like Ray
map_batches
does transforming the data frame.
r
You can use
apply
to use a custom transformation on a single column
k
Will the UDF apply in batches. I have large data set and want to really parallelize the processing across all ray nodes.
Also the output of UDF has to be single column? Can I return another dataframe out of it?
c
Yes it will be single column, but one way to get around is via structs. Example:
Copy code
@daft.udf(
    return_dtype=daft.DataType.struct(
        {
            "x": daft.DataType.int64(),
            "y": daft.DataType.int64(),
        }
    )
)
def my_udf(a, b):
    # simple UDF that just returns the two inputs as a struct column
    result = []
    for a_elem, b_elem in zip(a.to_pylist(), b.to_pylist()):
        result.append({"x": a_elem, "y": b_elem})
    return result

df = daft.from_pydict({"a": [1, 2, 3], "b": [4, 5, 6]})
# call UDF
df = df.select(my_udf(df["a"], df["b"]).alias("udf_result"))
# unnest struct fields
df = df.select("udf_result.*")
df.show()
Though not the most ergonomic, it should work See https://github.com/Eventual-Inc/Daft/issues/3567#issuecomment-2542203904 for more detail
And yes, the UDF is applied in parallel across batches in the ray cluster