viirya commented on code in PR #6824:
URL: https://github.com/apache/datafusion-comet/pull/6824#discussion_r4233131403


##########
native/shuffle/src/spark_unsafe/row.rs:
##########
@@ -668,8 +668,12 @@ fn append_list_column_batch(
         DataType::Date32 => {
             process_primitive_lists!(Date32Builder, append_dates_to_builder);
         }
-        DataType::Timestamp(TimeUnit::Microsecond, _) => {
-            process_primitive_lists!(TimestampMicrosecondBuilder, 
append_timestamps_to_builder);
+        DataType::Timestamp(TimeUnit::Microsecond, tz) => {
+            process_primitive_lists!(
+                TimestampMicrosecondBuilder,
+                append_timestamps_to_builder,

Review Comment:
   Forwarding `tz` here is what keeps JVM columnar shuffle working for 
`array<timestamp>` columns, but no test reaches it. I changed this argument to 
`&None` and all 176 shuffle crate tests still passed. The round-trip tests only 
go through `append_to_builder`, and `CometColumnarShuffleSuite` has no 
`array<timestamp>` case.
   
   If this forwarding ever breaks, `append_array`'s `assert_eq!` panics in 
release builds too. It only does so for a list of 64 or more elements that 
holds a null, which is exactly the shape small tests miss. Could you add a test 
that drives this path with a `Timestamp(Microsecond, Some(tz))` list column of 
at least 64 elements with a null? A `CometColumnarShuffleSuite` case on 
`array<timestamp>` would also check the schema end to end. A native test around 
`append_list_column_batch` would work too.



##########
native/shuffle/src/spark_unsafe/list.rs:
##########
@@ -39,55 +44,90 @@ use datafusion_comet_jni_bridge::errors::CometError;
 /// time with the copy, and arrays of 50 about a fifth.
 const MIN_BULK_APPEND_ELEMENTS: usize = 8;
 
+/// Element count from which an aligned nullable array that holds a null is 
appended with one copy
+/// of its values and a validity buffer built from its null bitset, instead of 
element by element.
+/// Building the two buffers takes four allocations, which cost as much as 
appending 40 to 55
+/// elements one by one, so shorter arrays stay on the loop: with the copy, 
arrays of 16 elements
+/// took 2.4 times as long and arrays of 32 elements 1.3 times. Arrays of 64 
elements take 65% of
+/// the time with the copy for 4-byte elements and 85% for 8-byte ones, and 
arrays of 1024 a ninth
+/// and a third.
+const MIN_BULK_NULLABLE_APPEND_ELEMENTS: usize = 64;
+
 /// Generates bulk append methods for primitive types in SparkUnsafeArray.
 ///
 /// # Safety invariants for all generated methods:
 /// - `element_offset` points to contiguous element data of length 
`num_elements`
 /// - `null_bitset_ptr()` returns a pointer to `ceil(num_elements/64)` i64 
words
 /// - These invariants are guaranteed by the SparkUnsafeArray layout from the 
JVM
+///
+/// A timestamp appender also takes the column's timezone, because 
`append_array` requires the

Review Comment:
   A timezone that does not match the builder's only panics on the bulk path, 
so a new caller that gets it wrong passes every short-array test. Could this 
doc state the contract directly? For example: "`timezone` must be the builder's 
timezone. A mismatch panics in `append_array` for arrays of 
`MIN_BULK_NULLABLE_APPEND_ELEMENTS` or more elements that hold a null."
   
   Also, the comment in `has_null()` (line 223) says a stray padding bit "would 
only send the array down the per-element path". With this change, it can send 
an aligned array of 64 or more elements to the bulk path instead. That is still 
correct, because the `BooleanBuffer` length excludes the padding bits. Could 
that comment be updated to match?



##########
native/shuffle/src/spark_unsafe/list.rs:
##########
@@ -39,55 +44,90 @@ use datafusion_comet_jni_bridge::errors::CometError;
 /// time with the copy, and arrays of 50 about a fifth.
 const MIN_BULK_APPEND_ELEMENTS: usize = 8;
 
+/// Element count from which an aligned nullable array that holds a null is 
appended with one copy
+/// of its values and a validity buffer built from its null bitset, instead of 
element by element.
+/// Building the two buffers takes four allocations, which cost as much as 
appending 40 to 55
+/// elements one by one, so shorter arrays stay on the loop: with the copy, 
arrays of 16 elements
+/// took 2.4 times as long and arrays of 32 elements 1.3 times. Arrays of 64 
elements take 65% of
+/// the time with the copy for 4-byte elements and 85% for 8-byte ones, and 
arrays of 1024 a ninth
+/// and a third.
+const MIN_BULK_NULLABLE_APPEND_ELEMENTS: usize = 64;

Review Comment:
   The cutoff was measured on macOS with the system allocator. Comet mostly 
runs on Linux, behind `AccountingAllocator`. At 64 elements the `i64` margin is 
about 14% (198 vs 170 ns), so the crossover could move on a different 
allocator. The doc comment cites measurements that come from an out-of-tree 
harness, so the next person who bumps arrow or changes the allocator can't 
re-check them.
   
   Could `array_element_append` take a size parameter, say 32, 64 and 128 
elements with one null per array? That would let anyone rerun the measurement 
in-tree. A Linux number for `i64` at 64 elements would also confirm that the 
cutoff doesn't regress there.



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