Rich-T-kid commented on code in PR #11174:
URL: https://github.com/apache/arrow-rs/pull/11174#discussion_r4077965989


##########
arrow-select/src/coalesce.rs:
##########


Review Comment:
   similar point here. we need to be able to expose this to the 
`inProgressArrays` trait so it can avoid this copy. ofc the generic path can 
still fall back to this.



##########
arrow-select/src/coalesce.rs:
##########
@@ -298,6 +299,43 @@ impl BatchCoalescer {
         self.push_batch(taken_batch)
     }
 
+    /// Push rows gathered from multiple [`RecordBatch`]es into the Coalescer.
+    ///
+    /// This is semantically equivalent to calling [`interleave_record_batch`]
+    /// followed by [`Self::push_batch`], but avoids allocating an intermediate
+    /// [`RecordBatch`].
+    ///
+    /// Each element of `indices` is a `(batch_index, row_index)` pair that
+    /// identifies a row in `batches`.
+    ///
+    /// # Example
+    /// ```
+    /// # use arrow_array::record_batch;
+    /// # use arrow_select::coalesce::BatchCoalescer;
+    /// let batch1 = record_batch!(("a", Int32, [1, 2, 3])).unwrap();
+    /// let batch2 = record_batch!(("a", Int32, [4, 5, 6])).unwrap();
+    /// let indices = vec![(0, 2), (1, 0), (0, 0), (1, 2)]; // rows: 3, 4, 1, 6
+    /// let mut coalescer = BatchCoalescer::new(batch1.schema(), 1000);
+    /// coalescer
+    ///     .push_batch_interleaved(&[&batch1, &batch2], &indices)
+    ///     .unwrap();
+    /// coalescer.finish_buffered_batch().unwrap();
+    /// let out = coalescer.next_completed_batch().unwrap();
+    /// let expected = record_batch!(("a", Int32, [3, 4, 1, 6])).unwrap();
+    /// assert_eq!(out, expected);
+    /// ```
+    pub fn push_batch_interleaved(
+        &mut self,
+        batches: &[&RecordBatch],
+        indices: &[(usize, usize)],
+    ) -> Result<(), ArrowError> {
+        if indices.is_empty() {
+            return Ok(());
+        }
+        let interleaved = interleave_record_batch(batches, indices)?;
+        self.push_batch(interleaved)

Review Comment:
   we should probably expose this `inProgressArrays` so they can take advantage 
of it. 



##########
arrow/benches/coalesce_kernels.rs:
##########
@@ -611,9 +611,99 @@ fn add_all_take_benchmarks(c: &mut Criterion) {
     }
 }
 
-criterion_group!(benches, add_all_filter_benchmarks, add_all_take_benchmarks);
+criterion_group!(
+    benches,
+    add_all_filter_benchmarks,
+    add_all_take_benchmarks,
+    add_all_interleave_benchmarks
+);
 criterion_main!(benches);
 
+fn add_all_interleave_benchmarks(c: &mut Criterion) {
+    let batch_size = 8192;
+    let num_source_batches = 4;
+
+    let primitive_schema = SchemaRef::new(Schema::new(vec![
+        Field::new("int32_val", DataType::Int32, true),
+        Field::new("float_val", DataType::Float64, true),
+        Field::new(
+            "timestamp_val",
+            DataType::Timestamp(TimeUnit::Nanosecond, Some("UTC".into())),
+            true,
+        ),
+    ]));
+
+    let single_utf8view_schema = SchemaRef::new(Schema::new(vec![Field::new(
+        "value",
+        DataType::Utf8View,
+        true,
+    )]));
+
+    let mixed_fsb_schema = SchemaRef::new(Schema::new(vec![
+        Field::new("fsb16_val", DataType::FixedSizeBinary(16), true),
+        Field::new("fsb32_val", DataType::FixedSizeBinary(32), true),
+    ]));
+
+    // Binary, boolean, and dictionary all go through the generic path; group
+    // them in one schema so a single benchmark covers the whole family.
+    let mixed_generic_schema = SchemaRef::new(Schema::new(vec![
+        Field::new("binary_val", DataType::Binary, true),
+        Field::new("bool_val", DataType::Boolean, true),
+        Field::new(
+            "dict_val",
+            DataType::Dictionary(Box::new(DataType::Int32), 
Box::new(DataType::Utf8)),
+            true,
+        ),
+    ]));
+
+    for null_density in [0.0, 0.1] {
+        for selectivity in [0.001, 0.01, 0.1, 0.8] {
+            for scenario in [

Review Comment:
   I may remove this.  we may just want to benchmark. from my understanding 
filtering isnt really a component for the two functions we are trying to 
optimize.



##########
arrow-select/src/coalesce.rs:
##########
@@ -298,6 +299,43 @@ impl BatchCoalescer {
         self.push_batch(taken_batch)
     }
 
+    /// Push rows gathered from multiple [`RecordBatch`]es into the Coalescer.
+    ///
+    /// This is semantically equivalent to calling [`interleave_record_batch`]
+    /// followed by [`Self::push_batch`], but avoids allocating an intermediate
+    /// [`RecordBatch`].
+    ///
+    /// Each element of `indices` is a `(batch_index, row_index)` pair that
+    /// identifies a row in `batches`.
+    ///
+    /// # Example
+    /// ```
+    /// # use arrow_array::record_batch;
+    /// # use arrow_select::coalesce::BatchCoalescer;
+    /// let batch1 = record_batch!(("a", Int32, [1, 2, 3])).unwrap();
+    /// let batch2 = record_batch!(("a", Int32, [4, 5, 6])).unwrap();
+    /// let indices = vec![(0, 2), (1, 0), (0, 0), (1, 2)]; // rows: 3, 4, 1, 6
+    /// let mut coalescer = BatchCoalescer::new(batch1.schema(), 1000);
+    /// coalescer
+    ///     .push_batch_interleaved(&[&batch1, &batch2], &indices)
+    ///     .unwrap();
+    /// coalescer.finish_buffered_batch().unwrap();
+    /// let out = coalescer.next_completed_batch().unwrap();
+    /// let expected = record_batch!(("a", Int32, [3, 4, 1, 6])).unwrap();
+    /// assert_eq!(out, expected);
+    /// ```
+    pub fn push_batch_interleaved(
+        &mut self,
+        batches: &[&RecordBatch],
+        indices: &[(usize, usize)],
+    ) -> Result<(), ArrowError> {

Review Comment:
   it may be akward needing to add tuples like this when instead we should 
maybe expect a slice
   
   ```suggestion
       pub fn push_batch_interleaved(
           &mut self,
           batches: &[&RecordBatch],
           indices: &[(usize, (usize,usize))], // array , (start,end)
       ) -> Result<(), ArrowError> {
   ```
   
   🤔 



-- 
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]

Reply via email to