Skip to content
Closed
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
6 changes: 6 additions & 0 deletions docs/source/user-guide/03-worker.md
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,12 @@ It receives a `WorkerQueryContext` with two fields:
- `headers` — the HTTP headers from the incoming request, handy for metadata like
authentication tokens or per-query configuration.

Implementations that need to trace or account for worker plan setup can also
override `WorkerSessionBuilder::run_plan_setup`. Its callback covers physical
plan decoding, worker plan hooks, and initial sampler startup. A wrapper that
succeeds must call the callback exactly once and return its result. Builders composed with
`MappedWorkerSessionBuilderExt::map` retain this behavior automatically.

```{note}
A worker only *executes* fragments — it never plans queries. So it needs your
codecs (to decode any custom nodes) but **not** the distributed planner or
Expand Down
1 change: 1 addition & 0 deletions src/test_utils/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,3 +10,4 @@ pub mod routing;
pub mod session_context;
pub mod test_work_unit_feed;
pub mod work_unit_file_scan;
pub mod worker_plan_setup;
8 changes: 8 additions & 0 deletions src/test_utils/worker_plan_setup.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
use crate::execution_plans::SamplerExec;
use datafusion::physical_plan::ExecutionPlan;
use std::sync::Arc;

/// Wraps a plan in the internal sampler node for worker setup integration tests.
pub fn wrap_in_sampler(plan: Arc<dyn ExecutionPlan>) -> Arc<dyn ExecutionPlan> {
Arc::new(SamplerExec::new(plan))
}
34 changes: 24 additions & 10 deletions src/worker/impl_coordinator_channel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -83,16 +83,30 @@ impl Worker {
})
.await?;

let codec = DistributedCodec::new_combined_with_user(session_state.config());
let task_ctx = session_state.task_ctx();
let proto_node = PhysicalPlanNode::try_decode(request.plan_proto.as_ref())?;
let mut plan = proto_node.try_into_physical_plan(&task_ctx, &codec)?;

for hook in self.hooks.on_plan.iter() {
plan = hook(plan, session_state.config())?;
}
load_info_rxs =
SamplerExec::kick_off_first_sampler(Arc::clone(&plan), Arc::clone(&task_ctx))?;
let mut plan_setup = None;
self.session_builder
.run_plan_setup(&session_state, &mut || {
let codec = DistributedCodec::new_combined_with_user(session_state.config());
let task_ctx = session_state.task_ctx();
let proto_node = PhysicalPlanNode::try_decode(request.plan_proto.as_ref())?;
let mut plan = proto_node.try_into_physical_plan(&task_ctx, &codec)?;

for hook in self.hooks.on_plan.iter() {
plan = hook(plan, session_state.config())?;
}
let sampler_receivers = SamplerExec::kick_off_first_sampler(
Arc::clone(&plan),
Arc::clone(&task_ctx),
)?;
plan_setup = Some((plan, task_ctx, sampler_receivers));
Ok(())
})?;
let Some((plan, task_ctx, sampler_receivers)) = plan_setup else {
return internal_err!(
"WorkerSessionBuilder::run_plan_setup did not run the setup callback"
);
};
load_info_rxs = sampler_receivers;

// Initialize partition count to the number of partitions in the stage
Ok::<_, DataFusionError>(TaskData {
Expand Down
99 changes: 99 additions & 0 deletions src/worker/session_builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,20 @@ pub trait WorkerSessionBuilder {
&self,
ctx: WorkerQueryContext,
) -> Result<SessionState, DataFusionError>;

/// Runs the synchronous setup for a plan received by the worker.
///
/// The callback decodes the physical plan, applies worker plan hooks, and starts any initial
/// sampling work. Implementations may override this method to wrap that work with tracing or
/// resource accounting. An implementation that returns `Ok(())` must invoke `setup` exactly
/// once and return its result; it may reject setup by returning an error without invoking it.
fn run_plan_setup(
&self,
_session_state: &SessionState,
setup: &mut dyn FnMut() -> Result<(), DataFusionError>,
) -> Result<(), DataFusionError> {
setup()
}
}

/// Noop implementation of the [WorkerSessionBuilder]. Used by default if no [WorkerSessionBuilder]
Expand Down Expand Up @@ -156,4 +170,89 @@ where
let builder = SessionStateBuilder::new_from_existing(state);
(self.f)(builder)
}

fn run_plan_setup(
&self,
session_state: &SessionState,
setup: &mut dyn FnMut() -> Result<(), DataFusionError>,
) -> Result<(), DataFusionError> {
self.inner.run_plan_setup(session_state, setup)
}
}

#[cfg(test)]
mod tests {
use super::*;
use datafusion::common::{assert_contains, internal_err};
use std::sync::atomic::{AtomicUsize, Ordering};

#[test]
fn default_plan_setup_runs_callback_once() {
let state = SessionStateBuilder::new().build();
let mut callback_calls = 0;

DefaultSessionBuilder
.run_plan_setup(&state, &mut || {
callback_calls += 1;
Ok(())
})
.unwrap();

assert_eq!(callback_calls, 1);
}

#[test]
fn default_plan_setup_propagates_callback_error() {
let state = SessionStateBuilder::new().build();

let error = DefaultSessionBuilder
.run_plan_setup(&state, &mut || internal_err!("plan setup failed"))
.expect_err("the callback error should be returned");

assert_contains!(error.to_string(), "plan setup failed");
}

#[test]
fn mapped_plan_setup_delegates_to_inner_builder() {
let wrapper_calls = Arc::new(AtomicUsize::new(0));
let builder = RecordingSessionBuilder {
wrapper_calls: Arc::clone(&wrapper_calls),
}
.map(|builder| Ok(builder.build()));
let state = SessionStateBuilder::new().build();
let mut callback_calls = 0;

builder
.run_plan_setup(&state, &mut || {
callback_calls += 1;
Ok(())
})
.unwrap();

assert_eq!(wrapper_calls.load(Ordering::Relaxed), 1);
assert_eq!(callback_calls, 1);
}

struct RecordingSessionBuilder {
wrapper_calls: Arc<AtomicUsize>,
}

#[async_trait]
impl WorkerSessionBuilder for RecordingSessionBuilder {
async fn build_session_state(
&self,
ctx: WorkerQueryContext,
) -> Result<SessionState, DataFusionError> {
Ok(ctx.builder.build())
}

fn run_plan_setup(
&self,
_session_state: &SessionState,
setup: &mut dyn FnMut() -> Result<(), DataFusionError>,
) -> Result<(), DataFusionError> {
self.wrapper_calls.fetch_add(1, Ordering::Relaxed);
setup()
}
}
}
Loading
Loading