Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
444 changes: 244 additions & 200 deletions Cargo.lock

Large diffs are not rendered by default.

18 changes: 9 additions & 9 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,8 @@ members = [
]

[workspace.dependencies]
datafusion = { version = "54.0.0", default-features = false }
datafusion-proto = { version = "54.0.0", default-features = false }
datafusion = { version = "55", default-features = false }
datafusion-proto = { version = "55", default-features = false }

[package]
name = "datafusion-distributed"
Expand Down Expand Up @@ -50,16 +50,16 @@ bincode = "1"
tonic-prost = "0.14.2"

# grpc-specific features
arrow-flight = { version = "58", optional = true }
arrow-select = { version = "58", optional = true }
arrow-ipc = { version = "58", features = ["zstd"], optional = true }
arrow-flight = { version = "59", optional = true }
arrow-select = { version = "59", optional = true }
arrow-ipc = { version = "59", features = ["zstd"], optional = true }
tonic = { version = "0.14.1", features = ["transport"], optional = true }
tower = { version = "0.5.2", optional = true }

# integration_tests deps
insta = { version = "1.46.0", features = ["filters"], optional = true }
parquet = { version = "58", optional = true }
arrow = { version = "58", optional = true, features = ["test_utils"] }
parquet = { version = "59", optional = true }
arrow = { version = "59", optional = true, features = ["test_utils"] }
hyper-util = { version = "0.1.16", optional = true }

