andygrove commented on issue #2520: URL: https://github.com/apache/datafusion-ballista/issues/2520#issuecomment-6062385556
I checked A3 against `main` and the picture is smaller than the table suggests. - The stage plan is not shipped whole to every task. `restrict_plan_to_partitions` slices scan file groups and shuffle reader locations first, so the bytes summed across a stage's tasks are about 1x the plan. Only broadcast readers, which stay intact by design, are repeated per task. - I timed restrict plus encode in a debug build on a synthetic filter over a scan or reader, for a 256-partition stage. A scan with about 10k files took 378 ms for 256 tasks and 57 ms for 32 tasks. A reader with 256 partitions and 34 locations each took 56 ms and 41 ms. Release builds should be several times faster. - The one superlinear piece is that each restrict clones the whole `FileScanConfig`, so cost is tasks times files. It only matters at much larger file counts than this. - The per-stage cache was removed by #2038 because plans now differ per task. #619 is the same idea and is already done. I'd drop A3 from the High list. If a plan does get too big, the broadcast reader's location list could be a cause, which may be related to #2290 and #2165 but I haven't confirmed that. -- 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]
