Expand description
ArrowMetal inside Apache DataFusion: a physical optimizer rule and a custom ExecutionPlan.
ArrowMetalRule walks DataFusion’s optimized physical plan and replaces the nodes ArrowMetal
can run with the same answer by a MetalExec: SortExec (with or without fetch), a hash
AggregateExec over column keys with sum/min/max/count/avg or none (DISTINCT),
a FilterExec whose predicate translates, and a HashJoinExec (inner, left, right) on equal
int32, int64 or Utf8 keys. MetalExec collects its input’s partitions, runs the operation
through ArrowMetal’s plan runner on the GPU, and emits RecordBatches with the schema
DataFusion expects. Anything it does not support is left unchanged, and the reason is recorded
in the rule’s Report.
The default config (ArrowMetalConfig::default) takes full sorts from 250,000 input rows,
aggregates where the measured table (src/agg_table.rs) takes their shape (a replaced
aggregate estimates its group count from a sample of its keys when it runs and either runs on
ArrowMetal or hands the node back to DataFusion’s own operators, AggregateChoice), and
joins where the measured join table (src/join_table.rs) takes them (JoinChoice). Top-k
and filters are left.
use datafusion::prelude::*;
use datafusion_arrowmetal::{session_context, ArrowMetalConfig, ArrowMetalRule};
// The measured take-list; `ArrowMetalConfig::all()` takes every shape the rule can translate.
let rule = ArrowMetalRule::new(ArrowMetalConfig::default());
let ctx = session_context(SessionConfig::new(), rule.clone());
// register tables, run SQL ...
for d in rule.report().decisions() { println!("{d}"); }Structs§
- AggSpec
- One aggregate of a
MetalOp::Aggregate. - Arrow
Metal Config - When the rule takes a node.
- Arrow
Metal Rule - Replaces
SortExec(+ itsSortPreservingMergeExec), hashAggregateExec(aSinglenode, or aFinal/Partialpair) andFilterExecwithMetalExecwhen ArrowMetal gives the same answer and the input is big enough. Cloning shares the report. - Decision
- One node the rule looked at, one replaced aggregate’s run-time choice, or one runtime fallback.
- Group
Choice - What a replaced aggregate’s run-time choice was made from.
- Group
Estimate - One probe’s answer.
- Metal
Exec - Runs one
MetalOpon ArrowMetal. - Report
- What the rule decided for the last
report_plansplans since it was created or last cleared. - SortKey
- One sort key: an input column, a direction, and where its nulls go.
Enums§
- AggKind
- An aggregate function
MetalExeccomputes. - Aggregate
Choice - Who runs an aggregate the rule replaced.
- Join
Choice - Which hash joins the rule replaces (decided at plan time, from both inputs’ exact row counts).
- JoinHow
- Which rows of an equi-join come out (the
JoinTypes of DataFusion’sHashJoinExecthatMetalExecruns). - MetalOp
- What a
MetalExeccomputes, in terms of its input’s column indices.
Functions§
- physical_
optimizer_ rules - DataFusion’s default physical optimizer rules with
ruleinserted before the last two: the post-optimizationFilterPushdown(which wires dynamic filters to the operators that own them, so a node replaced after it would leave a dynamic filter nobody updates) andSanityCheckPlan(so the rewritten plan is still checked for distribution and ordering). - session_
context - A
SessionContextwith DataFusion’s defaults plusrule. - with_
arrowmetal - Registers
ruleon aSessionStateBuilder(seephysical_optimizer_rulesfor where).