Hi Team, I exectuted the below query against a pa...
# general
s
Hi Team, I exectuted the below query against a parquet file of 1 GB size, which contains 140 million rows and 14 columns in a windows machine having 64 GB RAM and 20 cores. DAFT(v0.4) takes 50 secs to execute whereas Duckdb takes only 5 secs. Is it the expected behaviour or I miss something on the setup config or need to tweak settings daft.context.set_runner_native() start_time = time.time() query = """ SELECT t3."LOT", t3."WAFER", t3."Measure_Time", ( (count() - count(t3.over_all_pat_status))/count() )* 100 as yield FROM ( SELECT t2."LOT", t2."WAFER", t2."Measure_Time", t2."DIE_X", t2."DIE_Y", Min(t2.pat_status) as over_all_pat_status FROM ( SELECT t1.*, t1.VALUE BETWEEN t1.LL AND t1.UL as pat_status FROM ( SELECT * FROM read_parquet('test.parquet') t JOIN ( SELECT PARAMETER , (avg(VALUE) - (3 * STDDEV(VALUE))) as LL, (avg(VALUE) + (3 * STDDEV(VALUE))) as UL FROM read_parquet('test.parquet') GROUP BY PARAMETER ) t0 ON t.PARAMETER = t0.PARAMETER )t1 )t2 GROUP BY "LOT", "WAFER", "Measure_Time", "DIE_X", "DIE_Y" )t3 GROUP BY "LOT", "WAFER", "Measure_Time" ORDER BY WAFER """ df =daft.sql(query).collect() end_time = time.time() print("Execution Time : ", (end_time - start_time), "sec")
d
Could you share the query plan from df.explain(show_all=True)? I see a couple of nested queries, so I want to confirm if we manage to unnest them
s
* Source: | Number of partitions = 1 | Output schema = LOT#Utf8, WAFER#Utf8, Measure_Time#Timestamp(Nanoseconds, None), | yield#Float64 However here is the logical plan used to produce this result: == Unoptimized Logical Plan == * Sort: Sort by = (col(WAFER), ascending, nulls last) | * Project: col(LOT), col(WAFER), col(Measure_Time), col(yield) | * Aggregation: [[count(col(LOT), All) - count(col(over_all_pat_status), Valid)] / | count(col(LOT), All)] * lit(100) as yield | Group by = col(LOT), col(WAFER), col(Measure_Time) | Output schema = LOT#Utf8, WAFER#Utf8, Measure_Time#Timestamp(Nanoseconds, None), | yield#Float64 | * Project: col(LOT), col(WAFER), col(Measure_Time), col(DIE_X), col(DIE_Y), | col(over_all_pat_status) | * Aggregation: min(col(pat_status)) as over_all_pat_status | Group by = col(LOT), col(WAFER), col(Measure_Time), col(DIE_X), col(DIE_Y) | Output schema = LOT#Utf8, WAFER#Utf8, Measure_Time#Timestamp(Nanoseconds, None), | DIE_X#Int32, DIE_Y#Int32, over_all_pat_status#Boolean | * Project: col(LOT), col(WAFER), col(PARAMETER), col(TEST_TYPE), | col(TEST_PROGRAM), col(MODULE), col(Measure_Time), col(Temperature), | col(DIE_X), col(DIE_Y), col(PRODUCT), col(UNIT), col(PASS_FAIL), col(VALUE), | col(t0.PARAMETER), col(LL), col(UL), col(VALUE) in [col(LL),col(UL)] as | pat_status | * Project: col(LOT), col(WAFER), col(PARAMETER), col(TEST_TYPE), | col(TEST_PROGRAM), col(MODULE), col(Measure_Time), col(Temperature), | col(DIE_X), col(DIE_Y), col(PRODUCT), col(UNIT), col(PASS_FAIL), col(VALUE), | col(t0.PARAMETER), col(LL), col(UL) | * Join: Type = Inner | Strategy = Auto | Left on = col(PARAMETER) | Right on = col(t0.PARAMETER) | Null equals Nulls = [false] | Output schema = LOT#Utf8, WAFER#Utf8, PARAMETER#Utf8, TEST_TYPE#Utf8, | TEST_PROGRAM#Utf8, MODULE#Utf8, Measure_Time#Timestamp(Nanoseconds, None), | Temperature#Float64, DIE_X#Int32, DIE_Y#Int32, PRODUCT#Utf8, UNIT#Utf8, | PASS_FAIL#Utf8, VALUE#Float64, t0.PARAMETER#Utf8, LL#Float64, UL#Float64 |\ | * Project: col(PARAMETER) as t0.PARAMETER, col(LL), col(UL) | | | * Project: col(PARAMETER), col(LL), col(UL) | | | * Aggregation: mean(col(VALUE)) - [lit(3) * stddev(col(VALUE))] as LL, | | mean(col(VALUE)) + [lit(3) * stddev(col(VALUE))] as UL | | Group by = col(PARAMETER) | | Output schema = PARAMETER#Utf8, LL#Float64, UL#Float64 | | | * GlobScanOperator | | Glob paths = [test.parquet] | | Coerce int96 timestamp unit = Nanoseconds | | Use multithreading = true | | File schema = LOT#Utf8, WAFER#Utf8, PARAMETER#Utf8, TEST_TYPE#Utf8, | | TEST_PROGRAM#Utf8, MODULE#Utf8, Measure_Time#Timestamp(Nanoseconds, None), | | Temperature#Float64, DIE_X#Int32, DIE_Y#Int32, PRODUCT#Utf8, UNIT#Utf8, | | PASS_FAIL#Utf8, VALUE#Float64 | | Partitioning keys = [] | | Output schema = LOT#Utf8, WAFER#Utf8, PARAMETER#Utf8, TEST_TYPE#Utf8, | | TEST_PROGRAM#Utf8, MODULE#Utf8, Measure_Time#Timestamp(Nanoseconds, None), | | Temperature#Float64, DIE_X#Int32, DIE_Y#Int32, PRODUCT#Utf8, UNIT#Utf8, | | PASS_FAIL#Utf8, VALUE#Float64 | * GlobScanOperator | Glob paths = [test.parquet] | Coerce int96 timestamp unit = Nanoseconds | Use multithreading = true | File schema = LOT#Utf8, WAFER#Utf8, PARAMETER#Utf8, TEST_TYPE#Utf8, | TEST_PROGRAM#Utf8, MODULE#Utf8, Measure_Time#Timestamp(Nanoseconds, None), | Temperature#Float64, DIE_X#Int32, DIE_Y#Int32, PRODUCT#Utf8, UNIT#Utf8, | PASS_FAIL#Utf8, VALUE#Float64 | Partitioning keys = [] | Output schema = LOT#Utf8, WAFER#Utf8, PARAMETER#Utf8, TEST_TYPE#Utf8, | TEST_PROGRAM#Utf8, MODULE#Utf8, Measure_Time#Timestamp(Nanoseconds, None), | Temperature#Float64, DIE_X#Int32, DIE_Y#Int32, PRODUCT#Utf8, UNIT#Utf8, | PASS_FAIL#Utf8, VALUE#Float64 == Optimized Logical Plan == * Sort: Sort by = (col(WAFER), ascending, nulls last) | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = 0, | Upper bound bytes = 452014383 } | * Project: col(LOT), col(WAFER), col(Measure_Time), [[col(LOT.local_count(All)) | - col(over_all_pat_status.local_count(Valid))] / col(LOT.local_count(All))] * | lit(100) as yield | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = 0, | Upper bound bytes = 452014383 } | * Aggregation: count(col(LOT), All) as LOT.local_count(All), | count(col(over_all_pat_status), Valid) as over_all_pat_status.local_count(Valid) | Group by = col(LOT), col(WAFER), col(Measure_Time) | Output schema = LOT#Utf8, WAFER#Utf8, Measure_Time#Timestamp(Nanoseconds, None), | LOT.local_count(All)#UInt64, over_all_pat_status.local_count(Valid)#UInt64 | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = 0, | Upper bound bytes = 452014383 } | * Project: col(LOT), col(over_all_pat_status), col(WAFER), col(Measure_Time) | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = 0, | Upper bound bytes = 452014383 } | * Aggregation: min(col(pat_status)) as over_all_pat_status | Group by = col(LOT), col(WAFER), col(Measure_Time), col(DIE_X), col(DIE_Y) | Output schema = LOT#Utf8, WAFER#Utf8, Measure_Time#Timestamp(Nanoseconds, None), | DIE_X#Int32, DIE_Y#Int32, over_all_pat_status#Boolean | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = 0, | Upper bound bytes = 452014383 } | * Project: col(VALUE) in [col(LL),col(UL)] as pat_status, col(LOT), col(WAFER), | col(Measure_Time), col(DIE_X), col(DIE_Y) | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = 0, | Upper bound bytes = 452014383 } | * Join: Type = Inner | Strategy = Auto | Left on = col(PARAMETER) | Right on = col(t0.PARAMETER) | Null equals Nulls = [false] | Output schema = PARAMETER#Utf8, LOT#Utf8, WAFER#Utf8, | Measure_Time#Timestamp(Nanoseconds, None), DIE_X#Int32, DIE_Y#Int32, | VALUE#Float64, t0.PARAMETER#Utf8, LL#Float64, UL#Float64 | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = 0, | Upper bound bytes = 452014383 } |\ | * Project: col(PARAMETER) as t0.PARAMETER, col(VALUE.local_mean()) - | | col((Literal(Int64(3)) * VALUE.local_stddev())) as LL, col(VALUE.local_mean()) | | + col((Literal(Int64(3)) * VALUE.local_stddev())) as UL | | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = | | 452014383, Upper bound bytes = 452014383 } | | | * Project: col(PARAMETER) as PARAMETER, col(VALUE.local_mean()) as | | VALUE.local_mean(), lit(3) * col(VALUE.local_stddev()) as (Literal(Int64(3)) * | | VALUE.local_stddev()) | | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = | | 452014383, Upper bound bytes = 452014383 } | | | * Aggregation: mean(col(VALUE)) as VALUE.local_mean(), stddev(col(VALUE)) as | | VALUE.local_stddev() | | Group by = col(PARAMETER) | | Output schema = PARAMETER#Utf8, VALUE.local_mean()#Float64, | | VALUE.local_stddev()#Float64 | | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = | | 452014383, Upper bound bytes = 452014383 } | | | * Project: col(VALUE), col(PARAMETER) | | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = | | 452014383, Upper bound bytes = 452014383 } | | | * Num Scan Tasks = 1 | | File schema = LOT#Utf8, WAFER#Utf8, PARAMETER#Utf8, TEST_TYPE#Utf8, | | TEST_PROGRAM#Utf8, MODULE#Utf8, Measure_Time#Timestamp(Nanoseconds, None), | | Temperature#Float64, DIE_X#Int32, DIE_Y#Int32, PRODUCT#Utf8, UNIT#Utf8, | | PASS_FAIL#Utf8, VALUE#Float64 | | Partitioning keys = [] | | Projection pushdown = [VALUE, PARAMETER] | | Output schema = PARAMETER#Utf8, VALUE#Float64 | | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = | | 452014383, Upper bound bytes = 452014383 } | * Project: col(PARAMETER), col(LOT), col(WAFER), col(Measure_Time), col(DIE_X), | col(DIE_Y), col(VALUE) | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = | 271608643, Upper bound bytes = 271608643 } | * Num Scan Tasks = 1 | File schema = LOT#Utf8, WAFER#Utf8, PARAMETER#Utf8, TEST_TYPE#Utf8, | TEST_PROGRAM#Utf8, MODULE#Utf8, Measure_Time#Timestamp(Nanoseconds, None), | Temperature#Float64, DIE_X#Int32, DIE_Y#Int32, PRODUCT#Utf8, UNIT#Utf8, | PASS_FAIL#Utf8, VALUE#Float64 | Partitioning keys = [] | Projection pushdown = [PARAMETER, LOT, WAFER, Measure_Time, DIE_X, DIE_Y, VALUE] | Filter pushdown = [not(is_null(col(PARAMETER))) & not(is_null(col(PARAMETER)))] | & not(is_null(col(PARAMETER))) | Output schema = LOT#Utf8, WAFER#Utf8, PARAMETER#Utf8, | Measure_Time#Timestamp(Nanoseconds, None), DIE_X#Int32, DIE_Y#Int32, | VALUE#Float64 | Stats = { Lower bound rows = 0, Upper bound rows = None, Lower bound bytes = | 271608643, Upper bound bytes = 271608643 }
There are some warnings, Warning: Column '(Literal(Int64(3)) * VALUE.local_stddev())' contains *, preventing potential wildcard match Warning: Column '(Literal(Int64(3)) * VALUE.local_stddev())' contains *, preventing potential wildcard match Warning: Column '(Literal(Int64(3)) * VALUE.local_stddev())' contains *, preventing potential wildcard match Warning: Column '(Literal(Int64(3)) * VALUE.local_stddev())' contains *, preventing potential wildcard match Warning: Column '(Literal(Int64(3)) * VALUE.local_stddev())' contains *, preventing potential wildcard match Warning: Column '(Literal(Int64(3)) * VALUE.local_stddev())' contains *, preventing potential wildcard match Warning: Column '(Literal(Int64(3)) * VALUE.local_stddev())' contains *, preventing potential wildcard match Warning: Column '(Literal(Int64(3)) * VALUE.local_stddev())' contains *, preventing potential wildcard match Warning: Column '(Literal(Int64(3)) * VALUE.local_stddev())' contains *, preventing potential wildcard match Warning: Column '(Literal(Int64(3)) * VALUE.local_stddev())' contains *, preventing potential wildcard match Warning: Column '(Literal(Int64(3)) * VALUE.local_stddev())' contains *, preventing potential wildcard match Warning: Column '(Literal(Int64(3)) * VALUE.local_stddev())' contains *, preventing potential wildcard match
👍 1
j
Also, we have a really nifty tool for analyzing query performance. Could you run your query with the explain analyze environment variable?
Copy code
DAFT_DEV_ENABLE_EXPLAIN_ANALYZE=1
s
Copy code
mermaid
flowchart BT
daft_local_execution::sinks::blocking_sink::BlockingSinkNode0["SortResult
rows received =  25
rows emitted =  25
CPU Time = 0.01ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode1["ProjectOperator
rows received =  25
rows emitted =  25
CPU Time = 0.01ms
"]
daft_local_execution::sinks::blocking_sink::BlockingSinkNode2["GroupedAggregateSink
rows received =  292,325
rows emitted =  25
CPU Time = 0.03ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode3["ProjectOperator
rows received =  292,325
rows emitted =  292,325
CPU Time = 0.01ms
"]
daft_local_execution::sinks::blocking_sink::BlockingSinkNode4["GroupedAggregateSink
rows received =  143,997,800
rows emitted =  292,325
CPU Time = 0.18ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode5["ProjectOperator
rows received =  143,997,800
rows emitted =  143,997,800
CPU Time = 0.01ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode6["InnerHashJoinProbeOperator
rows received =  496
rows emitted =  143,997,800
CPU Time = 0.01ms
"]
daft_local_execution::sinks::blocking_sink::BlockingSinkNode7["HashJoinBuildSink
rows received =  143,997,800
rows emitted =  0
CPU Time = 0.70ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode8["ProjectOperator
rows received =  143,997,800
rows emitted =  143,997,800
CPU Time = 1.14ms
"]
daft_local_execution::sources::source::SourceNode9["ScanTask

rows emitted =  143,997,800
bytes read = 655.43 MiB
"]
daft_local_execution::sources::source::SourceNode9 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode8
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode8 --> daft_local_execution::sinks::blocking_sink::BlockingSinkNode7
daft_local_execution::sinks::blocking_sink::BlockingSinkNode7 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode6
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode10["ProjectOperator
rows received =  496
rows emitted =  496
CPU Time = 0.01ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode11["ProjectOperator
rows received =  496
rows emitted =  496
CPU Time = 0.01ms
"]
daft_local_execution::sinks::blocking_sink::BlockingSinkNode12["GroupedAggregateSink
rows received =  143,997,800
rows emitted =  496
CPU Time = 0.13ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode13["ProjectOperator
rows received =  143,997,800
rows emitted =  143,997,800
CPU Time = 0.71ms
"]
daft_local_execution::sources::source::SourceNode14["ScanTask

rows emitted =  143,997,800
bytes read = 388.04 MiB
"]
daft_local_execution::sources::source::SourceNode14 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode13
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode13 --> daft_local_execution::sinks::blocking_sink::BlockingSinkNode12
daft_local_execution::sinks::blocking_sink::BlockingSinkNode12 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode11
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode11 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode10
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode10 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode6
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode6 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode5
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode5 --> daft_local_execution::sinks::blocking_sink::BlockingSinkNode4
daft_local_execution::sinks::blocking_sink::BlockingSinkNode4 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode3
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode3 --> daft_local_execution::sinks::blocking_sink::BlockingSinkNode2
daft_local_execution::sinks::blocking_sink::BlockingSinkNode2 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode1
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode1 --> daft_local_execution::sinks::blocking_sink::BlockingSinkNode0
j
Nice! Here’s a visualization after pasting it into
<https://mermaid.live>
s
Nice Visualization !!!
j
Interestingly the times seems to be really really fast for computation. I’m guessing most of the time here is being spent on the Scan of Parquet. Is this data sensitive? Would love to poke around the data and query if you’re open to sharing it!
s
Sure Jay, I will share it shortly
daft_test.parquet
c
The cpu times for explain analyze needs fixing, PR here. https://github.com/Eventual-Inc/Daft/pull/3511 . Regardless, it showed that the hash join build sink processed 143 M rows, and the probe was only 496 rows. This is probably the bottleneck
s
Also, I could observe that DuckDB takes up to 3 GB Memory whereas Daft takes 25 GB to execute the query. This is the peak memory consumption during query execution.
Out of curiosity, I tried this with Daft 0.3.15 using python runner, it takes 25 secs. DuckDB v1.1.3 => 5 secs Daft v0.3.15 - Python runner => 25 secs Daft v0.4 - Native runner => 50 secs
👍 1
c
Can confirm that switching the join side works. Runtime reduced from 35.2s -> 4.4s. Currently the join side is determined by comparing the size bytes for each side, using the row count instead leads to choosing the other side for the join. Will have a PR out for this soon
❤️ 1
s
Here are results using Daft 0.4.3, Could see improvement in timings, but not as you get
Daft v0.4.1 - Native runner => 50 secs Daft v0.4.3 - Native runner => 20 secs DuckDB v1.1.3 => 5 secs However, I could see very good improvement on memory usage, it has come down from 20 GB to 3GB
Copy code
mermaid
flowchart BT
daft_local_execution::sinks::blocking_sink::BlockingSinkNode0["Sort: Sort by = (col(WAFER), ascending, nulls last)
Stats = { Approx num rows = 2,023,113, Approx size bytes = 162.07 MiB,
Accumulated selectivity = 0.46 }
Rows received =  25
Rows emitted =  25
CPU Time = 2.81ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode1["Project: col(LOT), col(WAFER), col(Measure_Time), [[col(LOT.local_count(All))
- col(over_all_pat_status.local_count(Valid))] / col(LOT.local_count(All))] *
lit(100) as yield
Resource request = None
Stats = { Approx num rows = 2,023,113, Approx size bytes = 162.07 MiB,
Accumulated selectivity = 0.46 }

Rows received =  25
Rows emitted =  25
CPU Time = 1.32ms
"]
daft_local_execution::sinks::blocking_sink::BlockingSinkNode2["GroupedAggregate: count(col(LOT) as LOT.local_count(All), All),
count(col(over_all_pat_status) as over_all_pat_status.local_count(Valid), Valid)
Group by: col(LOT), col(WAFER), col(Measure_Time)
Stats = { Approx num rows = 2,023,113, Approx size bytes = 162.07 MiB,
Accumulated selectivity = 0.46 }
Rows received =  292,325
Rows emitted =  25
CPU Time = 152.09ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode3["Project: col(LOT), col(over_all_pat_status), col(WAFER), col(Measure_Time)
Resource request = None
Stats = { Approx num rows = 2,528,892, Approx size bytes = 202.59 MiB,
Accumulated selectivity = 0.58 }

Rows received =  292,325
Rows emitted =  292,325
CPU Time = 0.33ms
"]
daft_local_execution::sinks::blocking_sink::BlockingSinkNode4["GroupedAggregate: min(col(pat_status) as over_all_pat_status)
Group by: col(LOT), col(WAFER), col(Measure_Time), col(DIE_X), col(DIE_Y)
Stats = { Approx num rows = 2,528,892, Approx size bytes = 202.59 MiB,
Accumulated selectivity = 0.58 }
Rows received =  143,997,800
Rows emitted =  292,325
CPU Time = 73973.23ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode5["Project: [col(VALUE) <= col(UL)] & [col(VALUE) >= col(LL)] as pat_status,
col(LOT), col(WAFER), col(Measure_Time), col(DIE_X), col(DIE_Y)
Resource request = None
Stats = { Approx num rows = 3,161,116, Approx size bytes = 255.87 MiB,
Accumulated selectivity = 0.72 }

Rows received =  143,997,800
Rows emitted =  143,997,800
CPU Time = 1124.42ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode6["InnerHashJoinProbe:
Probe on: [col(PARAMETER)]
Build on left: false
Stats = { Approx num rows = 3,161,116, Approx size bytes = 255.87 MiB,
Accumulated selectivity = 0.72 }

Rows received =  143,997,800
Rows emitted =  143,997,800
CPU Time = 102378.30ms
"]
daft_local_execution::sinks::blocking_sink::BlockingSinkNode7["HashJoinBuild:
Track Indices: true
Key Schema: t0.PARAMETER#Utf8
Null equals Nulls = [false]
Stats = { Approx num rows = 3,512,351, Approx size bytes = 93.79 MiB,
Accumulated selectivity = 0.80 }
Rows received =  496
Rows emitted =  0
CPU Time = 0.30ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode8["Project: col(PARAMETER) as t0.PARAMETER, col(VALUE.local_mean()) -
col((Literal(Int64(3)) * VALUE.local_stddev())) as LL, col(VALUE.local_mean()) +
col((Literal(Int64(3)) * VALUE.local_stddev())) as UL
Resource request = None
Stats = { Approx num rows = 3,512,351, Approx size bytes = 93.79 MiB,
Accumulated selectivity = 0.80 }

Rows received =  496
Rows emitted =  496
CPU Time = 0.43ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode9["Project: col(PARAMETER) as PARAMETER, col(VALUE.local_mean()) as
VALUE.local_mean(), lit(3) * col(VALUE.local_stddev()) as (Literal(Int64(3)) *
VALUE.local_stddev())
Resource request = None
Stats = { Approx num rows = 3,512,351, Approx size bytes = 93.79 MiB,
Accumulated selectivity = 0.80 }

Rows received =  496
Rows emitted =  496
CPU Time = 0.49ms
"]
daft_local_execution::sinks::blocking_sink::BlockingSinkNode10["GroupedAggregate: mean(col(VALUE) as VALUE.local_mean()), stddev(col(VALUE) as
VALUE.local_stddev())
Group by: col(PARAMETER)
Stats = { Approx num rows = 3,512,351, Approx size bytes = 93.79 MiB,
Accumulated selectivity = 0.80 }
Rows received =  143,997,800
Rows emitted =  496
CPU Time = 24968.10ms
"]
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode11["Project: col(VALUE), col(PARAMETER)
Resource request = None
Stats = { Approx num rows = 4,390,439, Approx size bytes = 118.28 MiB,
Accumulated selectivity = 1.00 }

Rows received =  143,997,800
Rows emitted =  143,997,800
CPU Time = 69.99ms
"]
daft_local_execution::sources::source::SourceNode12["ScanTaskSource:
Num Scan Tasks = 1
Estimated Scan Bytes = 312818807
Pushdowns: {projection: [VALUE, PARAMETER]}
Schema: {LOT#Utf8, WAFER#Utf8, PARAMETER#Utf8, TEST_TYPE#Utf8,
TEST_PROGRAM#Utf8, MODULE#Utf8, Measure_Time#Timestamp(Microseconds, None),
Temperature#Float64, DIE_X#Int32, DIE_Y#Int32, PRODUCT#Utf8, UNIT#Utf8,
PASS_FAIL#Utf8, VALUE#Float64}
Scan Tasks: [
{File {<file://C>:/_Architecture/Data/daft_test.parquet}}
]
Stats = { Approx num rows = 4,390,439, Approx size bytes = 118.28 MiB,
Accumulated selectivity = 1.00 }

Rows emitted =  143,997,800
Bytes read = 127.97 MiB
"]
daft_local_execution::sources::source::SourceNode12 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode11
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode11 --> daft_local_execution::sinks::blocking_sink::BlockingSinkNode10
daft_local_execution::sinks::blocking_sink::BlockingSinkNode10 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode9
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode9 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode8
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode8 --> daft_local_execution::sinks::blocking_sink::BlockingSinkNode7
daft_local_execution::sinks::blocking_sink::BlockingSinkNode7 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode6
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode13["Project: col(PARAMETER), col(LOT), col(WAFER), col(Measure_Time), col(DIE_X),
col(DIE_Y), col(VALUE)
Resource request = None
Stats = { Approx num rows = 3,951,395, Approx size bytes = 319.84 MiB,
Accumulated selectivity = 0.90 }

Rows received =  143,997,800
Rows emitted =  143,997,800
CPU Time = 62.93ms
"]
daft_local_execution::sources::source::SourceNode14["ScanTaskSource:
Num Scan Tasks = 1
Estimated Scan Bytes = 312818807
Pushdowns: {projection: [PARAMETER, LOT, WAFER, Measure_Time, DIE_X, DIE_Y,
VALUE], filter: not(is_null(col(PARAMETER)))}
Schema: {LOT#Utf8, WAFER#Utf8, PARAMETER#Utf8, TEST_TYPE#Utf8,
TEST_PROGRAM#Utf8, MODULE#Utf8, Measure_Time#Timestamp(Microseconds, None),
Temperature#Float64, DIE_X#Int32, DIE_Y#Int32, PRODUCT#Utf8, UNIT#Utf8,
PASS_FAIL#Utf8, VALUE#Float64}
Scan Tasks: [
{File {<file://C>:/_Architecture/Data/daft_test.parquet}}
]
Stats = { Approx num rows = 3,951,395, Approx size bytes = 319.84 MiB,
Accumulated selectivity = 0.90 }

Rows emitted =  143,997,800
Bytes read = 246.69 MiB
"]
daft_local_execution::sources::source::SourceNode14 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode13
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode13 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode6
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode6 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode5
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode5 --> daft_local_execution::sinks::blocking_sink::BlockingSinkNode4
daft_local_execution::sinks::blocking_sink::BlockingSinkNode4 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode3
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode3 --> daft_local_execution::sinks::blocking_sink::BlockingSinkNode2
daft_local_execution::sinks::blocking_sink::BlockingSinkNode2 --> daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode1
daft_local_execution::intermediate_ops::intermediate_op::IntermediateNode1 --> daft_local_execution::sinks::blocking_sink::BlockingSinkNode0
@Colin Ho, @jay Do I need to tweak any settings to achieve 4s timings that you had mentioned? Please confirm
j
@Colin Ho could you chime in here on join-side selection?
c
I actually ran the same query on duckdb and daft on my laptop (12 core, 36gb macbook pro), and i get daft: 4.40s duckdb: 1.89s The ratio looks almost the same as yours, so I think it is probably just the difference in the hardware.
s
@Colin Ho I tried it on my laptop (20 core, 64GB, Windows 10, 12th Gen Intel(R) Core(TM) i9-12900H 2.50 GHz Processor), Daft 0.4.1.3 is four times slower than DuckDB. I had run it with default Daft settings
c
Yes, that's likely the best Daft can do as of today. Looking at the profile, the biggest time consumers are the grouped aggregate and hash join probe. We would need to optimize those kernels and make them faster
👍 1