This is an automated email from the ASF dual-hosted git repository.
morningman pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 1ef44c86c10 branch-4.1: [fix](fe) Fix INSERT INTO local TVF ignoring
backend_id during scheduling #61732 (#61735)
1ef44c86c10 is described below
commit 1ef44c86c109953b33fe484f4cac8cd397510236
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Thu Mar 26 22:14:17 2026 -0700
branch-4.1: [fix](fe) Fix INSERT INTO local TVF ignoring backend_id during
scheduling #61732 (#61735)
Cherry-picked from #61732
Co-authored-by: Mingyu Chen (Rayner) <[email protected]>
---
.../java/org/apache/doris/planner/TVFTableSink.java | 16 ++++++++++++++++
.../main/java/org/apache/doris/qe/Coordinator.java | 19 +++++++++++++++++++
2 files changed, 35 insertions(+)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/planner/TVFTableSink.java
b/fe/fe-core/src/main/java/org/apache/doris/planner/TVFTableSink.java
index c511336767c..b8c184ca401 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/planner/TVFTableSink.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/planner/TVFTableSink.java
@@ -60,6 +60,22 @@ public class TVFTableSink extends DataSink {
this.cols = cols;
}
+ public String getTvfName() {
+ return tvfName;
+ }
+
+ /**
+ * Returns the backend_id specified in properties, or -1 if not set.
+ * For local TVF, this indicates the specific BE node where data should be
written.
+ */
+ public long getBackendId() {
+ String backendIdStr = properties.get("backend_id");
+ if (backendIdStr != null) {
+ return Long.parseLong(backendIdStr);
+ }
+ return -1;
+ }
+
public void bindDataSink() throws AnalysisException {
TTVFTableSink tSink = new TTVFTableSink();
tSink.setTvfName(tvfName);
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
b/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
index 8c924eace5a..c2227b280b8 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
@@ -72,6 +72,7 @@ import org.apache.doris.planner.ScanNode;
import org.apache.doris.planner.SchemaScanNode;
import org.apache.doris.planner.SetOperationNode;
import org.apache.doris.planner.SortNode;
+import org.apache.doris.planner.TVFTableSink;
import org.apache.doris.planner.UnionNode;
import org.apache.doris.proto.InternalService;
import org.apache.doris.proto.InternalService.PExecPlanFragmentResult;
@@ -1767,6 +1768,24 @@ public class Coordinator implements CoordInterface {
// TODO: rethink the whole function logic. could All BE sink
naturally merged into other judgements?
return;
}
+ // For local TVF sink with a specific backend_id, we must execute
the sink fragment
+ // on the designated backend. Otherwise, data would be written to
the wrong node's local disk.
+ if (fragment.getSink() instanceof TVFTableSink) {
+ TVFTableSink tvfSink = (TVFTableSink) fragment.getSink();
+ if ("local".equals(tvfSink.getTvfName()) &&
tvfSink.getBackendId() != -1) {
+ Backend targetBackend =
Env.getCurrentSystemInfo().getBackend(tvfSink.getBackendId());
+ if (targetBackend == null || !targetBackend.isAlive()) {
+ throw new UserException("Backend " +
tvfSink.getBackendId()
+ + " is not available for local TVF sink");
+ }
+ TNetworkAddress execHostport = new TNetworkAddress(
+ targetBackend.getHost(),
targetBackend.getBePort());
+ this.addressToBackendID.put(execHostport,
targetBackend.getId());
+ FInstanceExecParam instanceParam = new
FInstanceExecParam(null, execHostport, params);
+ params.instanceExecParams.add(instanceParam);
+ continue;
+ }
+ }
if (fragment.getDataPartition() == DataPartition.UNPARTITIONED) {
Reference<Long> backendIdRef = new Reference<Long>();
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]