Uploaded image for project: 'Apache Arrow'
  1. Apache Arrow
  2. ARROW-9464

[Rust] [DataFusion] Physical plan refactor to support optimization rules and more efficient use of threads

    XMLWordPrintableJSON

    Details

      Description

      I would like to propose a refactor of the physical/execution planning based on the experience I have had in implementing distributed execution in Ballista.

      This will likely need subtasks but here is an overview of the changes I am proposing.

      Introduce physical plan optimization rule to insert "shuffle" operators

      We should extend the ExecutionPlan trait so that each operator can specify its input and output partitioning needs, and then have an optimization rule that can insert any repartitioning or reordering steps required.

      For example, these are the methods to be added to ExecutionPlan. This design is based on Apache Spark.

       

      /// Specifies how data is partitioned across different nodes in the cluster
      fn output_partitioning(&self) -> Partitioning {
          Partitioning::UnknownPartitioning(0)
      }
      
      /// Specifies the data distribution requirements of all the children for this operator
      fn required_child_distribution(&self) -> Distribution {
          Distribution::UnspecifiedDistribution
      }
      
      /// Specifies how data is ordered in each partition
      fn output_ordering(&self) -> Option<Vec<SortOrder>> {
          None
      }
      
      /// Specifies the data distribution requirements of all the children for this operator
      fn required_child_ordering(&self) -> Option<Vec<Vec<SortOrder>>> {
          None
      }
       

      A good example of applying this rule would be in the case of hash aggregates where we perform a partial aggregate in parallel across partitions and then coalesce the results and apply a final hash aggregate.

      Another example would be a SortMergeExec specifying the sort order required for its children.

       

       

        Attachments

          Issue Links

            Activity

              People

              • Assignee:
                andygrove Andy Grove
                Reporter:
                andygrove Andy Grove
              • Votes:
                0 Vote for this issue
                Watchers:
                4 Start watching this issue

                Dates

                • Created:
                  Updated:
                  Resolved:

                  Time Tracking

                  Estimated:
                  Original Estimate - Not Specified
                  Not Specified
                  Remaining:
                  Remaining Estimate - 0h
                  0h
                  Logged:
                  Time Spent - 9.5h
                  9.5h