Skip to main content

Crate datafusion_arrowmetal

Crate datafusion_arrowmetal 

Source
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.
ArrowMetalConfig
When the rule takes a node.
ArrowMetalRule
Replaces SortExec (+ its SortPreservingMergeExec), hash AggregateExec (a Single node, or a Final/Partial pair) and FilterExec with MetalExec when 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.
GroupChoice
What a replaced aggregate’s run-time choice was made from.
GroupEstimate
One probe’s answer.
MetalExec
Runs one MetalOp on ArrowMetal.
Report
What the rule decided for the last report_plans plans 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 MetalExec computes.
AggregateChoice
Who runs an aggregate the rule replaced.
JoinChoice
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’s HashJoinExec that MetalExec runs).
MetalOp
What a MetalExec computes, in terms of its input’s column indices.

Functions§

physical_optimizer_rules
DataFusion’s default physical optimizer rules with rule inserted before the last two: the post-optimization FilterPushdown (which wires dynamic filters to the operators that own them, so a node replaced after it would leave a dynamic filter nobody updates) and SanityCheckPlan (so the rewritten plan is still checked for distribution and ordering).
session_context
A SessionContext with DataFusion’s defaults plus rule.
with_arrowmetal
Registers rule on a SessionStateBuilder (see physical_optimizer_rules for where).