[features]
Expand Down Expand Up @@ -98,8 +98,8 @@ datafusion = { workspace = true, features = ["parquet", "sql"] }
datafusion-proto = { workspace = true, features = ["parquet"] }
structopt = "0.3"
insta = { version = "1.46.0", features = ["filters"] }
parquet = "58"
arrow = { version = "58", features = ["test_utils"] }
parquet = "59"
arrow = { version = "59", features = ["test_utils"] }
tokio-stream = { version = "0.1.17", features = ["sync"] }
hyper-util = "0.1.16"
pretty_assertions = "1.4"
Expand Down
8 changes: 4 additions & 4 deletions benchmarks/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -14,17 +14,17 @@ datafusion-distributed = { path = "..", features = [
] }
tokio = { version = "1.48", features = ["full"] }
url = "2.5.7"
parquet = { version = "58" }
parquet = { version = "59" }
structopt = { version = "0.3.26" }
log = "0.4.27"
serde = { version = "1.0.219", features = ["derive"] }
serde_json = "1.0.141"
env_logger = "0.11.8"
futures = "0.3.31"
arrow = { version = "58", features = ["test_utils"] }
arrow = { version = "59", features = ["test_utils"] }
tonic = { version = "0.14.1", features = ["transport"] }
tpchgen = { git = "https://github.com/clflushopt/tpchgen-rs", rev = "438e9c2dbc25b2fff82c0efc08b3f13b5707874f" }
tpchgen-arrow = { git = "https://github.com/clflushopt/tpchgen-rs", rev = "438e9c2dbc25b2fff82c0efc08b3f13b5707874f" }
tpchgen = { git = "https://github.com/clflushopt/tpchgen-rs", rev = "f42708c32e7fd816b8ad8d2a880c1802a71b669f" }
tpchgen-arrow = { git = "https://github.com/clflushopt/tpchgen-rs", rev = "f42708c32e7fd816b8ad8d2a880c1802a71b669f" }
reqwest = "0.12"
zip = "6.0"
sketches-ddsketch = "0.3"
Expand Down
10 changes: 9 additions & 1 deletion benchmarks/benches/broadcast_cache_scenarios.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,9 +3,10 @@ use datafusion::arrow::array::UInt8Array;
use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use datafusion::arrow::record_batch::RecordBatch;
use datafusion::common::Statistics;
use datafusion::common::tree_node::TreeNodeRecursion;
use datafusion::error::Result;
use datafusion::execution::{SendableRecordBatchStream, TaskContext};
use datafusion::physical_expr::EquivalenceProperties;
use datafusion::physical_expr::{EquivalenceProperties, PhysicalExpr};
use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
use datafusion::physical_plan::{
Expand Down Expand Up @@ -133,6 +134,13 @@ impl ExecutionPlan for SyntheticExec {
vec![]
}

fn apply_expressions(
&self,
_f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> Result<TreeNodeRecursion>,
) -> Result<TreeNodeRecursion> {
Ok(TreeNodeRecursion::Continue)
}

fn with_new_children(
self: Arc<Self>,
_children: Vec<Arc<dyn ExecutionPlan>>,
Expand Down
2 changes: 1 addition & 1 deletion cli/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ edition = "2024"
[dependencies]
datafusion = { workspace = true }
datafusion-distributed = { path = "..", features = ["avro", "integration"] }
datafusion-cli = { version = "54", default-features = false }
datafusion-cli = { version = "55", default-features = false }
tokio = { version = "1.48", features = ["full"] }
clap = { version = "4", features = ["derive"] }
env_logger = "0.11"
Expand Down
4 changes: 2 additions & 2 deletions cli/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@
// File mainly copied from https://github.com/apache/datafusion/blob/main/datafusion-cli/src/main.rs

use clap::Parser;
use datafusion::common::config_err;
use datafusion::common::{config::ConfigNonZeroUsize, config_err};
use datafusion::config::ConfigOptions;
use datafusion::error::{DataFusionError, Result};
use datafusion::execution::SessionStateBuilder;
Expand Down Expand Up @@ -214,7 +214,7 @@ fn get_session_config(args: &Args) -> Result<SessionConfig> {
if batch_size == 0 {
return config_err!("batch_size must be greater than 0");
}
config_options.execution.batch_size = batch_size;
config_options.execution.batch_size = ConfigNonZeroUsize::try_new(batch_size)?;
};

// use easier to understand "tree" mode by default
Expand Down
2 changes: 1 addition & 1 deletion console/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ url = "2.5.7"
tokio-stream = "0.1.18"

[dev-dependencies]
arrow = "58"
arrow = "59"
async-trait = "0.1.89"
datafusion-distributed-benchmarks = { path = "../benchmarks" }
futures = "0.3.31"
Expand Down
8 changes: 4 additions & 4 deletions console/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -63,10 +63,10 @@ async fn run_app(
terminal.draw(|frame| ui::render(frame, app))?;

// Check for keyboard input (16ms timeout ~ 60fps responsiveness)
if event::poll(Duration::from_millis(16))? {
Comment thread
gabotechs marked this conversation as resolved.
if let Event::Key(key) = event::read()? {
input::handle_key_event(app, key);
}
if event::poll(Duration::from_millis(16))?
&& let Event::Key(key) = event::read()?
{
input::handle_key_event(app, key);
}

if app.should_quit {
Expand Down
50 changes: 25 additions & 25 deletions console/src/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -176,32 +176,32 @@ impl WorkerConn {

// Detect completed tasks: tasks that were running but disappeared
for old_task in &self.tasks {
if old_task.status == TaskStatus::Running as i32 {
Comment thread
gabotechs marked this conversation as resolved.
if let Some(sk) = &old_task.task_key {
let key = (sk.query_id.clone(), sk.stage_id, sk.task_number);
if !new_task_keys.contains(&key) {
// Task disappeared — assume completed
let observed_duration = self
.task_first_seen
.get(&key)
.map(|first| first.elapsed())
.unwrap_or_default();

self.completed_tasks.push_front(CompletedTaskRecord {
query_id: sk.query_id.clone(),
stage_id: sk.stage_id,
task_number: sk.task_number,
observed_duration,
});

// Maintain bounded size
while self.completed_tasks.len() > MAX_COMPLETED_TASKS {
self.completed_tasks.pop_back();
}

// Remove from first_seen tracking
self.task_first_seen.remove(&key);
if old_task.status == TaskStatus::Running as i32
&& let Some(sk) = &old_task.task_key
{
let key = (sk.query_id.clone(), sk.stage_id, sk.task_number);
if !new_task_keys.contains(&key) {
// Task disappeared — assume completed
let observed_duration = self
.task_first_seen
.get(&key)
.map(|first| first.elapsed())
.unwrap_or_default();

self.completed_tasks.push_front(CompletedTaskRecord {
query_id: sk.query_id.clone(),
stage_id: sk.stage_id,
task_number: sk.task_number,
observed_duration,
});

// Maintain bounded size
while self.completed_tasks.len() > MAX_COMPLETED_TASKS {
self.completed_tasks.pop_back();
}

// Remove from first_seen tracking
self.task_first_seen.remove(&key);
}
}
}
Expand Down
26 changes: 26 additions & 0 deletions docs/upgrade/4.0.0.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
# Upgrading from 3.0.0 to 4.0.0

This guide covers the source and behavioural changes merged after the
`v3.0.0` tag.

## 1. Introduction of PhysicalProtoConverterExtension to the DistributedCodec

Custom `PhysicalExtensionCodec` implementations must use the new
`&dyn PhysicalProtoConverterExtension` argument in both `try_encode` and
`try_decode` when encoding or decoding DataFusion
plans and expressions. See
[Distribute a custom execution plan](../source/user-guide/04-distribute-custom-plan.md)
for the complete signatures.

## 2. `NetworkBoundary::producer_head` now returns `Result`

Custom `NetworkBoundary` implementations may now return errors when
constructing producer heads.

```rust
fn producer_head(&self, consumer_tasks: usize) -> Result<ProducerHead> {
Ok(ProducerHead::BroadcastExec {
output_partitions: self.output_partitions * consumer_tasks,
})
}
```
20 changes: 17 additions & 3 deletions examples/custom_execution_plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,14 +23,15 @@ use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use datafusion::arrow::record_batch::RecordBatchOptions;
use datafusion::arrow::util::pretty::pretty_format_batches;
use datafusion::catalog::{Session, TableFunctionImpl};
use datafusion::common::tree_node::TreeNodeRecursion;
use datafusion::common::{
DataFusionError, Result, ScalarValue, exec_err, extensions_options, internal_err, plan_err,
};
use datafusion::config::ConfigExtension;
use datafusion::datasource::{TableProvider, TableType};
use datafusion::execution::{SendableRecordBatchStream, SessionStateBuilder, TaskContext};
use datafusion::logical_expr::Expr;
use datafusion::physical_expr::EquivalenceProperties;
use datafusion::physical_expr::{EquivalenceProperties, PhysicalExpr};
use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
use datafusion::physical_plan::{DisplayAs, DisplayFormatType, ExecutionPlan, PlanProperties};
Expand All @@ -43,7 +44,7 @@ use datafusion_distributed::{
ScaleUpLeafNodeEvent, ScaleUpLeafNodeEventResponse, SessionStateBuilderExt, WorkerQueryContext,
display_plan_ascii,
};
use datafusion_proto::physical_plan::PhysicalExtensionCodec;
use datafusion_proto::physical_plan::{PhysicalExtensionCodec, PhysicalProtoConverterExtension};
use datafusion_proto::protobuf;
use datafusion_proto::protobuf::proto_error;
use futures::{TryStreamExt, stream};
Expand Down Expand Up @@ -172,6 +173,13 @@ impl ExecutionPlan for NumbersExec {
vec![]
}

fn apply_expressions(
&self,
_f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> Result<TreeNodeRecursion>,
) -> Result<TreeNodeRecursion> {
Ok(TreeNodeRecursion::Continue)
}

fn with_new_children(
self: Arc<Self>,
_: Vec<Arc<dyn ExecutionPlan>>,
Expand Down Expand Up @@ -242,6 +250,7 @@ impl PhysicalExtensionCodec for NumbersExecCodec {
buf: &[u8],
inputs: &[Arc<dyn ExecutionPlan>],
_ctx: &TaskContext,
_proto_converter: &dyn PhysicalProtoConverterExtension,
) -> Result<Arc<dyn ExecutionPlan>> {
if !inputs.is_empty() {
return internal_err!("NumbersExec should have no children, got {}", inputs.len());
Expand All @@ -262,7 +271,12 @@ impl PhysicalExtensionCodec for NumbersExecCodec {
)))
}

fn try_encode(&self, node: Arc<dyn ExecutionPlan>, buf: &mut Vec<u8>) -> Result<()> {
fn try_encode(
&self,
node: Arc<dyn ExecutionPlan>,
buf: &mut Vec<u8>,
_proto_converter: &dyn PhysicalProtoConverterExtension,
) -> Result<()> {
let Some(exec) = node.downcast_ref::<NumbersExec>() else {
return internal_err!("Expected plan to be NumbersExec, but was {}", node.name());
};
Expand Down
20 changes: 17 additions & 3 deletions examples/custom_worker_url_routing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ use datafusion::common::{Result, internal_err};
use datafusion::config::ConfigOptions;
use datafusion::datasource::physical_plan::{FileGroup, FileScanConfig};
use datafusion::execution::{SendableRecordBatchStream, SessionStateBuilder, TaskContext};
use datafusion::physical_expr::PhysicalExpr;
use datafusion::physical_optimizer::PhysicalOptimizerRule;
use datafusion::physical_plan::stream::{
RecordBatchReceiverStreamBuilder, RecordBatchStreamAdapter,
Expand All @@ -45,7 +46,7 @@ use datafusion_distributed::{
DistributedLeafExec, RouteTasksEvent, RouteTasksEventResponse, ScaleUpLeafNodeEvent,
ScaleUpLeafNodeEventResponse, SessionStateBuilderExt, WorkerQueryContext, display_plan_ascii,
};
use datafusion_proto::physical_plan::PhysicalExtensionCodec;
use datafusion_proto::physical_plan::{PhysicalExtensionCodec, PhysicalProtoConverterExtension};
use datafusion_proto::protobuf;
use futures::TryStreamExt;
use prost::Message;
Expand Down Expand Up @@ -93,6 +94,13 @@ impl ExecutionPlan for CacheExec {
vec![&self.child]
}

fn apply_expressions(
&self,
_f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> Result<TreeNodeRecursion>,
) -> Result<TreeNodeRecursion> {
Ok(TreeNodeRecursion::Continue)
}

fn with_new_children(
self: Arc<Self>,
mut children: Vec<Arc<dyn ExecutionPlan>>,
Expand Down Expand Up @@ -151,7 +159,7 @@ impl ExecutionPlan for CacheExec {
fn hash_key(file_group: &FileGroup) -> usize {
let mut hasher = DefaultHasher::new();
for file in file_group.files() {
let serialized: protobuf::PartitionedFile = file.try_into().unwrap();
let serialized = protobuf::PartitionedFile::try_from(file).unwrap();
hasher.write(&serialized.encode_to_vec());
}
hasher.finish() as usize
Expand Down Expand Up @@ -247,14 +255,20 @@ impl PhysicalExtensionCodec for CachedFileScanCodec {
_buf: &[u8],
inputs: &[Arc<dyn ExecutionPlan>],
_ctx: &TaskContext,
_proto_converter: &dyn PhysicalProtoConverterExtension,
) -> Result<Arc<dyn ExecutionPlan>> {
let [child] = inputs else {
return internal_err!("CacheExec expects exactly 1 child, got {}", inputs.len());
};
Ok(CacheExec::new(Arc::clone(child)))
}

fn try_encode(&self, node: Arc<dyn ExecutionPlan>, _buf: &mut Vec<u8>) -> Result<()> {
fn try_encode(
&self,
node: Arc<dyn ExecutionPlan>,
_buf: &mut Vec<u8>,
_proto_converter: &dyn PhysicalProtoConverterExtension,
) -> Result<()> {
if node.downcast_ref::<CacheExec>().is_none() {
return internal_err!("Expected CacheExec, got {}", node.name());
}
Expand Down
Loading
Loading