morningman opened a new pull request, #68743:
URL: https://github.com/apache/doris/pull/68743

   ### What problem does this PR solve?
   
   Issue Number: None
   
   Related PR: #68338 (introduced `ScanNode.cannotBeRedispatched()`, which 
until now only a remote Doris scan answered)
   
   Problem Summary:
   
   **In short.** By default an ADBC catalog reads a scan in partitions. Each 
partition is a ticket for a remote result stream, and reading the ticket drains 
that stream. When an attempt of a query fails on an RPC error, Doris retries 
the query by dispatching the same plan again, with the same tickets. The retry 
therefore reads whatever the failed attempt left of each stream, or nothing at 
all, and the query **succeeds with rows missing**. This PR lets a connector 
scan range declare that it can be read only once. ADBC partitions declare it, 
and a plan containing them is no longer dispatched again: the query fails 
instead of returning a wrong result.
   
   **Background**
   
   - **How a partitioned read works.** An ADBC catalog uses 
`partitioned_read=auto` by default. While FE plans the scan, 
`AdbcScanPlanProvider` calls the driver's `executePartitioned()`. On a Flight 
SQL source that call is `GetFlightInfo`: the source runs the query, and each 
endpoint it returns becomes one scan range carrying a partition descriptor, 
i.e. a ticket. BE reads a range with `ConnectionReadPartition`, which is a 
`DoGet` on the ticket.
   - **A DoGet drains the stream.** On a Doris source, 
`ArrowFlightResultBlockBuffer::get_arrow_batch` pops every batch it hands out. 
After the sink closes, the drained buffer stays registered for 
`result_buffer_cancelled_interval_time` (300 s). A second `DoGet` of the same 
ticket in that window gets the rest of the stream, or an immediate end of 
stream. It gets no error either way.
   - **How a failed attempt is retried.** An attempt can fail with an 
`RpcException`: the `fetch_data` RPC to the result BE fails, or a fragment 
start RPC fails. If nothing has been sent to the MySQL client yet, 
`StmtExecutor.handleQueryWithRetry` retries the query up to 
`max_query_retry_time` (3) times. It does so by dispatching the same plan again 
under a new query id. The only thing that stops this is 
`ScanNode.cannotBeRedispatched()`. #68338 added it for remote Doris scans, 
whose ranges point at a session that `stop()` closes. `PluginDrivenScanNode` 
kept the default `false`.
   
   **1. The problem, and what it cost**
   
   - Any query over an ADBC catalog in partitioned mode (`auto`, which is the 
default, or `required`) can hit this. It happens when the query meets a 
transient RPC failure after its scans have read data but before its first row 
reached the client. An aggregation is the typical case.
   - The reproduction below uses a single-FE, single-BE cluster. An ADBC 
catalog points at the cluster's own Arrow Flight SQL port, the table has 
1,000,000 rows, and one `fetch_data` response is lost on purpose (with the 
debug point this PR adds):
   
   | Query (one `fetch_data` response lost) | Before | After |
   |---|---|---|
   | `SELECT count(k), sum(k), max(v)`, ADBC catalog, `partitioned_read=auto` | 
`0 / NULL / NULL`, **no error**; the audit log records `State=EOF`, 
`ScanRows=0` | `ERROR 1105: RpcException, msg: fetch result rpc failed ...` |
   | `SELECT k` (streaming), same catalog | **950,912 rows** (49,088 missing), 
no error | the same error |
   | `partitioned_read=disabled` (the range carries the statement) | retried, 
correct result | retried, correct result |
   | remote Doris catalog (`use_arrow_flight=true`) | retry refused, error | 
unchanged |
   
   - The BE log shows two Arrow Flight readers of the same remote query. The 
failed attempt's reader took 64 packets; the retry's reader took 0.
   
   **2. How this PR fixes it**
   
   - **SPI.** `ConnectorScanRange.isSingleUse()` is new and defaults to 
`false`. A range answers `true` when reading it consumes what it points at, so 
it can be read only once.
   - **ADBC.** `AdbcScanRange.isSingleUse()` answers `true` for a partition and 
`false` for a statement, because a statement range runs its query again every 
time it is read. Whether some other source lets a ticket be read twice is 
nothing the connector can ask, so every partition counts as single-use.
   - **Engine.** `PluginDrivenScanNode` turns ranges into splits in three 
places: `getSplits`, the partition-batch generation, and the streaming 
batch-mode generation. All three now go through `toSplit()`, which records a 
single-use range in a volatile flag (the batch paths run on other threads). 
`cannotBeRedispatched()` returns that flag. `handleQueryWithRetry` already 
consults it, so such a query now fails instead of being dispatched again.
   - **Docs.** The `ScanNode.cannotBeRedispatched()` contract and the retry log 
line now cover both reasons a plan cannot be dispatched again: `stop()` 
released what the ranges point at, or reading them consumed it.
   - **Debug point.** `ResultReceiver.getNext.dropDataBatch` loses the first 
