corgy-w opened a new pull request, #12325:
URL: https://github.com/apache/seatunnel/pull/12325
### What
The Doris source reader on the asynchronous Arrow path
(`doris.deserialize.arrow.async=true`) now fails the job when the backend fetch
fails, instead of hanging in `hasNext()` forever.
### Why
`DorisValueReader` fetches batches on a dedicated thread and the consumer
waits for either a batch or `eos`:
- `DorisValueReader.java:147-181` (before): the thread body is `try { while
(!eos.get()) { ... client.getNext(nextBatchParams) ... } } finally {
clientLock.unlock(); }`. There is no `catch`. If `client.getNext()` throws —
Doris BE restart, request timeout, scan context closed — the thread dies with
`eos` still `false` and the failure only reaches the thread's
uncaught-exception handler.
- `DorisValueReader.java:197-219`: the consumer loop is `while (!eos.get()
|| !rowBatchBlockingQueue.isEmpty())` with `Thread.sleep(5)` in the empty
branch. The queue stays empty and `eos` never flips, so the reader polls
forever: the task never completes, never fails, and never checkpoints again.
- `close()` (line 268) neither set `eos` nor interrupted the fetch thread,
and it blocks on `clientLock`, which the fetch thread holds for the whole scan
— so even cancelling the job could not break the loop.
The ClickHouse sibling gets this right: `ClickhouseValueReader.java:483-486`
sets `eos` in a `finally` block.
### Fix
- Extract the fetch body into `asyncFetchBatches()`, record any failure in a
`volatile Throwable asyncFailure`, always set `eos` in `finally`, and restore
the interrupt flag on `InterruptedException`.
- `hasNext()`: once the queue is drained, throw the recorded failure instead
of reporting an end of stream. The polling loop no longer swallows
`InterruptedException` (the same method's `take()` branch already threw on it).
- `close()`: signal `eos` and interrupt the fetch thread before taking
`clientLock`.
### How verified
- New `DorisValueReaderTest` (the first test for this class): a failing
`BackendClient.getNext` must make `hasNext()` throw inside the test timeout,
and `close()` must set `eos` and interrupt the fetch thread.
- With this change: `Tests run: 2, Failures: 0, Errors: 0, Skipped: 0` (`mvn
-o -pl seatunnel-connectors-v2/connector-doris test
-Dtest=DorisValueReaderTest`).
- Reverting only `DorisValueReader.java`:
- `testCloseSignalsTheAsyncFetchToStop` → `close must stop the
asynchronous fetch ==> expected: <true> but was: <false>`
- the failure-reporting test cannot pass (there is no recorded failure to
surface).
### Compatibility
No option name, default value, or public API change. Only the opt-in
asynchronous path is affected; the synchronous path is untouched.
--
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]