jayzhan211 opened a new issue, #26003:
URL: https://github.com/apache/datafusion/issues/26003

   ### Describe the bug
   
   Sort-merge semi / anti / mark joins (`EXISTS`, `IN`, `NOT EXISTS`, ...) fail 
on several common join key types once a key group reaches the end of an input 
batch:
   
   ```
   This feature is not implemented: Unsupported data type in sort merge join 
comparator: Timestamp(µs, "UTC")
   ```
   
   Affected key types include timezone-aware timestamps (`Timestamp(_, 
Some(tz))`, e.g. Parquet timestamps with `isAdjustedToUTC=true`), `Dictionary` 
(e.g. Hive partition columns), `Decimal256`, and `Time32` / `Time64`. An inner 
sort-merge join on the same keys works, and so does hash join. With the default 
`batch_size` (8192), any partition with more than one input batch can hit it.
   
   ### To Reproduce
   
   Default settings except forcing sort-merge join:
   
   ```sql
   SET datafusion.optimizer.prefer_hash_join = false;
   
   CREATE TABLE events AS
     SELECT arrow_cast(to_timestamp(value), 'Timestamp(Microsecond, 
Some("UTC"))') AS ts
     FROM generate_series(1, 100000);
   CREATE TABLE alerts AS
     SELECT arrow_cast(to_timestamp(value * 7), 'Timestamp(Microsecond, 
Some("UTC"))') AS ts
     FROM generate_series(1, 20000);
   
   SELECT count(*) FROM events e WHERE EXISTS (SELECT 1 FROM alerts a WHERE 
a.ts = e.ts);
   -- This feature is not implemented: Unsupported data type in sort merge join 
comparator: Timestamp(µs, "UTC")
   
   SELECT count(*) FROM events e WHERE e.ts IN (SELECT ts FROM alerts);
   -- same error
   
   SELECT count(*) FROM events e WHERE NOT EXISTS (SELECT 1 FROM alerts a WHERE 
a.ts = e.ts);
   -- same error
   
   SELECT count(*) FROM events e JOIN alerts a ON a.ts = e.ts;
   -- 14285 (inner join works)
   ```
   
   `Dictionary(Int32, Utf8)`, `Decimal256(40, 2)` and `Time64(Nanosecond)` keys 
fail the same way (with small inputs this needs e.g. `SET 
datafusion.execution.batch_size = 2` so a key group crosses a batch boundary).
   
   ### Expected behavior
   
   The same results as hash join (`SET datafusion.optimizer.prefer_hash_join = 
true`): `14285` for `EXISTS` / `IN` and `85715` for `NOT EXISTS`.
   
   ### Additional context
   
   Reproduced on current `main` (Oct 3 2026).
   
   The main merge scan compares keys with `JoinKeyComparator`, which supports 
these types, which is why inner joins work. When a key group reaches the end of 
a batch, the bitwise (semi/anti/mark) stream checks whether it continues into 
the next batch with `keys_match` in `sort_merge_join/bitwise_stream.rs`, which 
calls `compare_join_arrays` in `joins/utils.rs` for non-float keys. That 
function only handles a fixed list of types (integers, floats, strings/binary, 
`Decimal128`, `Date32/64`, timezone-naive `Timestamp`) and returns 
`not_impl_err` for everything else.
   
   #25133 was the same problem for floating-point keys; its fix (#25134) routed 
float keys in `keys_match` through `JoinKeyComparator`. Doing that for all key 
types would fix the rest.
   


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