Hi everyone! I have been working on benchmark comp...
# general
s
Hi everyone! I have been working on benchmark comparing frameworks for distributed data processing in Python (such as Dask, PySpark, Modin, Bodo) using a sample workload on NYC Taxi data. I recently found out about Daft and thought it would interesting to include as well. I gave it my best effort at an implementation and so far the results look really good but I am looking for feedback on how I could improve my approach or any feedback about the benchmark in general. Particularly, I am curious about tips and best practices around using UDFs and choosing the right partitioning strategy. Note: I am a Bodo developer, so I am a bit biased towards our system I'll include links in the thread for more context 🙂. Thanks in advanced!
👍 1
For more context on the workload, it basically reads the NYC taxi dataset from S3, does some join to include weather data, applies a UDF, and does a groupby aggregate and sorts. Here is the top level readme for more information about the benchmark and other frameworks we tested (code included here as well). Here is my implementation in Daft + config I used for reference.
c
just took a quick look. I believe you could rewrite this as a native daft expression
Copy code
def get_time_bucket(t):
        bucket = "other"
        if t in (8, 9, 10):
            bucket = "morning"
        elif t in (11, 12, 13, 14, 15):
            bucket = "midday"
        elif t in (16, 17, 18):
            bucket = "afternoon"
        elif t in (19, 20, 21):
            bucket = "evening"
        return bucket
IMO, the easiest way would be to use a
daft.sql_expr
here.
Copy code
daft.sql_expr("""
CASE
    WHEN hour IN (8, 9, 10) THEN 'morning'
    WHEN hour IN (11, 12, 13, 14, 15) THEN 'midday'
    WHEN hour IN (16, 17, 18) THEN 'afternoon'
    WHEN hour IN (19, 20, 21) THEN 'evening'
    ELSE 'other'
END
""")
👍 1
s
Thanks for the fast reply! I wasn't aware of the
sql_expr
API but that's a great suggestion. Although part of the goal of this benchmark was to stay in Python as much as possible and utilize UDFs, this does make that part of the workload quite a bit faster, so it would be an interesting comparison.
c
you can also express this using the
if_else
expression if you wanted it to be pure python. The
sql_expr
provided above gets parsed into several
if_else
statements under the hood, so they are logically equivalent. I just think the syntax for sql is more intuitive for this specific kind of operation.
👍 1
r
In predicate with number ranges could also be a less than + greater than check which might be a minor optimization too.