andygrove commented on code in PR #1916:
URL: 
https://github.com/apache/datafusion-ballista/pull/1916#discussion_r4115929052


##########
ballista/scheduler/src/state/execution_stage.rs:
##########
@@ -665,6 +668,24 @@ impl RunningStage {
             warn!("The metrics for stage {} should not be none", 
self.stage_id);
             vec![]
         });
+        let mut output_column_stats: Vec<ColumnStatistics> = vec![];
+        for info in self.task_infos.iter() {
+            if let task_status::Status::Successful(task_status) = 
&info.task_status {
+                if output_column_stats.is_empty() {
+                    output_column_stats = vec![
+                        ColumnStatistics::new_unknown();
+                        task_status.task_column_stats.len()
+                    ];
+                }
+                for col_stats in &task_status.task_column_stats {
+                    let slot = &mut output_column_stats[col_stats.column as 
usize];

Review Comment:
   This indexes into a vec sized from the first successful task's report, using 
a `column` value that came over the wire. If a later task sends a wider list, 
or a sparse one (which the `column` field seems designed to allow), this panics 
with index out of bounds. I reproduced that with a small test. It runs inside 
`succeed_stage` on the query stage event loop, after the stage has already been 
removed from `self.stages`, so a panic here would stop the scheduler from 
processing events.
   
   Could we size the vec from `self.plan.schema()` and use `get_mut` so a bad 
index gets skipped? The shuffle writer's `schema()` is the data schema, so the 
width lines up.
   
   There's a related case too. A successful task that sent no stats, like an 
older executor during a rolling upgrade, is skipped silently. Once the seed is 
`Exact(0)` that turns into an undercount that still says `Exact`. Marking any 
column a task didn't report as `Precision::Absent` keeps it honest, and 
`Absent` then sticks for the rest of the fold. Something like this passes the 
existing scheduler state tests for me:
   
   ```rust
   let num_columns = self.plan.schema().fields().len();
   let mut output_column_stats = vec![
       ColumnStatistics {
           null_count: Precision::Exact(0),
           ..ColumnStatistics::new_unknown()
       };
       num_columns
   ];
   for info in &self.task_infos {
       let task_status::Status::Successful(task) = &info.task_status else {
           continue;
       };
       let mut reported = vec![false; num_columns];
       for stats in &task.task_column_stats {
           let column = stats.column as usize;
           if let Some(slot) = output_column_stats.get_mut(column) {
               slot.null_count = slot
                   .null_count
                   .add(&Precision::Exact(stats.null_count as usize));
               reported[column] = true;
           }
       }
       for (slot, reported) in output_column_stats.iter_mut().zip(reported) {
           if !reported {
               slot.null_count = Precision::Absent;
           }
       }
   }
   ```



##########
ballista/core/src/execution_plans/sort_shuffle/writer.rs:
##########
@@ -630,10 +644,14 @@ impl SortShuffleWriterExec {
             // `MemoryPool` as the sole spill trigger.
             let memory_limit = config.memory_limit_per_task_bytes;
             let per_task_budget_enabled = memory_limit > 0;
+            let mut null_counts = vec![0u64; schema.fields().len()];
 
             while let Some(result) = stream.next().await {
                 let input_batch = result?;
                 metrics.input_rows.add(input_batch.num_rows());
+                for (i, col) in input_batch.columns().iter().enumerate() {
+                    null_counts[i] += col.null_count() as u64;

Review Comment:
   `null_count()` is the physical count from the null buffer, so a 
`DataType::Null` column reports 0 even though every row is null. Dictionary and 
run-end encoded arrays can undercount the same way. I checked with a 
`NullArray::new(3)` column and got 0. DataFusion's own 
`compute_record_batch_statistics` uses `logical_nulls()` for this reason. Could 
this be `col.logical_null_count()`? It matters once these stats feed a plan, 
because DataFusion's `AggregateStatistics` answers `COUNT(col)` as `num_rows - 
null_count` when the count is `Exact`.



##########
ballista/executor/src/execution_engine.rs:
##########
@@ -84,7 +85,7 @@ pub trait QueryStageExecutor: Sync + Send + Debug + Display {
         &self,
         task_id: usize,
         context: Arc<TaskContext>,
-    ) -> Result<Vec<ShuffleWritePartition>>;
+    ) -> Result<ShuffleWriteResult>;

Review Comment:
   Changing the return type here breaks anyone implementing 
`QueryStageExecutor` outside the repo, and it ripples into `as_task_status` and 
`Executor::execute_query_stage` as well. The trait already has a pattern for 
this kind of side data. `collect_runtime_stats_reports()` and 
`collect_window_state_reports()` are default methods that get called right 
after `execute_query_stage`, and their results ride along in 
`TaskCompletionExtras`. That struct is `#[non_exhaustive]` with a `Default` so 
it can grow without breaking callers. Would a default 
`collect_column_stats(&self) -> Vec<TaskColumnStats>` plus a `column_stats` 
field on `TaskCompletionExtras` work here? Then `execute_query_stage` keeps its 
signature and `ShuffleWriteResult` isn't needed. If you'd rather keep the new 
return type, it'll need the `api-change` label and an entry in 
`docs/source/upgrading/55.0.0.md`.



##########
ballista/scheduler/src/state/execution_stage.rs:
##########
@@ -665,6 +668,24 @@ impl RunningStage {
             warn!("The metrics for stage {} should not be none", 
self.stage_id);
             vec![]
         });
+        let mut output_column_stats: Vec<ColumnStatistics> = vec![];
+        for info in self.task_infos.iter() {
+            if let task_status::Status::Successful(task_status) = 
&info.task_status {
+                if output_column_stats.is_empty() {
+                    output_column_stats = vec![
+                        ColumnStatistics::new_unknown();

Review Comment:
   I think this fold always ends up as `Absent`. 
`ColumnStatistics::new_unknown()` seeds `null_count` with `Precision::Absent`, 
and `Precision::add` returns `Absent` whenever either side is `Absent`, so the 
sum never gets started. I tried a quick test where two successful tasks report 
`[1, 2]` and `[3, 4]`, and `to_successful()` returned `[Absent, Absent]` 
instead of `[Exact(4), Exact(6)]`. Seeding each slot with `Precision::Exact(0)` 
fixes it. It might be worth adding a test like that here, since nothing 
exercises `to_successful` with column stats yet and that's how this slipped 
past CI.



##########
ballista/core/proto/ballista.proto:
##########
@@ -859,6 +875,7 @@ message SuccessfulTask {
   // (`sort_shuffle::get_index_path`), and send only a reference here. That
   // keeps the completion message fixed-size regardless of aggregate.
   repeated WindowStateReport window_state = 4;
+  repeated TaskColumnStats taskColumnStats = 5;

Review Comment:
   Small one: the rest of this file uses snake_case field names, so 
`task_column_stats` would match. Prost already generates that name, and since 
this hasn't shipped the rename is free on the wire. Short doc comments on 
`TaskColumnStats` would help the follow-up PRs too, for example what `column` 
indexes into and that an empty list means stats weren't collected.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to