viirya opened a new issue, #5629:
URL: https://github.com/apache/datafusion-comet/issues/5629

   ### Describe the bug
   
   `mapsort` on a sliced `MapArray` fails the query instead of sorting it:
   
   ```
   org.apache.comet.CometNativeException: Invalid argument error: Max offset of 
40 exceeds length of entries 30
        at org.apache.comet.Native.executePlan(Native Method)
        at 
org.apache.comet.CometExecIterator.$anonfun$getNextBatch$2(CometExecIterator.scala:205)
        ...
   ```
   
   A native OFFSET slices the batch (DataFusion's limit uses `batch.slice(skip, 
...)`),
   and Arrow keeps a sliced `MapArray`'s original entry offsets. So 
`spark_map_sort`
   receives a map whose first entry offset is nonzero.
   
   `native/spark-expr/src/map_funcs/map_sort.rs` builds its `take` indices only 
for the
   visible maps, so the sorted entries are indexed from zero, but it then 
passes the
   *input* offsets to `MapArray::try_new`:
   
   ```rust
   let sorted_entries = take(maps_arg_entries, &indices, None)?;
   ...
   let sorted_map_arr = Arc::new(MapArray::try_new(
       Arc::clone(map_field),
       maps_arg.offsets().clone(),   // original offsets, not rebased
       sorted_map_struct.clone(),
       maps_arg.nulls().cloned(),
       *is_sorted,
   )?);
   ```
   
   For two two-entry maps, slicing away the first leaves offsets `[2, 4]` while 
`take`
   produced only 2 entries, and Arrow rejects the result.
   
   Spark 4.0+ inserts `MapSort` automatically, so no explicit `mapsort` call is 
needed
   to reach this — `InsertMapSortInGroupingExpressions` (group-by on a map key) 
and
   `InsertMapSortInRepartitionExpressions` (repartition on a map key) both do 
it.
   
   ### Steps to reproduce
   
   Spark 4.0 or later:
   
   ```scala
   withParquetTable(
     (0 until 20).map(i => (i, Map(s"b${i % 5}" -> i, s"a${i % 5}" -> (i + 
1)))),
     "tbl") {
     // OFFSET below a GROUP BY on the map column
     sql("SELECT _2, count(*) FROM (SELECT * FROM tbl ORDER BY _1 LIMIT 15 
OFFSET 5) GROUP BY _2")
       .collect()
   }
   ```
   
   The plan puts the aggregate's `_groupingmapsort` key above a native offset:
   
   ```
   HashAggregate(keys=[_groupingmapsort#3400], functions=[count(1)])
   +- HashAggregate(keys=[_groupingmapsort#3400], functions=[partial_count(1)])
      +- CometTakeOrderedAndProjectExec(limit=20, offset=5, orderBy=[_1 ASC 
NULLS FIRST])
         +- CometNativeScan parquet [_1,_2]
   ```
   
   A unit-level reproduction, independent of Spark:
   
   ```rust
   let full = /* two maps of two entries each */;
   let sliced = full.slice(1, 1);
   assert_eq!(sliced.offsets().first().copied(), Some(2));
   spark_map_sort(&[ColumnarValue::Array(Arc::new(sliced))]).unwrap();
   // Arrow error: Invalid argument error: Max offset of 4 exceeds length of 
entries 2
   ```
   
   ### Expected behavior
   
   The query should sort the map entries and return the same answer as Spark, 
rather
   than failing. `spark_map_sort` should rebase the output offsets to match the 
entries
   it actually took.
   
   ### Additional context
   
   Found while reviewing #5567, which admits nested types as native shuffle hash
   partitioning keys and would add `InsertMapSortInRepartitionExpressions` as a 
second
   route to this code. The bug is not caused by that change: the group-by path 
above
   reproduces on current `main` with no configuration changes. Credit to 
@sunchao for
   spotting it in review.
   


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