Skip to content
Merged
Show file tree
Hide file tree
Changes from 18 commits
Commits
Show all changes
27 commits
Select commit Hold shift + click to select a range
d417dc6
ensure apply_expressions is implemented for DistributedLeafExec
jayshrivastava Aug 12, 2026
99a9ab6
remove extra comment
jayshrivastava Aug 12, 2026
429fdec
display dynamic filters during execution
jayshrivastava Aug 11, 2026
b8c282b
remove viz variants
jayshrivastava Aug 11, 2026
64f9f5d
store refactor
jayshrivastava Aug 11, 2026
1115fb4
dont use vec<u8>
jayshrivastava Aug 11, 2026
00e3f5b
rename
jayshrivastava Aug 11, 2026
72a8e6c
refactor tests
jayshrivastava Aug 11, 2026
730266f
remove allowlist from set plan request
jayshrivastava Aug 21, 2026
f69728c
remove plan_for_display
jayshrivastava Aug 21, 2026
92fb135
make displaying align with metrics
jayshrivastava Aug 21, 2026
cdd153e
remove type aliases and use select_all
jayshrivastava Aug 25, 2026
886cf13
qualified import
jayshrivastava Aug 25, 2026
7a14f8d
refactors
jayshrivastava Aug 31, 2026
96ec2c2
handle error
jayshrivastava Aug 31, 2026
0f3fe49
unnest proto
jayshrivastava Aug 31, 2026
96d9835
lints and comments
jayshrivastava Aug 31, 2026
26297c8
fix: expose prepared visualization plan to rewrites
jayshrivastava Sep 1, 2026
7fcf1e4
refactor: centralize physical plan protobuf round trips
jayshrivastava Sep 1, 2026
c7db23c
move task context out of PreparedExecution
jayshrivastava Sep 2, 2026
42b9409
tests: use insta, change function visibility
jayshrivastava Sep 2, 2026
942a424
do not store trait in DiscoveredDynamicFilters
jayshrivastava Sep 2, 2026
c563f64
unify wait_for API for metrics with dynamic filters api
jayshrivastava Sep 2, 2026
6c7e24a
fix: include task-level metrics in rewritten stages
jayshrivastava Sep 2, 2026
2f10b0b
refactor prepared_execution`
jayshrivastava Sep 2, 2026
b183523
comments
jayshrivastava Sep 3, 2026
7f00630
snapshots
jayshrivastava Sep 3, 2026
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
5 changes: 5 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,11 @@ license = "Apache-2.0"
documentation = "https://datafusion-contrib.github.io/datafusion-distributed/"
repository = "https://github.com/datafusion-contrib/datafusion-distributed"

[[test]]
name = "dynamic_filtering"
path = "tests/dynamic_filtering/main.rs"
required-features = ["integration"]

[dependencies]
chrono = { version = "0.4.44" }
datafusion = { workspace = true, features = [
Expand Down
13 changes: 10 additions & 3 deletions docs/source/user-guide/05-metrics.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,11 @@ channel, so they are not lost even if the result stream is dropped early (for ex

## Rendering a plan with metrics

Two functions, both exported from the crate root, do the work:
These functions, all exported from the crate root, do the work:

- `rewrite_distributed_plan_with_dynamic_filters(plan)` — folds the completed dynamic filters
reported by each worker task into an isolated copy of the plan. When displaying both dynamic
filters and metrics, apply the dynamic-filter rewrite first.
- `rewrite_distributed_plan_with_metrics(plan, format)` — folds every task's metrics back into the
coordinator's copy of the plan. It waits for all worker metrics to arrive, so the result is always
complete. The `format` is a `DistributedMetricsFormat`:
Comment thread
gabotechs marked this conversation as resolved.
Outdated
Expand All @@ -54,11 +57,15 @@ execute_stream(plan.clone(), ctx.task_ctx())?
.try_collect::<Vec<_>>()
.await?;

// 3. Fold the per-task metrics back into the plan...
// 3. Fold the completed per-task dynamic filters back into the plan...
let plan =
rewrite_distributed_plan_with_dynamic_filters(plan).await?;

// 4. Fold the per-task metrics back into the plan...
let plan =
rewrite_distributed_plan_with_metrics(plan, DistributedMetricsFormat::Aggregated).await?;

// 4. ...and render it.
// 5. ...and render it.
println!("{}", display_plan_ascii(plan.as_ref(), true));
```

Expand Down
3 changes: 2 additions & 1 deletion src/codec/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,8 @@ mod user_codec;

