rluvaton commented on code in PR #25877:
URL: https://github.com/apache/datafusion/pull/25877#discussion_r4177002145


##########
datafusion/functions-aggregate/src/count.rs:
##########
@@ -824,6 +838,125 @@ impl GroupsAccumulator for CountGroupsAccumulator {
     }
 }
 
+/// [`CountGroupsAccumulator`] with the counts stored in blocks of
+/// `block_size` groups, so growing past a block never reallocates and
+/// copies the existing counts.
+#[derive(Debug)]
+struct BlockedCountGroupsAccumulator {
+    /// Count per group, see [`CountGroupsAccumulator::counts`].
+    counts: BlockedVec<i64>,
+}
+
+impl BlockedCountGroupsAccumulator {
+    fn new(block_size: usize) -> Self {
+        Self {
+            counts: BlockedVec::new(block_size),
+        }
+    }
+
+    fn counts_to_array(counts: Vec<i64>) -> ArrayRef {
+        // zero copy, count is never null
+        Arc::new(Int64Array::new(counts.into(), None))
+    }
+
+    fn blocks_to_arrays(blocks: impl IntoIterator<Item = Vec<i64>>) -> 
Vec<ArrayRef> {
+        blocks.into_iter().map(Self::counts_to_array).collect()
+    }
+
+    fn take(&mut self, emit_to: BlockedEmitTo) -> Vec<ArrayRef> {
+        match emit_to {
+            BlockedEmitTo::All => 
Self::blocks_to_arrays(self.counts.take_all()),
+            BlockedEmitTo::NextBlock => {
+                Self::blocks_to_arrays(self.counts.take_next_block())
+            }
+            BlockedEmitTo::First(n) => {
+                Self::blocks_to_arrays([self.counts.take_first(n)])
+            }
+        }
+    }
+}
+
+impl BlockedGroupsAccumulator for BlockedCountGroupsAccumulator {
+    fn block_size(&self) -> usize {
+        self.counts.block_size()
+    }
+
+    fn update_batch(
+        &mut self,
+        values: &[ArrayRef],
+        group_indices: &[BlocksIndex],
+        opt_filter: Option<&BooleanArray>,
+        total_num_groups: usize,
+    ) -> Result<()> {
+        assert_eq!(values.len(), 1, "single argument to update_batch");
+        let values = &values[0];
+        let nulls = values.logical_nulls().filter(|n| n.null_count() > 0);
+
+        self.counts.grow_to(total_num_groups, 0);
+
+        // Add one to each group's counter for each non null, non filtered 
value
+        // SAFETY: group_index is guaranteed to be in bounds and less than 
total_num_groups
+        unsafe {
+            self.counts.update_unchecked(

Review Comment:
   I had unchecked here to keep the same optimizations used as before so this 
can be as a follow up pr but not in this



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