This is an automated email from the ASF dual-hosted git repository.

JNSimba 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 cfb817f9c38 [fix](cdc) Use streaming job creator identity for CDC 
operations (#67851)
cfb817f9c38 is described below

commit cfb817f9c38856b4759dd0c56172acc8f69a2261
Author: wudi <[email protected]>
AuthorDate: Mon Sep 14 14:51:32 2026 +0800

    [fix](cdc) Use streaming job creator identity for CDC operations (#67851)
    
    ### What problem does this PR solve?
    
    Problem Summary:
    
    CDC schema changes and writes used a hard-coded `admin` identity instead
    of the identity of the user that created the streaming job. This
    bypasses the job creator's permission boundary and also makes CDC schema
    changes depend on the built-in admin credentials.
    
    This change adds streaming-specific schema-change and schema-inspection
    endpoints authenticated by the FE-issued cluster token. The FE resolves
    the streaming job from `jobId`, restores the job creator identity, and
    applies the existing ALTER or SHOW privilege checks. Both operations are
    accepted only by the master FE. The general query API and the
    password-authenticated table schema API keep their existing behavior.
    
    The CDC client also uses the job creator's Doris user for Stream Load
    instead of a hard-coded `admin` user.
---
 .../doris/job/cdc/request/WriteRecordRequest.java  |   1 +
 .../doris/httpv2/rest/StreamingJobAction.java      |  92 +++++++++++-
 .../doris/httpv2/rest/TableSchemaAction.java       |   7 +-
 .../insert/streaming/StreamingMultiTblTask.java    |   1 +
 .../rest/StreamingJobActionSchemaChangeTest.java   | 167 +++++++++++++++++++++
 .../cdcclient/service/PipelineCoordinator.java     |  10 +-
 .../doris/cdcclient/sink/DorisBatchStreamLoad.java |   2 +
 .../doris/cdcclient/sink/HttpPutBuilder.java       |   2 -
 .../org/apache/doris/cdcclient/utils/HttpUtil.java |   4 -
 .../doris/cdcclient/utils/SchemaChangeManager.java |  31 ++--
 .../cdcclient/itcase/CdcClientWriteHarness.java    |   1 +
 .../doris/cdcclient/itcase/MockDorisServer.java    |  16 +-
 .../cdcclient/utils/SchemaChangeManagerTest.java   |  20 ++-
 .../cdc/test_streaming_mysql_job_priv.groovy       |  41 ++++-
 14 files changed, 358 insertions(+), 37 deletions(-)

diff --git 
a/fe/fe-common/src/main/java/org/apache/doris/job/cdc/request/WriteRecordRequest.java
 
b/fe/fe-common/src/main/java/org/apache/doris/job/cdc/request/WriteRecordRequest.java
index 037ae137763..f397b0e4318 100644
--- 
a/fe/fe-common/src/main/java/org/apache/doris/job/cdc/request/WriteRecordRequest.java
+++ 
b/fe/fe-common/src/main/java/org/apache/doris/job/cdc/request/WriteRecordRequest.java
@@ -29,6 +29,7 @@ public class WriteRecordRequest extends JobBaseRecordRequest {
     private long taskTimeoutMs;
     private String targetDb;
     private String token;
+    private String dorisUser;
     private String taskId;
     private Map<String, String> streamLoadProps;
     // previous task ended abnormally, rebuild reader instead of reusing
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/StreamingJobAction.java 
b/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/StreamingJobAction.java
index a27973d5350..42d243b5754 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/StreamingJobAction.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/StreamingJobAction.java
@@ -18,26 +18,58 @@
 package org.apache.doris.httpv2.rest;
 
 import org.apache.doris.catalog.Env;
+import org.apache.doris.datasource.InternalCatalog;
 import org.apache.doris.httpv2.entity.ResponseEntityBuilder;
+import org.apache.doris.httpv2.exception.BadRequestException;
 import org.apache.doris.httpv2.exception.UnauthorizedException;
 import org.apache.doris.job.base.AbstractJob;
 import org.apache.doris.job.cdc.request.CommitOffsetRequest;
 import org.apache.doris.job.cdc.request.TaskFailureRequest;
 import org.apache.doris.job.extensions.insert.streaming.StreamingInsertJob;
+import org.apache.doris.mysql.privilege.PrivPredicate;
+import org.apache.doris.nereids.parser.NereidsParser;
+import org.apache.doris.nereids.trees.plans.commands.AlterTableCommand;
+import org.apache.doris.qe.AutoCloseConnectContext;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.QueryState;
+import org.apache.doris.qe.StmtExecutor;
 
 import com.google.common.base.Strings;
 import jakarta.servlet.http.HttpServletRequest;
 import org.apache.logging.log4j.LogManager;
 import org.apache.logging.log4j.Logger;
+import org.springframework.web.bind.annotation.PathVariable;
 import org.springframework.web.bind.annotation.RequestBody;
 import org.springframework.web.bind.annotation.RequestMapping;
 import org.springframework.web.bind.annotation.RequestMethod;
 import org.springframework.web.bind.annotation.RestController;
 
+import java.util.Map;
+
 @RestController
 public class StreamingJobAction extends RestBaseController {
     private static final Logger LOG = 
LogManager.getLogger(StreamingJobAction.class);
 
+    private final TableSchemaAction tableSchemaAction;
+
+    public StreamingJobAction(TableSchemaAction tableSchemaAction) {
+        this.tableSchemaAction = tableSchemaAction;
+    }
+
+    @RequestMapping(path = "/api/streaming/schema/{" + DB_KEY + "}/{" + 
TABLE_KEY + "}",
+            method = RequestMethod.GET)
+    public Object getTableSchema(@PathVariable(value = DB_KEY) String dbName,
+            @PathVariable(value = TABLE_KEY) String tblName, 
HttpServletRequest request) {
+        checkAuth(request);
+        if (!Env.getCurrentEnv().isMaster()) {
+            return ResponseEntityBuilder.okWithCommonError("Table schema must 
be queried on the master FE");
+        }
+        try (AutoCloseConnectContext ignored = new 
AutoCloseConnectContext(createJobContext(request))) {
+            return tableSchemaAction.getSchema(
+                    InternalCatalog.INTERNAL_CATALOG_NAME, dbName, tblName, 
PrivPredicate.SHOW);
+        }
+    }
+
     @RequestMapping(path = "/api/streaming/commit_offset", method = 
RequestMethod.PUT)
     public Object commitOffset(@RequestBody CommitOffsetRequest offsetRequest, 
HttpServletRequest request) {
         checkAuth(request);
@@ -50,14 +82,72 @@ public class StreamingJobAction extends RestBaseController {
         return failTask(failureRequest);
     }
 
+    @RequestMapping(path = "/api/streaming/schema_change", method = 
RequestMethod.POST)
+    public Object executeSchemaChange(@RequestBody Map<String, String> body, 
HttpServletRequest request) {
+        checkAuth(request);
+        if (!Env.getCurrentEnv().isMaster()) {
+            return ResponseEntityBuilder.okWithCommonError("Schema change must 
be executed on the master FE");
+        }
+        String stmt = body.get("stmt");
+        if (Strings.isNullOrEmpty(stmt)) {
+            return ResponseEntityBuilder.badRequest("Missing statement request 
body");
+        }
+
+        ConnectContext ctx = createJobContext(request);
+        try (AutoCloseConnectContext ignored = new 
AutoCloseConnectContext(ctx)) {
+            if (!(new NereidsParser().parseSingle(stmt) instanceof 
AlterTableCommand)) {
+                return ResponseEntityBuilder.badRequest("Only one ALTER TABLE 
statement is allowed");
+            }
+            StmtExecutor executor = new StmtExecutor(ctx, stmt);
+            executor.execute();
+            if (ctx.getState().getStateType() == 
QueryState.MysqlStateType.ERR) {
+                return 
ResponseEntityBuilder.okWithCommonError(ctx.getState().getErrorMessage());
+            }
+            return ResponseEntityBuilder.ok();
+        } catch (Exception e) {
+            LOG.warn("Failed to execute schema change", e);
+            return ResponseEntityBuilder.okWithCommonError(e.getMessage());
+        }
+    }
+
     private void checkAuth(HttpServletRequest request) {
         String authToken = request.getHeader("token");
         if (Strings.isNullOrEmpty(authToken)) {
             throw new UnauthorizedException("Miss token");
         }
         if (!checkClusterToken(authToken)) {
-            throw new UnauthorizedException("Invalid token: " + authToken);
+            throw new UnauthorizedException("Invalid token");
+        }
+    }
+
+    // Call only after validating the internal token. The caller owns the 
context's scope.
+    private static ConnectContext createJobContext(HttpServletRequest request) 
{
+        String jobIdHeader = request.getHeader("jobId");
+        if (Strings.isNullOrEmpty(jobIdHeader)) {
+            throw new BadRequestException("Missing jobId header; CDC client 
must send the streaming job ID");
+        }
+        long jobId;
+        try {
+            jobId = Long.parseLong(jobIdHeader);
+        } catch (NumberFormatException e) {
+            throw new BadRequestException("Invalid jobId header: " + 
jobIdHeader);
+        }
+        AbstractJob job = Env.getCurrentEnv().getJobManager().getJob(jobId);
+        if (!(job instanceof StreamingInsertJob)) {
+            throw new BadRequestException("Job " + jobId + " is not a 
streaming job or does not exist");
+        }
+        if (job.getCreateUser() == null) {
+            throw new BadRequestException("Streaming job " + jobId + " has no 
creator identity");
+        }
+        ConnectContext ctx = new ConnectContext();
+        ctx.setEnv(Env.getCurrentEnv());
+        ctx.setRemoteIP(request.getRemoteAddr());
+        ctx.setCurrentUserIdentity(job.getCreateUser());
+        if (!Strings.isNullOrEmpty(job.getCurrentDbName())) {
+            ctx.setDatabase(job.getCurrentDbName());
         }
+        ctx.getState().setInternal(true);
+        return ctx;
     }
 
     private Object failTask(TaskFailureRequest failureRequest) {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/TableSchemaAction.java 
b/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/TableSchemaAction.java
index 9f06ba6e115..6761e7fdc6f 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/TableSchemaAction.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/TableSchemaAction.java
@@ -111,6 +111,10 @@ public class TableSchemaAction extends RestBaseController {
             @PathVariable(value = TABLE_KEY) final String tblName,
             HttpServletRequest request, HttpServletResponse response) {
         executeCheckPassword(request, response);
+        return getSchema(catalogName, dbName, tblName, PrivPredicate.SELECT);
+    }
+
+    Object getSchema(String catalogName, String dbName, String tblName, 
PrivPredicate privilege) {
         // just allocate 2 slot for top holder map
         Map<String, Object> resultMap = new HashMap<>(2);
 
@@ -119,9 +123,8 @@ public class TableSchemaAction extends RestBaseController {
         }
 
         try {
-            // check privilege for select, otherwise return 401 HTTP status
             checkTblAuth(ConnectContext.get().getCurrentUserIdentity(), 
catalogName, dbName, tblName,
-                    PrivPredicate.SELECT);
+                    privilege);
             TableIf table;
             try {
                 CatalogIf catalog = StringUtils.isNotBlank(catalogName) ? 
Env.getCurrentEnv().getCatalogMgr()
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingMultiTblTask.java
 
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingMultiTblTask.java
index bf41d674afc..e652f7c9e1e 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingMultiTblTask.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingMultiTblTask.java
@@ -223,6 +223,7 @@ public class StreamingMultiTblTask extends 
AbstractStreamingTask {
         request.setDataSource(dataSourceType.name());
         request.setTaskId(getTaskId() + "");
         request.setToken(getToken());
+        request.setDorisUser(getUserIdentity().getQualifiedUser());
         request.setTargetDb(targetDb);
 
         Map<String, String> props = generateStreamLoadProps();
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/httpv2/rest/StreamingJobActionSchemaChangeTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/httpv2/rest/StreamingJobActionSchemaChangeTest.java
new file mode 100644
index 00000000000..90f77417096
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/httpv2/rest/StreamingJobActionSchemaChangeTest.java
@@ -0,0 +1,167 @@
+// 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.httpv2.rest;
+
+import org.apache.doris.analysis.UserIdentity;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.TokenManager;
+import org.apache.doris.httpv2.entity.ResponseBody;
+import org.apache.doris.job.extensions.insert.streaming.StreamingInsertJob;
+import org.apache.doris.qe.StmtExecutor;
+import org.apache.doris.utframe.TestWithFeService;
+
+import jakarta.servlet.http.HttpServletRequest;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedConstruction;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+import org.springframework.http.ResponseEntity;
+import org.springframework.util.AntPathMatcher;
+import org.springframework.web.bind.annotation.RequestMapping;
+
+import java.lang.reflect.Method;
+import java.nio.charset.StandardCharsets;
+import java.util.Base64;
+import java.util.Collections;
+import java.util.Map;
+
+public class StreamingJobActionSchemaChangeTest extends TestWithFeService {
+    private static final String DB_NAME = "streaming_schema_change_test";
+    private static final String TABLE_NAME = "token_auth_tbl";
+    private static final String LOAD_ONLY_USER = "streaming_schema_load_only";
+    private final StreamingJobAction action = new StreamingJobAction(new 
TableSchemaAction());
+
+    @Override
+    protected void runBeforeAll() throws Exception {
+        createDatabase(DB_NAME);
+        createTable("CREATE TABLE " + DB_NAME + "." + TABLE_NAME + " (k1 INT) "
+                + "DISTRIBUTED BY HASH(k1) BUCKETS 1 PROPERTIES 
('replication_num' = '1')");
+        StreamingInsertJob job = Mockito.mock(StreamingInsertJob.class);
+        Mockito.when(job.getJobId()).thenReturn(123L);
+        
Mockito.when(job.getCreateUser()).thenReturn(connectContext.getCurrentUserIdentity());
+        Mockito.when(job.getCurrentDbName()).thenReturn(DB_NAME);
+        Env.getCurrentEnv().getJobManager().createJobInternal(job, true);
+
+        addUser(LOAD_ONLY_USER, false);
+        grantPriv("GRANT LOAD_PRIV ON " + DB_NAME + "." + TABLE_NAME + " TO " 
+ LOAD_ONLY_USER + "@'%'");
+        StreamingInsertJob loadOnlyJob = 
Mockito.mock(StreamingInsertJob.class);
+        Mockito.when(loadOnlyJob.getJobId()).thenReturn(124L);
+        Mockito.when(loadOnlyJob.getCreateUser())
+                
.thenReturn(UserIdentity.createAnalyzedUserIdentWithIp(LOAD_ONLY_USER, "%"));
+        Mockito.when(loadOnlyJob.getCurrentDbName()).thenReturn(DB_NAME);
+        Env.getCurrentEnv().getJobManager().createJobInternal(loadOnlyJob, 
true);
+    }
+
+    @Test
+    public void testExecuteSchemaChangeWithToken() throws Exception {
+        HttpServletRequest request = tokenRequest();
+        Map<String, String> body = Collections.singletonMap(
+                "stmt", "ALTER TABLE " + DB_NAME + "." + TABLE_NAME + " ADD 
COLUMN added_col INT");
+
+        ResponseEntity<?> result = (ResponseEntity<?>) 
action.executeSchemaChange(body, request);
+
+        ResponseBody<?> responseBody = (ResponseBody<?>) result.getBody();
+        Assertions.assertEquals(RestApiStatusCode.OK.code, 
responseBody.getCode());
+    }
+
+    @Test
+    public void testExecuteSchemaChangeRejectsDropTable() throws Exception {
+        Map<String, String> body = Collections.singletonMap(
+                "stmt", "DROP TABLE " + DB_NAME + "." + TABLE_NAME);
+
+        ResponseEntity<?> result = (ResponseEntity<?>) 
action.executeSchemaChange(body, tokenRequest());
+
+        ResponseBody<?> responseBody = (ResponseBody<?>) result.getBody();
+        Assertions.assertEquals(RestApiStatusCode.BAD_REQUEST.code, 
responseBody.getCode());
+        Assertions.assertNotNull(Env.getCurrentEnv().getInternalCatalog()
+                
.getDbOrMetaException(DB_NAME).getTableOrMetaException(TABLE_NAME));
+    }
+
+    @Test
+    public void testExecuteSchemaChangeRejectsMultipleStatements() throws 
Exception {
+        Map<String, String> body = Collections.singletonMap("stmt",
+                "ALTER TABLE " + DB_NAME + "." + TABLE_NAME + " ADD COLUMN 
rejected_col INT; "
+                        + "DROP TABLE " + DB_NAME + "." + TABLE_NAME);
+
+        ResponseEntity<?> result = (ResponseEntity<?>) 
action.executeSchemaChange(body, tokenRequest());
+
+        ResponseBody<?> responseBody = (ResponseBody<?>) result.getBody();
+        Assertions.assertEquals(RestApiStatusCode.COMMON_ERROR.code, 
responseBody.getCode());
+        Assertions.assertNull(Env.getCurrentEnv().getInternalCatalog()
+                
.getDbOrMetaException(DB_NAME).getTableOrMetaException(TABLE_NAME).getColumn("rejected_col"));
+    }
+
+    @Test
+    public void testExecuteSchemaChangeRejectsNonMaster() throws Exception {
+        HttpServletRequest request = tokenRequest();
+        Map<String, String> body = Collections.singletonMap("stmt", "ALTER 
TABLE forwarded_tbl ADD COLUMN k2 INT");
+        Env follower = Mockito.mock(Env.class);
+        TokenManager tokenManager = Mockito.mock(TokenManager.class);
+        Mockito.when(follower.isMaster()).thenReturn(false);
+        
Mockito.when(tokenManager.checkAuthToken(Mockito.anyString())).thenReturn(true);
+        Mockito.when(follower.getTokenManager()).thenReturn(tokenManager);
+        try (MockedStatic<Env> env = Mockito.mockStatic(Env.class);
+                MockedConstruction<StmtExecutor> executors = 
Mockito.mockConstruction(StmtExecutor.class)) {
+            env.when(Env::getCurrentEnv).thenReturn(follower);
+            ResponseEntity<?> result = (ResponseEntity<?>) 
action.executeSchemaChange(body, request);
+
+            ResponseBody<?> responseBody = (ResponseBody<?>) result.getBody();
+            Assertions.assertEquals(RestApiStatusCode.COMMON_ERROR.code, 
responseBody.getCode());
+            Assertions.assertEquals("Schema change must be executed on the 
master FE", responseBody.getData());
+            Assertions.assertTrue(executors.constructed().isEmpty());
+        }
+    }
+
+    @Test
+    public void testGetTableSchemaWithToken() throws Exception {
+        HttpServletRequest request = tokenRequest(124L);
+        ResponseEntity<?> result = (ResponseEntity<?>) 
action.getTableSchema(DB_NAME, TABLE_NAME, request);
+
+        ResponseBody<?> responseBody = (ResponseBody<?>) result.getBody();
+        Assertions.assertEquals(RestApiStatusCode.OK.code, 
responseBody.getCode());
+        Assertions.assertEquals(200, ((Map<?, ?>) 
responseBody.getData()).get("status"));
+    }
+
+    @Test
+    public void testStreamingSchemaRouteDoesNotShadowCatalogRoute() throws 
Exception {
+        Method method = StreamingJobAction.class.getMethod(
+                "getTableSchema", String.class, String.class, 
HttpServletRequest.class);
+        String route = method.getAnnotation(RequestMapping.class).path()[0];
+        AntPathMatcher matcher = new AntPathMatcher();
+
+        Assertions.assertTrue(matcher.match(route, 
"/api/streaming/schema/db1/tbl1"));
+        Assertions.assertFalse(matcher.match(route, 
"/api/streaming/db1/tbl1/_schema"));
+    }
+
+    private HttpServletRequest tokenRequest() throws Exception {
+        return tokenRequest(123L);
+    }
+
+    private HttpServletRequest tokenRequest(long jobId) throws Exception {
+        HttpServletRequest request = Mockito.mock(HttpServletRequest.class);
+        Mockito.when(request.getRemoteAddr()).thenReturn("127.0.0.1");
+        
Mockito.when(request.getHeader("jobId")).thenReturn(String.valueOf(jobId));
+        Mockito.when(request.getHeader("token"))
+                
.thenReturn(Env.getCurrentEnv().getTokenManager().acquireToken());
+        String invalidBasic = Base64.getEncoder().encodeToString(
+                "admin:invalid-password".getBytes(StandardCharsets.UTF_8));
+        Mockito.when(request.getHeader("Authorization")).thenReturn("Basic " + 
invalidBasic);
+        return request;
+    }
+}
diff --git 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/service/PipelineCoordinator.java
 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/service/PipelineCoordinator.java
index 8f3de9fbf5c..9c133e1f172 100644
--- 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/service/PipelineCoordinator.java
+++ 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/service/PipelineCoordinator.java
@@ -452,6 +452,9 @@ public class PipelineCoordinator {
      * <p>Heartbeat events will carry the latest offset.
      */
     public void writeRecords(WriteRecordRequest writeRecordRequest) throws 
Exception {
+        Preconditions.checkArgument(
+                StringUtils.isNotBlank(writeRecordRequest.getDorisUser()),
+                "Missing dorisUser; FE must send the Doris job creator's user 
name");
         // Extract connection parameters up front for use throughout this 
method
         String feAddr = writeRecordRequest.getFrontendAddress();
         String targetDb = writeRecordRequest.getTargetDb();
@@ -611,7 +614,11 @@ public class PipelineCoordinator {
                             ddlCount += result.getSchemaChanges().size();
                         }
                         SchemaChangeManager.executeChanges(
-                                feAddr, targetDb, token, 
result.getSchemaChanges());
+                                feAddr,
+                                targetDb,
+                                token,
+                                writeRecordRequest.getJobId(),
+                                result.getSchemaChanges());
                         hasExecuteDDL = true;
                         
sourceReader.applySchemaChange(result.getUpdatedSchemas());
                         lastMessageIsHeartbeat = false;
@@ -780,6 +787,7 @@ public class PipelineCoordinator {
         batchStreamLoad.setCurrentTaskId(writeRecordRequest.getTaskId());
         
batchStreamLoad.setFrontendAddress(writeRecordRequest.getFrontendAddress());
         batchStreamLoad.setToken(writeRecordRequest.getToken());
+        batchStreamLoad.setDorisUser(writeRecordRequest.getDorisUser());
         batchStreamLoad.setLoadProps(writeRecordRequest.getStreamLoadProps());
         batchStreamLoad.getLoadStatistic().clear();
         return batchStreamLoad;
diff --git 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/sink/DorisBatchStreamLoad.java
 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/sink/DorisBatchStreamLoad.java
index 207813a523c..dbc68c680a9 100644
--- 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/sink/DorisBatchStreamLoad.java
+++ 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/sink/DorisBatchStreamLoad.java
@@ -100,6 +100,7 @@ public class DorisBatchStreamLoad implements Serializable {
     private String targetDb;
     private String jobId;
     @Setter private String token;
+    @Setter private String dorisUser;
     // stream load headers
     @Setter private Map<String, String> loadProps = new HashMap<>();
     @Getter private LoadStatistic loadStatistic;
@@ -404,6 +405,7 @@ public class DorisBatchStreamLoad implements Serializable {
                     .setUrl(loadUrl)
                     .addProperties(loadProps)
                     .addTokenAuth(token)
+                    .baseAuth(dorisUser, "")
                     .setLabel(finalLabel)
                     .formatJson()
                     .addCommonHeader()
diff --git 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/sink/HttpPutBuilder.java
 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/sink/HttpPutBuilder.java
index d24f61397a2..114485c99a3 100644
--- 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/sink/HttpPutBuilder.java
+++ 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/sink/HttpPutBuilder.java
@@ -18,7 +18,6 @@
 package org.apache.doris.cdcclient.sink;
 
 import org.apache.doris.cdcclient.common.Constants;
-import org.apache.doris.cdcclient.utils.HttpUtil;
 
 import org.apache.commons.codec.binary.Base64;
 import org.apache.commons.collections.MapUtils;
@@ -70,7 +69,6 @@ public class HttpPutBuilder {
     }
 
     public HttpPutBuilder addTokenAuth(String token) {
-        header.put(HttpHeaders.AUTHORIZATION, HttpUtil.getAuthHeader());
         header.put("token", token);
         return this;
     }
diff --git 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/utils/HttpUtil.java
 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/utils/HttpUtil.java
index 88c8accf65f..29800e4b66d 100644
--- 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/utils/HttpUtil.java
+++ 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/utils/HttpUtil.java
@@ -56,8 +56,4 @@ public class HttpUtil {
                 .addInterceptorLast(new RequestContent(true))
                 .build();
     }
-
-    public static String getAuthHeader() {
-        return "Basic YWRtaW46";
-    }
 }
diff --git 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/utils/SchemaChangeManager.java
 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/utils/SchemaChangeManager.java
index 1f74e548f0d..9480c1a873a 100644
--- 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/utils/SchemaChangeManager.java
+++ 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/utils/SchemaChangeManager.java
@@ -38,8 +38,8 @@ import org.slf4j.LoggerFactory;
 public class SchemaChangeManager {
 
     private static final Logger LOG = 
LoggerFactory.getLogger(SchemaChangeManager.class);
-    private static final String SCHEMA_CHANGE_API = 
"http://%s/api/query/default_cluster/%s";;
-    private static final String TABLE_SCHEMA_API = 
"http://%s/api/%s/%s/_schema";;
+    private static final String SCHEMA_CHANGE_API = 
"http://%s/api/streaming/schema_change";;
+    private static final String TABLE_SCHEMA_API = 
"http://%s/api/streaming/schema/%s/%s";;
     private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
     private static final String COLUMN_EXISTS_MSG = "Can not add column which 
already exists";
     private static final String COLUMN_NOT_EXISTS_MSG = "Column does not 
exists";
@@ -56,10 +56,15 @@ public class SchemaChangeManager {
      * @param feAddr Doris FE address (host:port)
      * @param db target database
      * @param token FE auth token
+     * @param jobId streaming job ID used by FE to resolve the creator identity
      * @param schemaChanges schema changes to execute
      */
     public static void executeChanges(
-            String feAddr, String db, String token, 
List<SchemaChangeOperation> schemaChanges)
+            String feAddr,
+            String db,
+            String token,
+            String jobId,
+            List<SchemaChangeOperation> schemaChanges)
             throws IOException {
         if (schemaChanges == null || schemaChanges.isEmpty()) {
             LOG.info("No DDL statements to execute");
@@ -67,27 +72,27 @@ public class SchemaChangeManager {
         }
         for (SchemaChangeOperation operation : schemaChanges) {
             LOG.info("Executing DDL on FE {}: {}", feAddr, operation.getSql());
-            execute(feAddr, db, token, operation);
+            execute(feAddr, db, token, jobId, operation);
         }
     }
 
     /**
-     * Execute a single SQL statement via the FE query API.
+     * Execute a single SQL statement via the FE streaming schema change API.
      *
      * <p>Known idempotent errors are swallowed directly. For other failures, 
the current Doris
      * schema is checked before the failure is propagated.
      */
     public static void execute(
-            String feAddr, String db, String token, SchemaChangeOperation 
operation)
+            String feAddr, String db, String token, String jobId, 
SchemaChangeOperation operation)
             throws IOException {
-        HttpPost post = buildHttpPost(feAddr, db, token, operation.getSql());
+        HttpPost post = buildHttpPost(feAddr, token, jobId, 
operation.getSql());
         try {
             String responseBody = handleResponse(post);
             LOG.info("Executed DDL {} with response: {}", operation.getSql(), 
responseBody);
             parseResponse(operation, responseBody);
         } catch (Exception ddlFailure) {
             try {
-                if (isAlreadyApplied(feAddr, db, token, operation)) {
+                if (isAlreadyApplied(feAddr, db, token, jobId, operation)) {
                     LOG.warn(
                             "[DDL-IDEMPOTENT] Doris schema already reflects {} 
{}. SQL: {}",
                             operation.getType(),
@@ -104,17 +109,17 @@ public class SchemaChangeManager {
 
     // ─── Internal helpers 
─────────────────────────────────────────────────────
 
-    private static HttpPost buildHttpPost(String feAddr, String db, String 
token, String sql)
+    private static HttpPost buildHttpPost(String feAddr, String token, String 
jobId, String sql)
             throws IOException {
-        String url = String.format(SCHEMA_CHANGE_API, feAddr, db);
+        String url = String.format(SCHEMA_CHANGE_API, feAddr);
         Map<String, Object> bodyMap = new HashMap<>();
         bodyMap.put("stmt", sql);
         String body = OBJECT_MAPPER.writeValueAsString(bodyMap);
 
         HttpPost post = new HttpPost(url);
         post.setHeader("Content-Type", "application/json;charset=UTF-8");
-        post.setHeader("Authorization", HttpUtil.getAuthHeader());
         post.setHeader("token", token);
+        post.setHeader("jobId", jobId);
         post.setEntity(new StringEntity(body, "UTF-8"));
         return post;
     }
@@ -130,12 +135,12 @@ public class SchemaChangeManager {
     }
 
     private static boolean isAlreadyApplied(
-            String feAddr, String db, String token, SchemaChangeOperation 
operation)
+            String feAddr, String db, String token, String jobId, 
SchemaChangeOperation operation)
             throws IOException {
         String url = String.format(TABLE_SCHEMA_API, feAddr, db, 
operation.getTableName());
         HttpGet request = new HttpGet(url);
-        request.setHeader("Authorization", HttpUtil.getAuthHeader());
         request.setHeader("token", token);
+        request.setHeader("jobId", jobId);
 
         String responseBody;
         try (CloseableHttpClient client = HttpUtil.getHttpClient();
diff --git 
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/CdcClientWriteHarness.java
 
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/CdcClientWriteHarness.java
index 6a93bca540a..e381c7d9b79 100644
--- 
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/CdcClientWriteHarness.java
+++ 
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/CdcClientWriteHarness.java
@@ -449,6 +449,7 @@ final class CdcClientWriteHarness implements AutoCloseable {
         req.setTaskId(String.valueOf(taskSeq.incrementAndGet()));
         req.setTargetDb(targetDb);
         req.setToken("test-token");
+        req.setDorisUser("cdc_job_user");
         req.setMaxInterval(3);
         req.setTaskTimeoutMs(60_000);
         req.setRebuildReader(rebuildReaderOnNextWrite);
diff --git 
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/MockDorisServer.java
 
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/MockDorisServer.java
index ba6898a9c49..e8a2ae8accf 100644
--- 
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/MockDorisServer.java
+++ 
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/MockDorisServer.java
@@ -38,14 +38,16 @@ import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicInteger;
 
 /**
- * A tiny in-process stand-in for the Doris BE stream-load endpoint and the FE 
commit-offset
- * endpoint, so the from-to {@code writeRecords} path can be exercised without 
a real Doris cluster.
+ * A tiny in-process stand-in for the Doris BE stream-load endpoint and FE 
streaming endpoints, so
+ * the from-to {@code writeRecords} path can be exercised without a real Doris 
cluster.
  *
  * <ul>
  *   <li>{@code PUT /api/{db}/{table}/_stream_load} — captures the 
newline-delimited JSON rows and
  *       replies with a Success stream-load result.
- *   <li>{@code POST /api/streaming/commit_offset} — captures the committed 
offset payload and
+ *   <li>{@code PUT /api/streaming/commit_offset} — captures the committed 
offset payload and
  *       replies {@code {"code":0}}.
+ *   <li>{@code POST /api/streaming/schema_change} — executes schema changes 
against mock metadata.
+ *   <li>{@code GET /api/streaming/schema/{db}/{table}} — returns mock table 
metadata.
  * </ul>
  */
 final class MockDorisServer implements AutoCloseable {
@@ -98,7 +100,7 @@ final class MockDorisServer implements AutoCloseable {
                 this.committedOffset = new String(body, 
StandardCharsets.UTF_8);
                 response = "{\"code\":0,\"msg\":\"ok\"}";
             }
-        } else if (path.endsWith("/_schema")) {
+        } else if (path.startsWith("/api/streaming/schema/")) {
             schemaRequestCount.incrementAndGet();
             List<String> properties = new ArrayList<>();
             synchronized (schemaColumns) {
@@ -110,7 +112,7 @@ final class MockDorisServer implements AutoCloseable {
                     "{\"code\":0,\"data\":{\"status\":200,\"properties\":["
                             + String.join(",", properties)
                             + "]}}";
-        } else if (path.contains("/api/query/")) {
+        } else if (path.equals("/api/streaming/schema_change")) {
             // FE schema-change endpoint: body is {"stmt":"<DDL>"}
             JsonNode node = MAPPER.readTree(body);
             executedDdls.add(node.path("stmt").asText(""));
@@ -141,6 +143,8 @@ final class MockDorisServer implements AutoCloseable {
                 response = applyDdlToMockSchema(node.path("stmt").asText(""));
             }
             ddlResponses.add(response);
+        } else if (path.equals("/api/streaming/report_task_failure")) {
+            response = "{\"code\":0,\"msg\":\"ok\"}";
         } else {
             response = "{\"code\":-1,\"msg\":\"unknown path " + path + "\"}";
         }
@@ -170,7 +174,7 @@ final class MockDorisServer implements AutoCloseable {
         return committedOffset;
     }
 
-    /** All DDL statements executed via the FE query endpoint, in arrival 
order. */
+    /** All DDL statements executed via the FE streaming endpoint, in arrival 
order. */
     List<String> executedDdls() {
         return new ArrayList<>(executedDdls);
     }
diff --git 
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/utils/SchemaChangeManagerTest.java
 
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/utils/SchemaChangeManagerTest.java
index 38d409994a8..fa343de437f 100644
--- 
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/utils/SchemaChangeManagerTest.java
+++ 
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/utils/SchemaChangeManagerTest.java
@@ -57,6 +57,7 @@ class SchemaChangeManagerTest {
                 feAddr,
                 "target_db",
                 "token",
+                "123",
                 SchemaChangeOperation.addColumn(
                         "target_table",
                         "new_col",
@@ -74,6 +75,7 @@ class SchemaChangeManagerTest {
                 feAddr,
                 "target_db",
                 "token",
+                "123",
                 SchemaChangeOperation.dropColumn(
                         "target_table",
                         "old_col",
@@ -95,7 +97,7 @@ class SchemaChangeManagerTest {
         assertThatThrownBy(
                         () ->
                                 SchemaChangeManager.execute(
-                                        feAddr, "target_db", "token", 
operation))
+                                        feAddr, "target_db", "token", "123", 
operation))
                 .isInstanceOf(IOException.class)
                 .hasMessageContaining("Failed to execute schema change");
     }
@@ -113,7 +115,7 @@ class SchemaChangeManagerTest {
         assertThatThrownBy(
                         () ->
                                 SchemaChangeManager.execute(
-                                        feAddr, "target_db", "token", 
operation))
+                                        feAddr, "target_db", "token", "123", 
operation))
                 .isInstanceOf(IOException.class)
                 .hasMessageContaining("Failed to execute schema change");
     }
@@ -122,7 +124,7 @@ class SchemaChangeManagerTest {
     void schemaQueryFailureKeepsOriginalDdlFailure() throws Exception {
         respondToDdlWithUnknownError();
         server.createContext(
-                "/api/target_db/target_table/_schema",
+                "/api/streaming/schema/target_db/target_table",
                 exchange -> respond(exchange, "{\"code\":1,\"msg\":\"schema 
unavailable\"}"));
         SchemaChangeOperation operation =
                 SchemaChangeOperation.addColumn(
@@ -133,7 +135,7 @@ class SchemaChangeManagerTest {
         assertThatThrownBy(
                         () ->
                                 SchemaChangeManager.execute(
-                                        feAddr, "target_db", "token", 
operation))
+                                        feAddr, "target_db", "token", "123", 
operation))
                 .isInstanceOf(IOException.class)
                 .hasMessageContaining("Column operation cannot be applied")
                 .satisfies(error -> 
assertThat(error.getSuppressed()).hasSize(1));
@@ -142,13 +144,14 @@ class SchemaChangeManagerTest {
     @Test
     void successfulDdlDoesNotQuerySchema() throws Exception {
         server.createContext(
-                "/api/query/default_cluster/target_db",
+                "/api/streaming/schema_change",
                 exchange -> respond(exchange, 
"{\"code\":0,\"msg\":\"success\"}"));
 
         SchemaChangeManager.execute(
                 feAddr,
                 "target_db",
                 "token",
+                "123",
                 SchemaChangeOperation.addColumn(
                         "target_table",
                         "new_col",
@@ -159,7 +162,7 @@ class SchemaChangeManagerTest {
 
     private void respondToDdlWithUnknownError() {
         server.createContext(
-                "/api/query/default_cluster/target_db",
+                "/api/streaming/schema_change",
                 exchange ->
                         respond(
                                 exchange,
@@ -168,7 +171,7 @@ class SchemaChangeManagerTest {
 
     private void respondToSchemaWithColumns(String... columns) {
         server.createContext(
-                "/api/target_db/target_table/_schema",
+                "/api/streaming/schema/target_db/target_table",
                 exchange -> {
                     schemaRequests.incrementAndGet();
                     StringBuilder properties = new StringBuilder();
@@ -187,6 +190,9 @@ class SchemaChangeManagerTest {
     }
 
     private static void respond(HttpExchange exchange, String body) throws 
IOException {
+        
assertThat(exchange.getRequestHeaders().getFirst("token")).isEqualTo("token");
+        
assertThat(exchange.getRequestHeaders().getFirst("jobId")).isEqualTo("123");
+        
assertThat(exchange.getRequestHeaders().getFirst("Authorization")).isNull();
         byte[] bytes = body.getBytes(StandardCharsets.UTF_8);
         exchange.sendResponseHeaders(200, bytes.length);
         exchange.getResponseBody().write(bytes);
diff --git 
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_priv.groovy
 
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_priv.groovy
index 540b96b0463..6100c7eee9e 100644
--- 
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_priv.groovy
+++ 
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_priv.groovy
@@ -100,7 +100,7 @@ suite("test_streaming_mysql_job_priv", 
"p0,external,mysql,external_docker,extern
             sql """INSERT INTO ${mysqlDb}.${tableName} (name, age) VALUES 
('B1', 2);"""
         }
 
-        // create streaming job by load_priv and create_priv
+        // Create the job without ALTER privilege to verify schema changes use 
the job creator's identity.
         sql """grant load_priv,create_priv on ${dbName}.* to ${user}"""
         connect(user, "${pwd}", url) {
             sql """CREATE JOB ${jobName}
@@ -141,6 +141,45 @@ suite("test_streaming_mysql_job_priv", 
"p0,external,mysql,external_docker,extern
 
         def jobResult = sql """select * from jobs("type"="insert") where 
Name='${jobName}'"""
         log.info("show jobResult: " + jobResult)
+        sql """REVOKE select_priv ON ${dbName}.* FROM ${user}"""
+
+        // A token-authenticated schema change must still check the job 
creator's ALTER privilege.
+        connect("root", "123456", 
"jdbc:mysql://${externalEnvIp}:${mysql_port}") {
+            sql """ALTER TABLE ${mysqlDb}.${tableName} ADD COLUMN cdc_auth_col 
VARCHAR(50)"""
+            sql """INSERT INTO ${mysqlDb}.${tableName} (name, age, 
cdc_auth_col)
+                    VALUES ('SchemaChangePriv', 30, 'created_by_job_user')"""
+        }
+
+        Awaitility.await().atMost(180, SECONDS).pollInterval(2, 
SECONDS).until({
+            def errors = sql """SELECT ErrorMsg FROM jobs("type"="insert") 
WHERE Name='${jobName}'"""
+            log.info("schema change privilege error: " + errors)
+            def columns = sql "DESC ${tableName}"
+            errors.size() == 1 && errors[0][0].toString().contains("ALTER 
TABLE command denied")
+                    && errors[0][0].toString().contains(user)
+                    && !columns.any { it[0] == "cdc_auth_col" }
+        })
+
+        // Let the existing automatic retry replay the failed DDL after ALTER 
is granted.
+        sql """GRANT alter_priv ON ${dbName}.* TO ${user}"""
+        Awaitility.await().atMost(180, SECONDS).pollInterval(2, 
SECONDS).until({
+            def columns = sql "DESC ${tableName}"
+            if (!columns.any { it[0] == "cdc_auth_col" }) {
+                return false
+            }
+            def rows = sql "SELECT cdc_auth_col FROM ${tableName} WHERE name = 
'SchemaChangePriv'"
+            rows.size() == 1 && rows[0][0] == "created_by_job_user"
+        })
+
+        connect("root", "123456", 
"jdbc:mysql://${externalEnvIp}:${mysql_port}") {
+            sql """DELETE FROM ${mysqlDb}.${tableName} WHERE name = 
'SchemaChangePriv'"""
+            sql """ALTER TABLE ${mysqlDb}.${tableName} DROP COLUMN 
cdc_auth_col"""
+        }
+
+        Awaitility.await().atMost(180, SECONDS).pollInterval(2, 
SECONDS).until({
+            def columns = sql "DESC ${tableName}"
+            def rows = sql "SELECT COUNT(*) FROM ${tableName} WHERE name = 
'SchemaChangePriv'"
+            !columns.any { it[0] == "cdc_auth_col" } && (rows[0][0] as int) == 0
+        })
 
         // create a new mysql user only has select priv
         def newMysqlUser = "mysql_job_priv"


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to