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]
