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
44 changes: 27 additions & 17 deletions vortex-array/src/executor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -476,6 +476,14 @@ impl Executable for ArrayRef {
for (slot_idx, slot) in array.slots().iter().enumerate() {
let Some(child) = slot else { continue };
if let Some(reduced_parent) = child.reduce_parent(&array, slot_idx)? {
log_parent_rewrite(
ctx,
"reduce_parent",
slot_idx,
&array,
child,
&reduced_parent,
);
reduced_parent.statistics().inherit_from(array.statistics());
trace_op!(record_single_step_applied(
"reduce_parent",
Expand All @@ -500,13 +508,6 @@ impl Executable for ArrayRef {
kernels,
ctx,
)? {
ctx.log(format_args!(
"execute_parent: slot[{}]({}) rewrote {} -> {}",
slot_idx,
child.encoding_id(),
array,
executed_parent
));
executed_parent
.statistics()
.inherit_from(array.statistics());
Expand Down Expand Up @@ -618,7 +619,7 @@ fn finalize_done(
}

fn execute_parent_for_child(
_phase: &'static str,
phase: &'static str,
parent: &ArrayRef,
child: &ArrayRef,
slot_idx: usize,
Expand All @@ -642,8 +643,9 @@ fn execute_parent_for_child(
"Executed parent canonical dtype mismatch"
);
}
log_parent_rewrite(ctx, phase, slot_idx, parent, child, &result);
trace_op!(record_session_execute_parent_applied(
_phase,
phase,
parent,
child,
slot_idx,
Expand All @@ -653,7 +655,7 @@ fn execute_parent_for_child(
return Ok(Some(result));
}
trace_op!(record_session_execute_parent_declined(
_phase,
phase,
parent,
child,
slot_idx,
Expand All @@ -665,6 +667,21 @@ fn execute_parent_for_child(
Ok(None)
}

/// Log a parent rewrite to the execution log.
fn log_parent_rewrite(
ctx: &mut ExecutionCtx,
phase: &'static str,
slot_idx: usize,
parent: &ArrayRef,
child: &ArrayRef,
output: &ArrayRef,
) {
ctx.log(format_args!(
"{phase}: slot[{slot_idx}]({}) rewrote {parent} -> {output}",
child.encoding_id(),
));
}

/// Try execute_parent on each occupied slot of the array.
fn try_execute_parent(
array: &ArrayRef,
Expand All @@ -676,13 +693,6 @@ fn try_execute_parent(
if let Some(executed_parent) =
execute_parent_for_child("child_execute_parent", array, child, slot_idx, kernels, ctx)?
{
ctx.log(format_args!(
"execute_parent: slot[{}]({}) rewrote {} -> {}",
slot_idx,
child.encoding_id(),
array,
executed_parent
));
executed_parent
.statistics()
.inherit_from(array.statistics());
Expand Down
47 changes: 18 additions & 29 deletions vortex-array/src/optimizer/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,27 +59,26 @@ pub trait ArrayOptimizer {

impl ArrayOptimizer for ArrayRef {
fn optimize(&self) -> VortexResult<ArrayRef> {
Ok(try_optimize(self, None)?.unwrap_or_else(|| self.clone()))
Ok(optimize_owned(self.clone(), None)?.0)
}

fn optimize_ctx(&self, session: &VortexSession) -> VortexResult<ArrayRef> {
Ok(try_optimize(self, Some(session))?.unwrap_or_else(|| self.clone()))
Ok(optimize_owned(self.clone(), Some(session))?.0)
}

fn optimize_recursive(&self, session: &VortexSession) -> VortexResult<ArrayRef> {
Ok(try_optimize_recursive(self, session)?.unwrap_or_else(|| self.clone()))
Ok(try_optimize_recursive(self.clone(), session)?.0)
}
}

fn try_optimize(
array: &ArrayRef,
fn optimize_owned(
mut current_array: ArrayRef,
session: Option<&VortexSession>,
) -> VortexResult<Option<ArrayRef>> {
let mut current_array = array.clone();
) -> VortexResult<(ArrayRef, bool)> {
let mut any_optimizations = false;
let session_kernels = session.map(|session| session.kernels());

trace_op!(record_optimize_start(array, session.is_some()));
trace_op!(record_optimize_start(&current_array, session.is_some()));

for _ in 0..=MAX_OPTIMIZER_REWRITE_PASS {
trace_op!(record_optimize_loop_start(&current_array));
Expand Down Expand Up @@ -127,7 +126,7 @@ fn try_optimize(

trace_op!(record_optimize_done(&current_array, any_optimizations));

return Ok(any_optimizations.then_some(current_array));
return Ok((current_array, any_optimizations));
}

vortex_bail!("Exceeded maximum optimization iterations (possible infinite loop)");
Expand Down Expand Up @@ -171,35 +170,29 @@ fn try_session_parent_reduce(
}

fn try_optimize_recursive(
array: &ArrayRef,
current_array: ArrayRef,
session: &VortexSession,
) -> VortexResult<Option<ArrayRef>> {
let mut current_array = array.clone();
let mut any_optimizations = false;
) -> VortexResult<(ArrayRef, bool)> {
trace_op!(record_optimize_recursive_start(&current_array));

trace_op!(record_optimize_recursive_start(array));

if let Some(new_array) = try_optimize(&current_array, Some(session))? {
current_array = new_array;
any_optimizations = true;
}
let (mut current_array, mut any_optimizations) = optimize_owned(current_array, Some(session))?;

let mut new_slots = SmallVec::with_capacity(current_array.slots().len());
let mut any_slot_optimized = false;
for slot in current_array.slots() {
match slot {
Some(child) => {
if let Some(new_child) = try_optimize_recursive(child, session)? {
let (new_child, new_slot_optimized) =
try_optimize_recursive(child.clone(), session)?;
if new_slot_optimized {
trace_op!(record_optimize_recursive_slot(
new_slots.len(),
child,
&new_child,
));
new_slots.push(Some(new_child));
any_slot_optimized = true;
} else {
new_slots.push(Some(child.clone()));
}
new_slots.push(Some(new_child));
any_slot_optimized |= new_slot_optimized;
}
None => new_slots.push(None),
}
Expand All @@ -212,9 +205,5 @@ fn try_optimize_recursive(
any_optimizations = true;
}

if any_optimizations {
Ok(Some(current_array))
} else {
Ok(None)
}
Ok((current_array, any_optimizations))
}
Loading