(&self, batch: &RecordBatch)
| 113 | } |
| 114 | |
| 115 | fn evaluate(&self, batch: &RecordBatch) -> Result<ArrayRef> { |
| 116 | let mut evaluator = self.expr.create_evaluator()?; |
| 117 | let num_rows = batch.num_rows(); |
| 118 | if evaluator.uses_window_frame() { |
| 119 | let sort_options = self.order_by.iter().map(|o| o.options).collect(); |
| 120 | let mut row_wise_results = vec![]; |
| 121 | |
| 122 | let mut values = self.evaluate_args(batch)?; |
| 123 | let order_bys = get_orderby_values(self.order_by_columns(batch)?); |
| 124 | let n_args = values.len(); |
| 125 | values.extend(order_bys); |
| 126 | let order_bys_ref = &values[n_args..]; |
| 127 | |
| 128 | let mut window_frame_ctx = |
| 129 | WindowFrameContext::new(Arc::clone(&self.window_frame), sort_options); |
| 130 | let mut last_range = Range { start: 0, end: 0 }; |
| 131 | // We iterate on each row to calculate window frame range and and window function result |
| 132 | for idx in 0..num_rows { |
| 133 | let range = window_frame_ctx.calculate_range( |
| 134 | order_bys_ref, |
| 135 | &last_range, |
| 136 | num_rows, |
| 137 | idx, |
| 138 | )?; |
| 139 | let value = evaluator.evaluate(&values, &range)?; |
| 140 | row_wise_results.push(value); |
| 141 | last_range = range; |
| 142 | } |
| 143 | ScalarValue::iter_to_array(row_wise_results) |
| 144 | } else if evaluator.include_rank() { |
| 145 | let columns = self.order_by_columns(batch)?; |
| 146 | let sort_partition_points = evaluate_partition_ranges(num_rows, &columns)?; |
| 147 | evaluator.evaluate_all_with_rank(num_rows, &sort_partition_points) |
| 148 | } else { |
| 149 | let values = self.evaluate_args(batch)?; |
| 150 | evaluator.evaluate_all(&values, num_rows) |
| 151 | } |
| 152 | } |
| 153 | |
| 154 | /// Evaluate the window function against the batch. This function facilitates |
| 155 | /// stateful, bounded-memory implementations. |
no test coverage detected