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]

Reply via email to