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]