andygrove commented on issue #1430:
URL: 
https://github.com/apache/datafusion-comet/issues/1430#issuecomment-5872382366

   Looking at this again, I don't think either suggestion would change 
anything. With AQE, Spark re-plans after each stage finishes and the re-plan 
runs `RewriteJoin` again, so for a join between two exchanges we already 
compare the measured size of each shuffle's output, with its exact row count as 
the tie-break. That output is exactly the set of columns that goes into the 
hash table, so estimating it from row count times column sizes would be less 
accurate than what we have now. Where there are no runtime stats, Spark's 
estimate already works the way the second suggestion describes: 
`LogicalRelation` stats come from `HadoopFsRelation.sizeInBytes`, and each 
projection scales that by the size of the columns it keeps. With CBO and table 
stats, Spark computes row count times row size directly. The rewrite is still 
opt-in via `spark.comet.exec.forceShuffledHashJoin`, and nobody has reported a 
bad build side since #1424. I think we can close this, and open a new issue 
with the plan and j
 oin metrics if we find a query where we pick the wrong side.


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to