This is an automated email from the ASF dual-hosted git repository.
morningman pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new cfb47188bf6 [fix](fe) Do not retry a query with the same plan after an
ADBC partitioned read (#68743)
cfb47188bf6 is described below
commit cfb47188bf644fc9994497007f7db015f9135b92
Author: Mingyu Chen (Rayner) <[email protected]>
AuthorDate: Thu Oct 8 17:57:44 2026 +0800
[fix](fe) Do not retry a query with the same plan after an ADBC partitioned
read (#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. Only a
session's own query can take a hit: an internal query, such as
auto-analyze fetching its rows at the same time, is filtered out before
the hit is spent.
- **CI.** The suite runs in the External regression pipeline, the only
one that provisions the ADBC drivers, but that pipeline's FE ran with
`enable_debug_points=false`. Its `fe.conf` now enables them. The only
other `external` suites that use FE debug points,
`test_jdbc_refresh_catalog_schema_refresh_non_blocking` and
`test_jdbc_refresh_catalog_table_names_refresh_non_blocking`, are
manual-validation cases that a guard on that setting used to skip. They
are now tagged `nonConcurrent`, which the framework requires of a suite
using debug points, and run only when
`enableJdbcRefreshNonBlockingTest=true`, which no pipeline sets. They
stay skipped, as before.
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, which only a session's own
(non-internal) query can take.
```
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.
---
.../apache/doris/connector/adbc/AdbcScanRange.java | 14 ++
.../doris/connector/adbc/AdbcScanRangeTest.java | 14 ++
.../connector/spi/scan/ConnectorScanRange.java | 15 +++
.../datasource/scan/PluginDrivenScanNode.java | 30 ++++-
.../java/org/apache/doris/planner/ScanNode.java | 15 ++-
.../java/org/apache/doris/qe/ResultReceiver.java | 18 +++
.../java/org/apache/doris/qe/StmtExecutor.java | 18 +--
.../scan/PluginDrivenScanNodeRedispatchTest.java | 80 ++++++++++++
.../doris/qe/ResultReceiverDebugPointTest.java | 115 +++++++++++++++++
.../adbc/test_adbc_partitioned_read_retry.out | 13 ++
regression-test/pipeline/external/conf/fe.conf | 3 +
.../adbc/test_adbc_partitioned_read_retry.groovy | 141 +++++++++++++++++++++
...resh_catalog_schema_refresh_non_blocking.groovy | 11 +-
...catalog_table_names_refresh_non_blocking.groovy | 11 +-
14 files changed, 479 insertions(+), 19 deletions(-)
diff --git
a/fe/fe-connector/fe-connector-adbc/src/main/java/org/apache/doris/connector/adbc/AdbcScanRange.java
b/fe/fe-connector/fe-connector-adbc/src/main/java/org/apache/doris/connector/adbc/AdbcScanRange.java
index a52d0b307fa..9dfc77fd01f 100644
---
a/fe/fe-connector/fe-connector-adbc/src/main/java/org/apache/doris/connector/adbc/AdbcScanRange.java
+++
b/fe/fe-connector/fe-connector-adbc/src/main/java/org/apache/doris/connector/adbc/AdbcScanRange.java
@@ -93,6 +93,20 @@ public class AdbcScanRange implements ConnectorScanRange {
return properties;
}
+ /**
+ * True for a partition, false for a statement.
+ *
+ * <p>A partition is a ticket for one result stream of a query the source
executed while FE planned the
+ * scan, and BE reads it by draining that stream: a Doris source hands
each batch out once, so a second
+ * read of the same ticket returns what the first one left, or nothing.
Whether another source lets a
+ * ticket be read twice is its own business, which nothing here can ask,
so every partition counts as
+ * single-use. A statement runs its query when it is read, so reading the
range again runs it again.
+ */
+ @Override
+ public boolean isSingleUse() {
+ return properties.containsKey(PARAM_PARTITION_DESCRIPTOR);
+ }
+
/**
* Writes the parameters into the ADBC slot of the range descriptor.
*
diff --git
a/fe/fe-connector/fe-connector-adbc/src/test/java/org/apache/doris/connector/adbc/AdbcScanRangeTest.java
b/fe/fe-connector/fe-connector-adbc/src/test/java/org/apache/doris/connector/adbc/AdbcScanRangeTest.java
index c25c2d0f4f2..b2db701e28f 100644
---
a/fe/fe-connector/fe-connector-adbc/src/test/java/org/apache/doris/connector/adbc/AdbcScanRangeTest.java
+++
b/fe/fe-connector/fe-connector-adbc/src/test/java/org/apache/doris/connector/adbc/AdbcScanRangeTest.java
@@ -106,6 +106,20 @@ class AdbcScanRangeTest {
Assertions.assertFalse(params.containsKey("query_sql"));
}
+ @Test
+ void partitionsCanBeReadOnlyOnceWhileStatementsRunOnEveryRead() {
+ // A partition is a ticket for a result stream the source produced
once, and reading it drains the
+ // stream: the engine must not dispatch the plan again on a retry, or
the retry reads nothing and
+ // the query succeeds with no rows. A statement is run again by every
read of its range.
+ Assertions.assertTrue(new AdbcScanRange.Builder()
+
.driverPath("/opt/doris/plugins/adbc_drivers/libadbc_driver_flightsql.so")
+ .uri("grpc://remote:9090")
+ .partitionDescriptor("Zm9vYmFy")
+ .build()
+ .isSingleUse());
+ Assertions.assertFalse(minimal().build().isSingleUse());
+ }
+
@Test
void refusesToCarryBothKindsOfWorkOrNeither() {
// The two are alternatives, and BE rejects a range that says both or
neither. Failing while
diff --git
a/fe/fe-connector/fe-connector-spi/src/main/java/org/apache/doris/connector/spi/scan/ConnectorScanRange.java
b/fe/fe-connector/fe-connector-spi/src/main/java/org/apache/doris/connector/spi/scan/ConnectorScanRange.java
index 6778bd3ffcb..572a78ba1e0 100644
---
a/fe/fe-connector/fe-connector-spi/src/main/java/org/apache/doris/connector/spi/scan/ConnectorScanRange.java
+++
b/fe/fe-connector/fe-connector-spi/src/main/java/org/apache/doris/connector/spi/scan/ConnectorScanRange.java
@@ -179,6 +179,21 @@ public interface ConnectorScanRange extends Serializable {
return false;
}
+ /**
+ * Whether reading this range consumes what it points at, so that it can
be read only once.
+ *
+ * <p>The engine may read a planned range a second time: a query whose
attempt failed on an RPC error
+ * is retried by dispatching the same plan again, ranges included. That is
right for a range the source
+ * serves afresh on every read (a file, a statement run when it is read),
and wrong for a handle to a
+ * result the source produced once -- a partition of a remote query that
already ran, whose stream the
+ * first read drains: read again, it yields what the first read left, or
nothing, and the query
+ * succeeds with rows missing. A range answering {@code true} keeps the
engine from dispatching its plan
+ * again, so such a query fails instead. The default {@code false} keeps
the retry.</p>
+ */
+ default boolean isSingleUse() {
+ return false;
+ }
+
/**
* Populates per-range Thrift params from this scan range's data.
*
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/scan/PluginDrivenScanNode.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/scan/PluginDrivenScanNode.java
index fac36819a0b..8da3e442ef6 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/scan/PluginDrivenScanNode.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/scan/PluginDrivenScanNode.java
@@ -178,6 +178,11 @@ public class PluginDrivenScanNode extends
FileQueryScanNode {
private int nativeReadSplitNum;
private int totalReadSplitNum;
+ // Whether the connector planned a range that can be read only once
(ConnectorScanRange.isSingleUse),
+ // which keeps the plan from being dispatched again
(cannotBeRedispatched). Set by toSplit, which the
+ // asynchronous batch-mode split generation runs on other threads as well.
+ private volatile boolean plannedSingleUseRange;
+
// Populated from ConnectorScanPlanProvider.getScanNodePropertiesResult()
private ScanNodePropertiesResult cachedPropertiesResult;
private Map<String, String> scanNodeProperties;
@@ -1742,7 +1747,7 @@ public class PluginDrivenScanNode extends
FileQueryScanNode {
List<Split> splits = new ArrayList<>(ranges.size());
for (ConnectorScanRange range : ranges) {
- splits.add(new PluginDrivenSplit(range));
+ splits.add(toSplit(range));
}
// FIX-E (explain gap): accumulate the native/total scan-range counts
(for the connector
// EXPLAIN line paimonNativeReadSplits) and, under COUNT(*) pushdown,
the precomputed merged row
@@ -1831,6 +1836,25 @@ public class PluginDrivenScanNode extends
FileQueryScanNode {
return ctx.getExecutor().getParsedStmt().isExplain();
}
+ // Every range the connector plans becomes a split here, on the planning
thread or on the batch-mode
+ // split generation threads, so that the plan's single-use ranges are
known (cannotBeRedispatched).
+ private Split toSplit(ConnectorScanRange range) {
+ if (range.isSingleUse()) {
+ plannedSingleUseRange = true;
+ }
+ return new PluginDrivenSplit(range);
+ }
+
+ /**
+ * True once the connector planned a range that can be read only once
(ConnectorScanRange#isSingleUse):
+ * a partition of a remote query that already ran, which the failed
attempt may have drained. The same
+ * plan dispatched again would read only what that attempt left of it.
+ */
+ @Override
+ public boolean cannotBeRedispatched() {
+ return plannedSingleUseRange;
+ }
+
/**
* Counts the scan ranges read by BE's native (ORC/Parquet) reader (vs
JNI), via the generic
* {@link ConnectorScanRange#isNativeReadRange()} (default false). Drives
the EXPLAIN
@@ -2118,7 +2142,7 @@ public class PluginDrivenScanNode extends
FileQueryScanNode {
connectorSession,
batchRequest, batch));
List<Split> batchSplits = new
ArrayList<>(ranges.size());
for (ConnectorScanRange range : ranges) {
- batchSplits.add(new
PluginDrivenSplit(range));
+ batchSplits.add(toSplit(range));
}
if (splitAssignment.needMoreSplit()) {
splitAssignment.addToQueue(batchSplits);
@@ -2209,7 +2233,7 @@ public class PluginDrivenScanNode extends
FileQueryScanNode {
// heap stays bounded for million-file scans.
while (splitAssignment.needMoreSplit() && source.hasNext()) {
List<Split> one = new ArrayList<>(1);
- one.add(new PluginDrivenSplit(source.next()));
+ one.add(toSplit(source.next()));
splitAssignment.addToQueue(one);
}
splitAssignment.finishSchedule();
diff --git a/fe/fe-core/src/main/java/org/apache/doris/planner/ScanNode.java
b/fe/fe-core/src/main/java/org/apache/doris/planner/ScanNode.java
index d594fbbdbae..7e8c7747e81 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/planner/ScanNode.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/planner/ScanNode.java
@@ -149,12 +149,15 @@ public abstract class ScanNode extends PlanNode
implements SplitGenerator {
}
/**
- * Whether {@link #stop()} has released something the BE would need again
if the same plan were
- * dispatched once more, so that a retry of the query has to plan again
rather than reuse this
- * node's scan ranges (StmtExecutor.handleQueryWithRetry re-dispatches the
plan of a failed
- * attempt whose coordinator was cancelled, and cancel() stops the scan
nodes). A remote Doris
- * scan's ranges are the endpoints of the query its Flight SQL session ran
on the other frontend,
- * gone with the session; a batch split source has the same property but
is left as it is here.
+ * Whether this node's scan ranges cannot be read again by the same plan,
so that a query whose
+ * attempt failed must not be retried by dispatching that plan once more
+ * (StmtExecutor.handleQueryWithRetry re-dispatches the plan of a failed
attempt whose coordinator
+ * was cancelled, and cancel() stops the scan nodes). Either {@link
#stop()} released what the
+ * ranges point at -- a remote Doris scan's ranges are the endpoints of
the query its Flight SQL
+ * session ran on the other frontend, gone with the session -- or reading
them consumed it -- a
+ * connector range that can be read only once
(ConnectorScanRange#isSingleUse), such as an ADBC
+ * partition, the failed attempt may have drained. A batch split source
has the same property but
+ * is left as it is here.
*/
public boolean cannotBeRedispatched() {
return false;
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/ResultReceiver.java
b/fe/fe-core/src/main/java/org/apache/doris/qe/ResultReceiver.java
index a67a5398f29..1159a2806a1 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/ResultReceiver.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/ResultReceiver.java
@@ -18,6 +18,7 @@
package org.apache.doris.qe;
import org.apache.doris.common.Status;
+import org.apache.doris.common.util.DebugPointUtil;
import org.apache.doris.common.util.DebugUtil;
import org.apache.doris.proto.InternalService;
import org.apache.doris.proto.Types;
@@ -158,6 +159,23 @@ public class ResultReceiver {
return null;
}
+ // Test hook: lose the first response that carries rows, the
way a failed fetch_data RPC
+ // does (THRIFT_RPC_ERROR, as the ExecutionException branch
below reports it), at a point
+ // where the BE has produced the rows -- and so the query's
scans have read their input.
+ // StmtExecutor.handleQueryWithRetry then retries the query,
which is what this exercises.
+ // Only a session's own query qualifies: the point is armed
for a number of hits, which
+ // isEnable spends, and an internal query (auto-analyze, say)
fetching rows at the same
+ // time would otherwise take the hit meant for the query under
test.
+ if (pResult.hasRowBatch() && pResult.getRowBatch().size() > 0
+ && ConnectContext.get() != null &&
!ConnectContext.get().getState().isInternal()
+ &&
DebugPointUtil.isEnable("ResultReceiver.getNext.dropDataBatch")) {
+ LOG.warn("debug point
ResultReceiver.getNext.dropDataBatch: dropping packet {} of finstId={}",
+ pResult.getPacketSeq(),
DebugUtil.printId(getRealFinstId()));
+ status.updateStatus(TStatusCode.THRIFT_RPC_ERROR,
+ "fetch result rpc failed (debug point
ResultReceiver.getNext.dropDataBatch)");
+ return null;
+ }
+
packetIdx++;
isDone = pResult.getEos();
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
index 9a8ae03257f..5d5d4e6b53f 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
@@ -817,9 +817,9 @@ public class StmtExecutor {
originStmt.originStmt, context.getSqlHash(),
context.getQualifiedUser());
}
- // Whether a scan node of the current plan released, when the failed
attempt was cancelled, what
- // the BE would scan with again if handleQueryWithRetry dispatched the
same plan once more
- // (ScanNode.cannotBeRedispatched).
+ // Whether a scan node of the current plan has ranges the BE could not
read again if
+ // handleQueryWithRetry dispatched the same plan once more: released when
the failed attempt was
+ // cancelled, or consumed by its reading them
(ScanNode.cannotBeRedispatched).
private boolean planCannotBeRedispatched() {
if (planner == null) {
return false;
@@ -1293,11 +1293,13 @@ public class StmtExecutor {
}
}
if (isNeedRetry && planCannotBeRedispatched()) {
- // The failed attempt's cancel() stopped the scan nodes,
and one of them released
- // what the BE scans with: a remote Doris scan's session
on the other frontend,
- // whose query the scan ranges point at. The same plan
cannot be dispatched again.
- LOG.warn("not retrying query {} with the same plan: a scan
node released what the backend"
- + " scans with when the failed attempt was
cancelled. stmt: {}",
+ // A scan node's ranges cannot be read again: the failed
attempt's cancel() released
+ // what they point at (a remote Doris scan's session on
the other frontend, whose
+ // query the ranges are the endpoints of), or reading them
consumed it (an ADBC
+ // partition, a result stream the failed attempt may have
drained). Dispatched again,
+ // the plan would read nothing, or only what the failed
attempt left, and succeed.
+ LOG.warn("not retrying query {} with the same plan: a scan
node's ranges cannot be read"
+ + " again by the backend. stmt: {}",
DebugUtil.printId(context.queryId()),
parsedStmt.getOrigStmt().originStmt);
throw e;
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/scan/PluginDrivenScanNodeRedispatchTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/scan/PluginDrivenScanNodeRedispatchTest.java
new file mode 100644
index 00000000000..a7d73f7bb5e
--- /dev/null
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/scan/PluginDrivenScanNodeRedispatchTest.java
@@ -0,0 +1,80 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.datasource.scan;
+
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.connector.spi.scan.ConnectorScanRange;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+import java.util.Collections;
+import java.util.Map;
+import java.util.Optional;
+
+/**
+ * A plugin-driven scan whose connector planned a range that can be read only
once keeps its plan from being
+ * dispatched again ({@link PluginDrivenScanNode#cannotBeRedispatched()}),
which is what makes
+ * StmtExecutor.handleQueryWithRetry refuse to retry a failed attempt with the
same plan. An ADBC partition
+ * is such a range: a ticket for a result stream of a remote query that
already ran, which the failed attempt
+ * may have drained, so a retry reading the same tickets returns only what
that attempt left -- or nothing --
+ * and the query succeeds with rows missing.
+ *
+ * <p>Driven on a partial ({@code CALLS_REAL_METHODS}) node, as in
+ * {@code PluginDrivenScanNodeScanProviderSelectionTest}: every range the node
plans goes through
+ * {@code toSplit}, on the planning thread and on the batch-mode split
generation threads alike.</p>
+ */
+public class PluginDrivenScanNodeRedispatchTest {
+
+ private static ConnectorScanRange range(boolean singleUse) {
+ return new ConnectorScanRange() {
+ @Override
+ public Optional<String> getPath() {
+ return Optional.of("/dummyPath");
+ }
+
+ @Override
+ public Map<String, String> getProperties() {
+ return Collections.emptyMap();
+ }
+
+ @Override
+ public boolean isSingleUse() {
+ return singleUse;
+ }
+ };
+ }
+
+ @Test
+ public void planWithSingleUseRangeCannotBeRedispatched() {
+ PluginDrivenScanNode node = Mockito.mock(PluginDrivenScanNode.class,
Mockito.CALLS_REAL_METHODS);
+ Assertions.assertFalse(node.cannotBeRedispatched());
+
+ // Ranges the source serves afresh on every read (a file, a statement
run when it is read) keep the
+ // same-plan retry.
+ Deencapsulation.invoke(node, "toSplit", range(false));
+ Assertions.assertFalse(node.cannotBeRedispatched());
+
+ // One range that can be read only once is enough, wherever it comes
in the plan.
+ Deencapsulation.invoke(node, "toSplit", range(true));
+ Assertions.assertTrue(node.cannotBeRedispatched());
+ Deencapsulation.invoke(node, "toSplit", range(false));
+ Assertions.assertTrue(node.cannotBeRedispatched());
+ }
+}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/qe/ResultReceiverDebugPointTest.java
b/fe/fe-core/src/test/java/org/apache/doris/qe/ResultReceiverDebugPointTest.java
new file mode 100644
index 00000000000..9c8f2670977
--- /dev/null
+++
b/fe/fe-core/src/test/java/org/apache/doris/qe/ResultReceiverDebugPointTest.java
@@ -0,0 +1,115 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.qe;
+
+import org.apache.doris.common.Config;
+import org.apache.doris.common.Status;
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.common.util.DebugPointUtil;
+import org.apache.doris.common.util.DebugPointUtil.DebugPoint;
+import org.apache.doris.proto.InternalService;
+import org.apache.doris.proto.Types;
+import org.apache.doris.thrift.TNetworkAddress;
+import org.apache.doris.thrift.TResultBatch;
+import org.apache.doris.thrift.TStatusCode;
+import org.apache.doris.thrift.TUniqueId;
+
+import com.google.protobuf.ByteString;
+import org.apache.thrift.TSerializer;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.nio.ByteBuffer;
+import java.util.Collections;
+import java.util.concurrent.CompletableFuture;
+
+/**
+ * The test hook ResultReceiver.getNext.dropDataBatch is armed for a number of
hits, and asking whether it is
+ * enabled spends one. Only a session's own query may take a hit: an internal
query fetching rows at the same
+ * time (auto-analyze, say) has to leave it to the query it was armed for, or
the suite that armed it passes
+ * without the failure it injects, or fails on a query that never failed.
+ */
+public class ResultReceiverDebugPointTest {
+
+ private static final String DEBUG_POINT =
"ResultReceiver.getNext.dropDataBatch";
+
+ private boolean debugPointsEnabled;
+
+ @BeforeEach
+ public void setUp() {
+ debugPointsEnabled = Config.enable_debug_points;
+ Config.enable_debug_points = true;
+ DebugPoint oneHit = new DebugPoint();
+ oneHit.executeLimit = 1;
+ DebugPointUtil.addDebugPoint(DEBUG_POINT, oneHit);
+ }
+
+ @AfterEach
+ public void tearDown() {
+ DebugPointUtil.clearDebugPoints();
+ Config.enable_debug_points = debugPointsEnabled;
+ ConnectContext.remove();
+ }
+
+ @Test
+ public void internalQueryLeavesTheHitToTheQueryUnderTest() throws
Exception {
+ runAs(true);
+ Status status = new Status();
+ RowBatch batch = receiverWithOneRow().getNext(status);
+
+ Assertions.assertTrue(status.ok());
+ Assertions.assertEquals(1, batch.getBatch().getRowsSize());
+ // Still armed: the query the hit was meant for gets it.
+ Assertions.assertTrue(DebugPointUtil.isEnable(DEBUG_POINT));
+ }
+
+ @Test
+ public void sessionQueryTakesTheHit() throws Exception {
+ runAs(false);
+ Status status = new Status();
+ RowBatch batch = receiverWithOneRow().getNext(status);
+
+ Assertions.assertNull(batch);
+ Assertions.assertEquals(TStatusCode.THRIFT_RPC_ERROR,
status.getErrorCode());
+ // Armed for one hit, and that hit is spent.
+ Assertions.assertFalse(DebugPointUtil.isEnable(DEBUG_POINT));
+ }
+
+ private static void runAs(boolean internal) {
+ ConnectContext ctx = new ConnectContext();
+ ctx.getState().setInternal(internal);
+ ctx.setThreadLocalInfo();
+ }
+
+ // A receiver whose fetch_data response has already arrived: the last
packet, carrying one row.
+ private static ResultReceiver receiverWithOneRow() throws Exception {
+ TResultBatch rows = new
TResultBatch(Collections.singletonList(ByteBuffer.wrap(new byte[] {1})), false,
0);
+ InternalService.PFetchDataResult result =
InternalService.PFetchDataResult.newBuilder()
+ .setStatus(Types.PStatus.newBuilder().setStatusCode(0))
+ .setPacketSeq(0)
+ .setEos(true)
+ .setRowBatch(ByteString.copyFrom(new
TSerializer().serialize(rows)))
+ .build();
+ ResultReceiver receiver = new ResultReceiver(new TUniqueId(1, 2), new
TUniqueId(3, 4), 1L,
+ new TNetworkAddress("127.0.0.1", 8060), Long.MAX_VALUE, 1 <<
20, false);
+ Deencapsulation.setField(receiver, "fetchDataAsyncFuture",
CompletableFuture.completedFuture(result));
+ return receiver;
+ }
+}
diff --git
a/regression-test/data/external_table_p0/adbc/test_adbc_partitioned_read_retry.out
b/regression-test/data/external_table_p0/adbc/test_adbc_partitioned_read_retry.out
new file mode 100644
index 00000000000..d267115247b
--- /dev/null
+++
b/regression-test/data/external_table_p0/adbc/test_adbc_partitioned_read_retry.out
@@ -0,0 +1,13 @@
+-- This file is automatically generated. You should know what you did if you
want to edit this
+-- !partitions --
+100000 4999950000 v99999
+
+-- !statement --
+100000 4999950000 v99999
+
+-- !statement_retried --
+100000 4999950000 v99999
+
+-- !partitions_after_retry --
+100000 4999950000 v99999
+
diff --git a/regression-test/pipeline/external/conf/fe.conf
b/regression-test/pipeline/external/conf/fe.conf
index 4466998ae67..43d563e78d9 100644
--- a/regression-test/pipeline/external/conf/fe.conf
+++ b/regression-test/pipeline/external/conf/fe.conf
@@ -93,6 +93,9 @@ dynamic_partition_check_interval_seconds=3
enable_feature_binlog=true
+# enable debug points, for the nonConcurrent suites that inject failures into
FE
+enable_debug_points=true
+
auth_token = 5ff161c3-2c08-4079-b108-26c8850b6598
infodb_support_ext_catalog=true
diff --git
a/regression-test/suites/external_table_p0/adbc/test_adbc_partitioned_read_retry.groovy
b/regression-test/suites/external_table_p0/adbc/test_adbc_partitioned_read_retry.groovy
new file mode 100644
index 00000000000..32f8c41ffe5
--- /dev/null
+++
b/regression-test/suites/external_table_p0/adbc/test_adbc_partitioned_read_retry.groovy
@@ -0,0 +1,141 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+// ############################################################################
+// A query whose attempt fails on an RPC error is retried by dispatching the
+// same plan again (StmtExecutor.handleQueryWithRetry). Under a partitioned
read
+// the plan's ranges are tickets for the result streams of a remote query that
+// ran while FE planned the scan, and the attempt that failed has drained them:
+// the same tickets read again return nothing, and the query used to succeed
+// with no rows. Such a plan must not be dispatched again, so the query fails
+// instead. A statement range runs its query on every read, so that query is
+// still retried and returns every row.
+//
+// The failure is the FE debug point ResultReceiver.getNext.dropDataBatch: it
+// loses the first fetch_data response that carries rows -- by then the scans
+// have read their input -- and reports THRIFT_RPC_ERROR, the status a failed
+// fetch_data RPC maps to. Debug points are global, hence nonConcurrent.
+//
+// Setup is the same as test_adbc_catalog_scan -- see its header. In short:
+// FE and every BE must be able to read libadbc_driver_flightsql.so at the same
+// absolute path.
+// ############################################################################
+
+suite("test_adbc_partitioned_read_retry", "p0,external,nonConcurrent") {
+ String repoRoot = new
File(context.config.suitePath).getParentFile().getParentFile()
+ .getAbsolutePath()
+ String thirdparty = System.getenv("DORIS_THIRDPARTY")
+ if (thirdparty == null || thirdparty.isEmpty()) {
+ thirdparty = "${repoRoot}/thirdparty"
+ }
+ String driverPath = context.config.otherConfigs.get("adbcDriverPath")
+ if (driverPath == null || driverPath.isEmpty()) {
+ driverPath =
"${thirdparty}/installed/lib64/libadbc_driver_flightsql.so"
+ }
+
+ if (!new File(driverPath).canRead()) {
+ // Not a pass. Nothing about ADBC has been exercised by this run.
+ logger.info("SKIPPED test_adbc_partitioned_read_retry: no readable
ADBC Flight SQL driver at "
+ + "${driverPath}. Install it with 'cd thirdparty &&
./build-thirdparty.sh arrow_adbc', "
+ + "or set adbcDriverPath in regression-conf.groovy. "
+ + "THE RETRY OF AN ADBC PARTITIONED READ IS NOT BEING TESTED.")
+ return
+ }
+
+ def frontends = sql "show frontends"
+ String arrowPort = frontends[0][6]
+
+ sql """DROP CATALOG IF EXISTS
test_adbc_partitioned_read_retry_partitions"""
+ sql """DROP CATALOG IF EXISTS test_adbc_partitioned_read_retry_statement"""
+ sql """DROP DATABASE IF EXISTS test_adbc_partitioned_read_retry_db FORCE"""
+ sql """CREATE DATABASE test_adbc_partitioned_read_retry_db"""
+
+ sql """
+ CREATE TABLE test_adbc_partitioned_read_retry_db.src (
+ `k` bigint NOT NULL,
+ `v` varchar(32) NOT NULL
+ ) DISTRIBUTED BY HASH(`k`) BUCKETS 8
+ PROPERTIES ("replication_num" = "1")
+ """
+ sql """
+ INSERT INTO test_adbc_partitioned_read_retry_db.src
+ SELECT number, concat('v', number) FROM numbers("number" = "100000")
+ """
+
+ // 'required', not the default 'auto': a driver that stopped partitioning
would be downgraded to the
+ // statement path silently, and the partitioned half of this suite would
test the statement half.
+ sql """
+ CREATE CATALOG test_adbc_partitioned_read_retry_partitions PROPERTIES (
+ "type" = "adbc",
+ "driver_url" = "${driverPath}",
+ "sql_dialect" = "doris",
+ "uri" = "grpc://127.0.0.1:${arrowPort}",
+ "user" = "root",
+ "password" = "",
+ "partitioned_read" = "required"
+ )
+ """
+ sql """
+ CREATE CATALOG test_adbc_partitioned_read_retry_statement PROPERTIES (
+ "type" = "adbc",
+ "driver_url" = "${driverPath}",
+ "sql_dialect" = "doris",
+ "uri" = "grpc://127.0.0.1:${arrowPort}",
+ "user" = "root",
+ "password" = "",
+ "partitioned_read" = "disabled"
+ )
+ """
+
+ // Without a failure, both paths return every row.
+ order_qt_partitions """
+ SELECT count(k), sum(k), max(v)
+ FROM
test_adbc_partitioned_read_retry_partitions.test_adbc_partitioned_read_retry_db.src
+ """
+ order_qt_statement """
+ SELECT count(k), sum(k), max(v)
+ FROM
test_adbc_partitioned_read_retry_statement.test_adbc_partitioned_read_retry_db.src
+ """
+
+ try {
+ // A statement runs again when the retry reads its range: one response
lost, every row returned.
+
GetDebugPoint().enableDebugPointForAllFEs("ResultReceiver.getNext.dropDataBatch",
[execute: 1])
+ order_qt_statement_retried """
+ SELECT count(k), sum(k), max(v)
+ FROM
test_adbc_partitioned_read_retry_statement.test_adbc_partitioned_read_retry_db.src
+ """
+ // The response lost was that query's, which proves its retry
happened: armed for one hit, the debug
+ // point would otherwise fail this partitioned read.
+ order_qt_partitions_after_retry """
+ SELECT count(k), sum(k), max(v)
+ FROM
test_adbc_partitioned_read_retry_partitions.test_adbc_partitioned_read_retry_db.src
+ """
+
+ // A partitioned read is not retried: the attempt that failed drained
its tickets, so the query fails
+ // rather than return what is left of them -- which is nothing, a
count of 0.
+
GetDebugPoint().enableDebugPointForAllFEs("ResultReceiver.getNext.dropDataBatch",
[execute: 1])
+ test {
+ sql """
+ SELECT count(k), sum(k), max(v)
+ FROM
test_adbc_partitioned_read_retry_partitions.test_adbc_partitioned_read_retry_db.src
+ """
+ exception "fetch result rpc failed"
+ }
+ } finally {
+
GetDebugPoint().disableDebugPointForAllFEs("ResultReceiver.getNext.dropDataBatch")
+ }
+}
diff --git
a/regression-test/suites/external_table_p0/jdbc/test_jdbc_refresh_catalog_schema_refresh_non_blocking.groovy
b/regression-test/suites/external_table_p0/jdbc/test_jdbc_refresh_catalog_schema_refresh_non_blocking.groovy
index 56ea2982070..2fb58b861bb 100644
---
a/regression-test/suites/external_table_p0/jdbc/test_jdbc_refresh_catalog_schema_refresh_non_blocking.groovy
+++
b/regression-test/suites/external_table_p0/jdbc/test_jdbc_refresh_catalog_schema_refresh_non_blocking.groovy
@@ -22,11 +22,20 @@ import java.util.UUID
// default Apache Doris CI pipeline, so the case can be skipped or may not run
// end-to-end there. That is expected. The case still serves as a valuable
// reference for manual validation of the schema refresh non-blocking behavior.
-suite("test_jdbc_refresh_catalog_schema_refresh_non_blocking", "p0,external") {
+// It is therefore opt-in: set enableJdbcRefreshNonBlockingTest=true in
+// regression-conf.groovy to run it. It is nonConcurrent because it injects FE
+// debug points.
+suite("test_jdbc_refresh_catalog_schema_refresh_non_blocking",
"p0,external,nonConcurrent") {
String enabled = context.config.otherConfigs.get("enableJdbcTest")
if (enabled == null || !enabled.equalsIgnoreCase("true")) {
return
}
+ // Manual validation only (see above): no regression pipeline sets this
switch.
+ String manualEnabled =
context.config.otherConfigs.get("enableJdbcRefreshNonBlockingTest")
+ if (manualEnabled == null || !manualEnabled.equalsIgnoreCase("true")) {
+ logger.info("skip: enableJdbcRefreshNonBlockingTest is not true
(manual validation case)")
+ return
+ }
String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
String mysqlPort = context.config.otherConfigs.get("mysql_57_port")
diff --git
a/regression-test/suites/external_table_p0/jdbc/test_jdbc_refresh_catalog_table_names_refresh_non_blocking.groovy
b/regression-test/suites/external_table_p0/jdbc/test_jdbc_refresh_catalog_table_names_refresh_non_blocking.groovy
index fe1e5b9b390..8e87ef5a098 100644
---
a/regression-test/suites/external_table_p0/jdbc/test_jdbc_refresh_catalog_table_names_refresh_non_blocking.groovy
+++
b/regression-test/suites/external_table_p0/jdbc/test_jdbc_refresh_catalog_table_names_refresh_non_blocking.groovy
@@ -22,11 +22,20 @@ import java.util.UUID
// default Apache Doris CI pipeline, so the case can be skipped or may not run
// end-to-end there. That is expected. The case still serves as a valuable
// reference for manual validation of the table-names refresh non-blocking
behavior.
-suite("test_jdbc_refresh_catalog_table_names_refresh_non_blocking",
"p0,external") {
+// It is therefore opt-in: set enableJdbcRefreshNonBlockingTest=true in
+// regression-conf.groovy to run it. It is nonConcurrent because it injects FE
+// debug points.
+suite("test_jdbc_refresh_catalog_table_names_refresh_non_blocking",
"p0,external,nonConcurrent") {
String enabled = context.config.otherConfigs.get("enableJdbcTest")
if (enabled == null || !enabled.equalsIgnoreCase("true")) {
return
}
+ // Manual validation only (see above): no regression pipeline sets this
switch.
+ String manualEnabled =
context.config.otherConfigs.get("enableJdbcRefreshNonBlockingTest")
+ if (manualEnabled == null || !manualEnabled.equalsIgnoreCase("true")) {
+ logger.info("skip: enableJdbcRefreshNonBlockingTest is not true
(manual validation case)")
+ return
+ }
String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
String mysqlPort = context.config.otherConfigs.get("mysql_57_port")
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]