`fetch_data` response that carries rows and reports `THRIFT_RPC_ERROR`, the 
status a failed `fetch_data` RPC maps to. A test can thus drive the retry after 
the scans have read their input.
   
   What it buys: a partitioned read never returns a silently incomplete result 
after a retry, and statement ranges keep the retry.
   
   The trade-off: a partitioned ADBC query that meets a transient RPC error now 
fails instead of being retried. Retrying with a fresh plan, which would fetch 
fresh tickets, is not done here. The existing re-plan path is reserved for 
cloud errors and is chosen by matching the error message, so it is not a sound 
place to add this case.
   
   **3. The classes, and how they call each other**
   
   - `ConnectorScanRange` (fe-connector-spi): `isSingleUse()`, default `false`.
   - `AdbcScanRange` (fe-connector-adbc): `isSingleUse()` is `true` when the 
range carries a partition descriptor.
   - `PluginDrivenScanNode` (fe-core): `toSplit(range)` sets 
`plannedSingleUseRange`, and `cannotBeRedispatched()` returns it.
   - `StmtExecutor.handleQueryWithRetry` / `planCannotBeRedispatched()`: 
unchanged logic; comment and log text updated.
   - `ResultReceiver.getNext`: the debug point.
   
   ```
   planning (FE)                                          a failed attempt (FE)
   PluginDrivenScanNode.getSplits / startSplit            
QueryProcessor.getNext <- ResultReceiver (fetch_data lost: THRIFT_RPC_ERROR)
     -> AdbcScanPlanProvider.planScan                       -> RpcException
          -> executePartitioned()   (the source runs it)  
StmtExecutor.handleQueryWithRetry
          -> AdbcScanRange{partition_descriptor} ...        -> 
planCannotBeRedispatched()
     -> toSplit(range)                                           -> 
ScanNode.cannotBeRedispatched()
          range.isSingleUse() -> plannedSingleUseRange              
PluginDrivenScanNode: plannedSingleUseRange
                                                                    
RemoteDorisScanNode:  session closed by stop()
                                                               true  -> the 
query fails
                                                               false -> the 
same plan is dispatched again
   ```
   
   Untouched: BE, the ADBC reader, `AdbcScanPlanProvider`, and the retry loop 
itself.
   
   ### Release note
   
   A query over an ADBC catalog that reads in partitions 
(`partitioned_read=auto`, the default, or `required`) now fails when one of its 
attempts fails on an RPC error. It is no longer retried, so it can no longer 
return incomplete results.
   
   ### Check List (For Author)
   
   - Test <!-- At least one of them must be included. -->
       - [x] Regression test
       - [x] Unit Test
       - [x] Manual test (add detailed scripts or steps below)
       - [ ] No need to test or manual test. Explain why:
           - [ ] This is a refactor/code format and no logic has been changed.
           - [ ] Previous test can cover this change.
           - [ ] No code files have been changed.
           - [ ] Other reason <!-- Add your reason?  -->
   
       **Regression test:** 
`external_table_p0/adbc/test_adbc_partitioned_read_retry` (new; in the 
nonConcurrent group because debug points are global). It reads a 100,000-row 
table through two catalogs, one with `partitioned_read=required` and one with 
`disabled`:
       - Both catalogs return every row.
       - With one injected failure, the statement read is retried and correct.
       - The following partitioned read is undisturbed, which proves the 
injection was spent on the statement read.
       - With one more injected failure, the partitioned read fails.
   
       Run against the build without this fix, the suite fails where it should: 
`Expect exception msg contains 'fetch result rpc failed', but meet 'null'`. On 
the fixed build, the `.out` was generated with `-forceGenOut`, and a normal run 
then passed.
   
       **Unit tests:** `PluginDrivenScanNodeRedispatchTest` (new) and 
`AdbcScanRangeTest` (+1). The unchanged neighbours 
`PluginDrivenScanNodeExplainStatsTest`, `PluginDrivenScanNodeBatchModeTest`, 
`PluginDrivenScanNodeScanProviderSelectionTest`, `RemoteDorisScanNodeTest` and 
`AdbcScanPlanProviderTest` also pass: 61/61, checkstyle 0.
   
       **Manual test:** the reproduction in the table above, before and after 
the fix. It used a local single-FE, single-BE cluster and the ADBC Flight SQL 
driver built from apache-arrow-adbc-24.
   
   - Behavior changed:
       - [ ] No.
       - [x] Yes. A query over an ADBC catalog in partitioned mode that meets a 
transient RPC error now fails instead of being retried. A retry could only 
return incomplete results.
   
   - Does this need documentation?
       - [x] No.
       - [ ] Yes. <!-- Add document PR link here. eg: 
https://github.com/apache/doris-website/pull/1214 -->
   
   ### Check List (For Reviewer who merge this PR)
   
   - [ ] Confirm the release note
   - [ ] Confirm test cases
   - [ ] Confirm document
   - [ ] Add branch pick label <!-- Add branch pick label that this PR should 
merge into -->
   


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