pub use distributed_codec::DistributedCodec;
pub(crate) use physical_plan::{
decode_execution_plan, decode_partitioning, encode_execution_plan, encode_partitioning,
decode_execution_plan, decode_partitioning, decode_physical_expr, encode_execution_plan,
encode_partitioning, encode_physical_expr,
};
pub(crate) use user_codec::{
get_distributed_user_codecs, set_distributed_user_codec, set_distributed_user_codec_arc,
Expand Down
28 changes: 25 additions & 3 deletions src/codec/physical_plan.rs
Original file line number Diff line number Diff line change
@@ -1,15 +1,17 @@
use super::DistributedCodec;
use datafusion::arrow::datatypes::SchemaRef;
use datafusion::arrow::datatypes::{Schema, SchemaRef};
use datafusion::common::Result;
use datafusion::execution::TaskContext;
use datafusion::physical_expr::Partitioning;
use datafusion::physical_expr::{Partitioning, PhysicalExpr};
use datafusion::physical_plan::ExecutionPlan;
use datafusion_proto::bytes::{
physical_plan_from_bytes_with_proto_converter, physical_plan_to_bytes_with_proto_converter,
};
use datafusion_proto::physical_plan::from_proto::parse_protobuf_partitioning;
use datafusion_proto::physical_plan::to_proto::serialize_partitioning;
use datafusion_proto::physical_plan::{DeduplicatingProtoConverter, PhysicalPlanDecodeContext};
use datafusion_proto::physical_plan::{
DeduplicatingProtoConverter, PhysicalPlanDecodeContext, PhysicalProtoConverterExtension,
};
use datafusion_proto::protobuf;
use datafusion_proto::protobuf::proto_error;
use prost::Message;
Expand Down Expand Up @@ -42,6 +44,26 @@ pub(crate) fn decode_execution_plan(
physical_plan_from_bytes_with_proto_converter(encoded, task_ctx, &codec, &converter)
}

pub(crate) fn encode_physical_expr(
expression: &Arc<dyn PhysicalExpr>,
task_ctx: &TaskContext,
) -> Result<protobuf::PhysicalExprNode> {
let codec = DistributedCodec::new_combined_with_user(task_ctx.session_config());
let converter = new_proto_converter();
converter.physical_expr_to_proto(expression, &codec)
}

pub(crate) fn decode_physical_expr(
proto: &protobuf::PhysicalExprNode,
input_schema: &Schema,
task_ctx: &TaskContext,
) -> Result<Arc<dyn PhysicalExpr>> {
let codec = DistributedCodec::new_combined_with_user(task_ctx.session_config());
let decode_ctx = PhysicalPlanDecodeContext::new(task_ctx, &codec);
let converter = new_proto_converter();
converter.proto_to_physical_expr(proto, input_schema, &decode_ctx)
}

pub(crate) fn encode_partitioning(
partitioning: &Partitioning,
task_ctx: &TaskContext,
Expand Down
50 changes: 46 additions & 4 deletions src/common/maybe_encoded.rs
Original file line number Diff line number Diff line change
@@ -1,17 +1,21 @@
use crate::codec::{
decode_execution_plan, decode_partitioning, encode_execution_plan, encode_partitioning,
decode_execution_plan, decode_partitioning, decode_physical_expr, encode_execution_plan,
encode_partitioning, encode_physical_expr,
};
use datafusion::arrow::datatypes::SchemaRef;
use datafusion::arrow::datatypes::{Schema, SchemaRef};
use datafusion::common::{Result, internal_err};
use datafusion::execution::TaskContext;
use datafusion::physical_expr::Partitioning;
use datafusion::physical_expr::{Partitioning, PhysicalExpr};
use datafusion::physical_plan::ExecutionPlan;
use datafusion_proto::protobuf::PhysicalExprNode;
use datafusion_proto::protobuf::proto_error;
use prost::Message;
use std::sync::Arc;

/// A value that a transport may either leave encoded or materialize in memory.
/// Users are free to pass [MaybeEncoded::Encoded] or [MaybeEncoded::Decoded] at any
/// moment and Distributed DataFusion's code will internally know how to handle it.
#[derive(Clone)]
#[derive(Clone, Debug)]
pub enum MaybeEncoded<T> {
Encoded(Vec<u8>),
Decoded(T),
Expand Down Expand Up @@ -82,6 +86,44 @@ impl MaybeEncoded<Partitioning> {
}
}

impl MaybeEncoded<Arc<dyn PhysicalExpr>> {
/// Returns the encoded [`PhysicalExpr`] as protobuf bytes:
/// - If in `Decoded` state, it encodes it using the codecs registered in the [`TaskContext`].
/// - If in `Encoded` state, it passes through the existing bytes.
pub fn encode(self, ctx: &Arc<TaskContext>) -> Result<Vec<u8>> {
match self {
Self::Encoded(encoded) => Ok(encoded),
Self::Decoded(expression) => {
Ok(encode_physical_expr(&expression, ctx)?.encode_to_vec())
}
}
}

/// Returns the decoded [`PhysicalExpr`].
/// - If in `Decoded` state, it passes through the expression.
/// - If in `Encoded` state, it decodes it using the provided schema and task context.
pub fn decode(
self,
input_schema: &Schema,
task_ctx: &TaskContext,
) -> Result<Arc<dyn PhysicalExpr>> {
self.decode_with(|encoded| {
let proto = PhysicalExprNode::decode(encoded.as_slice())
.map_err(|error| proto_error(error.to_string()))?;
decode_physical_expr(&proto, input_schema, task_ctx)
})
}

/// Materializes the expression's protobuf representation without changing the stored form.
pub(crate) fn to_proto(&self, task_ctx: &TaskContext) -> Result<PhysicalExprNode> {
match self {
Self::Encoded(encoded) => PhysicalExprNode::decode(encoded.as_slice())
.map_err(|error| proto_error(error.to_string())),
Self::Decoded(expression) => encode_physical_expr(expression, task_ctx),
}
}
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down
Loading
Loading