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]

Reply via email to