This is an automated email from the ASF dual-hosted git repository.
yiguolei 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 95265876df7 branch-4.1: [improvement](fe) Improve audit logs for S3
streaming insert jobs #67489 (#67756)
95265876df7 is described below
commit 95265876df795b248e893a3072eb0a8c9b1f55b3
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Thu Sep 10 21:28:16 2026 +0800
branch-4.1: [improvement](fe) Improve audit logs for S3 streaming insert
jobs #67489 (#67756)
Cherry-picked from #67489
---------
Co-authored-by: wudi <[email protected]>
---
.../insert/streaming/StreamingInsertTask.java | 51 +++-
.../job/offset/s3/S3SourceOffsetProvider.java | 2 +-
.../streaming/StreamingInsertTaskAuditTest.java | 270 +++++++++++++++++++++
3 files changed, 321 insertions(+), 2 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTask.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTask.java
index df23f107240..da74592287c 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTask.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTask.java
@@ -20,7 +20,9 @@ package org.apache.doris.job.extensions.insert.streaming;
import org.apache.doris.analysis.UserIdentity;
import org.apache.doris.catalog.Env;
import org.apache.doris.common.Config;
+import org.apache.doris.common.ErrorCode;
import org.apache.doris.common.FeConstants;
+import org.apache.doris.common.Pair;
import org.apache.doris.common.Status;
import org.apache.doris.common.util.Util;
import org.apache.doris.job.base.Job;
@@ -30,16 +32,22 @@ import org.apache.doris.job.extensions.insert.InsertTask;
import org.apache.doris.job.offset.SourceOffsetProvider;
import org.apache.doris.load.loadv2.LoadJob;
import org.apache.doris.nereids.StatementContext;
+import org.apache.doris.nereids.analyzer.UnboundTVFRelation;
import org.apache.doris.nereids.glue.LogicalPlanAdapter;
import org.apache.doris.nereids.parser.NereidsParser;
+import org.apache.doris.nereids.trees.plans.commands.info.BaseViewInfo;
import
org.apache.doris.nereids.trees.plans.commands.insert.InsertIntoTableCommand;
+import org.apache.doris.nereids.util.SqlLiteralUtils;
+import org.apache.doris.qe.AuditLogHelper;
import org.apache.doris.qe.ConnectContext;
import org.apache.doris.qe.QueryState;
import org.apache.doris.qe.StmtExecutor;
+import org.apache.doris.tablefunction.S3TableValuedFunction;
import org.apache.doris.thrift.TCell;
import org.apache.doris.thrift.TRow;
import org.apache.doris.thrift.TStatusCode;
+import com.google.common.base.Preconditions;
import lombok.Getter;
import lombok.extern.log4j.Log4j2;
import org.apache.commons.lang3.StringUtils;
@@ -49,6 +57,8 @@ import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Optional;
+import java.util.TreeMap;
+import java.util.stream.Collectors;
@Log4j2
@Getter
@@ -61,6 +71,8 @@ public class StreamingInsertTask extends
AbstractStreamingTask {
private StreamingJobProperties jobProperties;
private Map<String, String> originTvfProps;
private String cloudCluster;
+ private String auditSql;
+ private final boolean auditEnabled;
SourceOffsetProvider offsetProvider;
public StreamingInsertTask(long jobId,
@@ -79,10 +91,12 @@ public class StreamingInsertTask extends
AbstractStreamingTask {
this.jobProperties = jobProperties;
this.originTvfProps = originTvfProps;
this.cloudCluster = cloudCluster;
+ this.auditEnabled =
S3TableValuedFunction.NAME.equalsIgnoreCase(offsetProvider.getSourceType());
}
@Override
public void before() throws Exception {
+ auditSql = null;
if (getIsCanceled().get()) {
log.info("streaming insert task has been canceled, task id is {}",
getTaskId());
return;
@@ -111,6 +125,14 @@ public class StreamingInsertTask extends
AbstractStreamingTask {
this.taskCommand = offsetProvider.rewriteTvfParams(baseCommand,
runningOffset, getTaskId());
this.taskCommand.setLabelName(Optional.of(labelName));
this.stmtExecutor = new StmtExecutor(ctx, new
LogicalPlanAdapter(taskCommand, ctx.getStatementContext()));
+ if (auditEnabled) {
+ ctx.setExecutor(stmtExecutor);
+ try {
+ this.auditSql = buildAuditSql();
+ } catch (Exception e) {
+ log.warn("Failed to prepare audit SQL, label {}; skipping
audit for this attempt", labelName, e);
+ }
+ }
}
@Override
@@ -121,6 +143,9 @@ public class StreamingInsertTask extends
AbstractStreamingTask {
}
log.info("start to run streaming insert task, label {}, offset is {}",
labelName, runningOffset.toString());
String errMsg = null;
+ if (auditEnabled) {
+ ctx.setStartTime();
+ }
try {
taskCommand.run(ctx, stmtExecutor);
if (ctx.getState().getStateType() == QueryState.MysqlStateType.OK)
{
@@ -130,12 +155,36 @@ public class StreamingInsertTask extends
AbstractStreamingTask {
}
throw new JobException(errMsg);
} catch (Exception e) {
+ String errorMessage = Util.getRootCauseMessage(e);
+ if (auditEnabled && ctx.getState().getStateType() !=
QueryState.MysqlStateType.ERR) {
+ ctx.getState().setError(ErrorCode.ERR_INTERNAL_ERROR,
errorMessage);
+ }
log.warn("execute insert task error, label is {},offset is {}",
taskCommand.getLabelName(),
runningOffset.toString(), e);
- throw new JobException(Util.getRootCauseMessage(e));
+ throw new JobException(errorMessage);
+ } finally {
+ if (auditSql != null) {
+ AuditLogHelper.logAuditLog(ctx, auditSql,
stmtExecutor.getParsedStmt(),
+ stmtExecutor.getQueryStatisticsForAuditLog(), true);
+ }
}
}
+ private String buildAuditSql() {
+ TreeMap<Pair<Integer, Integer>, String> replacements = new
TreeMap<>(new Pair.PairComparator<>());
+ new NereidsParser().parseForEncryption(sql, replacements);
+ List<UnboundTVFRelation> tvfRelations =
taskCommand.getAllTVFRelation();
+ Preconditions.checkState(replacements.size() == 1 &&
tvfRelations.size() == 1,
+ "S3 streaming insert must contain exactly one TVF");
+ String rewrittenProperties =
tvfRelations.get(0).getProperties().getMap().entrySet().stream()
+ .map(entry ->
SqlLiteralUtils.quoteStringLiteral(entry.getKey()) + " = "
+ + SqlLiteralUtils.quoteStringLiteral(entry.getValue()))
+ .collect(Collectors.joining(", "));
+ Pair<Integer, Integer> tvfPropertiesRange = replacements.firstKey();
+ replacements.replace(tvfPropertiesRange, rewrittenProperties);
+ return BaseViewInfo.rewriteSql(replacements, sql);
+ }
+
@Override
public List<Long> getScanBackendIds() {
if (stmtExecutor != null && stmtExecutor.getCoord() != null) {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3SourceOffsetProvider.java
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3SourceOffsetProvider.java
index 630b423ce9d..a3eed67b378 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3SourceOffsetProvider.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3SourceOffsetProvider.java
@@ -135,7 +135,7 @@ public class S3SourceOffsetProvider implements
SourceOffsetProvider {
public InsertIntoTableCommand rewriteTvfParams(InsertIntoTableCommand
originCommand,
Offset runningOffset, long taskId) {
S3Offset offset = (S3Offset) runningOffset;
- Map<String, String> props = new HashMap<>();
+ Map<String, String> props =
Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER);
// rewrite plan
Plan rewritePlan = originCommand.getParsedPlan().get().rewriteUp(plan
-> {
if (plan instanceof UnboundTVFRelation) {
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTaskAuditTest.java
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTaskAuditTest.java
new file mode 100644
index 00000000000..e07caa6edd1
--- /dev/null
+++
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTaskAuditTest.java
@@ -0,0 +1,270 @@
+// 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.job.extensions.insert.streaming;
+
+import org.apache.doris.analysis.StmtType;
+import org.apache.doris.analysis.UserIdentity;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.common.profile.SummaryProfile;
+import org.apache.doris.datasource.CatalogMgr;
+import org.apache.doris.datasource.InternalCatalog;
+import org.apache.doris.job.exception.JobException;
+import org.apache.doris.job.extensions.insert.InsertTask;
+import org.apache.doris.job.offset.Offset;
+import org.apache.doris.job.offset.SourceOffsetProvider;
+import org.apache.doris.job.offset.jdbc.JdbcTvfSourceOffsetProvider;
+import org.apache.doris.job.offset.s3.S3Offset;
+import org.apache.doris.job.offset.s3.S3SourceOffsetProvider;
+import org.apache.doris.nereids.StatementContext;
+import org.apache.doris.nereids.analyzer.UnboundTVFRelation;
+import org.apache.doris.nereids.glue.LogicalPlanAdapter;
+import org.apache.doris.nereids.parser.NereidsParser;
+import org.apache.doris.nereids.trees.expressions.Properties;
+import org.apache.doris.nereids.trees.plans.RelationId;
+import
org.apache.doris.nereids.trees.plans.commands.insert.InsertIntoTableCommand;
+import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
+import org.apache.doris.plugin.AuditEvent;
+import org.apache.doris.qe.AuditLogHelper;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.QueryState;
+import org.apache.doris.qe.StmtExecutor;
+import org.apache.doris.resource.workloadschedpolicy.WorkloadRuntimeStatusMgr;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.MockedConstruction;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.Optional;
+import java.util.concurrent.atomic.AtomicLong;
+
+public class StreamingInsertTaskAuditTest {
+ private static final String ORIGIN_URI = "s3://bucket/input/*.csv";
+ private static final String RESOLVED_URI =
"s3://bucket/input/{1.csv,2.csv}";
+ private static final String S3_SQL = "insert into target_table select *
from s3("
+ + "\"uri\" = \"" + ORIGIN_URI + "\", "
+ + "\"s3.secret_key\" = \"private-value\", "
+ + "\"enclose\" = \"\\\"\")";
+
+ @Test
+ public void testS3RunSubmitsAuditEvent() throws Exception {
+ AuditEvent auditEvent = runS3Task(null);
+
+ Assertions.assertEquals(AuditEvent.EventType.AFTER_QUERY,
auditEvent.type);
+ Assertions.assertEquals(StmtType.INSERT.name(), auditEvent.stmtType);
+ Assertions.assertFalse(auditEvent.stmt.contains(ORIGIN_URI));
+ Assertions.assertTrue(auditEvent.stmt.contains(RESOLVED_URI));
+ Assertions.assertFalse(auditEvent.stmt.contains("private-value"));
+ Assertions.assertEquals("OK", auditEvent.state);
+ Assertions.assertTrue(auditEvent.isInternal);
+ }
+
+ @Test
+ public void testFailedS3RunSubmitsErrorAuditEvent() throws Exception {
+ AuditEvent auditEvent = runS3Task(new RuntimeException("insert
failed"));
+
+ Assertions.assertEquals(AuditEvent.EventType.AFTER_QUERY,
auditEvent.type);
+ Assertions.assertEquals(StmtType.INSERT.name(), auditEvent.stmtType);
+ Assertions.assertEquals("ERR", auditEvent.state);
+ Assertions.assertTrue(auditEvent.errorMessage.contains("insert
failed"));
+ Assertions.assertTrue(auditEvent.isInternal);
+ }
+
+ @Test
+ public void testCdcRunDoesNotSubmitAuditEvent() throws Exception {
+ String sql = "insert into target_table select * from cdc_stream("
+ + "\"type\" = \"mysql\", \"jdbc_url\" =
\"jdbc:mysql://127.0.0.1:3306\", "
+ + "\"table\" = \"source_table\", \"offset\" = \"latest\")";
+ runTask(null, sql, Mockito.mock(JdbcTvfSourceOffsetProvider.class),
Mockito.mock(Offset.class), false);
+ }
+
+ @Test
+ public void testAuditParsingFailureDoesNotFailInsert() throws Exception {
+ runWithAuditPreparationFailure(true);
+ }
+
+ @Test
+ public void testAuditRangeFailureDoesNotFailInsert() throws Exception {
+ runWithAuditPreparationFailure(false);
+ }
+
+ private void runWithAuditPreparationFailure(boolean parserFailure) throws
Exception {
+ ConnectContext ctx = Mockito.mock(ConnectContext.class);
+ QueryState state = new QueryState();
+ Mockito.when(ctx.getState()).thenReturn(state);
+ StreamingJobProperties properties =
Mockito.mock(StreamingJobProperties.class);
+ SourceOffsetProvider provider =
Mockito.mock(SourceOffsetProvider.class);
+ Mockito.when(provider.getSourceType()).thenReturn("s3");
+ S3Offset offset = new S3Offset();
+ offset.setFileLists(RESOLVED_URI);
+ Mockito.when(provider.getNextOffset(Mockito.eq(properties),
Mockito.anyMap())).thenReturn(offset);
+ InsertIntoTableCommand baseCommand =
Mockito.mock(InsertIntoTableCommand.class);
+
Mockito.when(baseCommand.getParsedPlan()).thenReturn(Optional.of(Mockito.mock(LogicalPlan.class)));
+ InsertIntoTableCommand taskCommand =
Mockito.mock(InsertIntoTableCommand.class);
+ UnboundTVFRelation tvf = Mockito.mock(UnboundTVFRelation.class);
+
Mockito.when(taskCommand.getAllTVFRelation()).thenReturn(Collections.singletonList(tvf));
+ Mockito.when(provider.rewriteTvfParams(Mockito.eq(baseCommand),
Mockito.eq(offset), Mockito.anyLong()))
+ .thenReturn(taskCommand);
+ Mockito.doAnswer(invocation -> {
+ state.setOk();
+ return null;
+ }).when(taskCommand).run(Mockito.eq(ctx),
Mockito.any(StmtExecutor.class));
+
+ try (MockedStatic<InsertTask> insertTask =
Mockito.mockStatic(InsertTask.class);
+ MockedConstruction<StmtExecutor> executors =
Mockito.mockConstruction(StmtExecutor.class);
+ MockedConstruction<NereidsParser> parsers =
Mockito.mockConstruction(NereidsParser.class,
+ (parser, construction) -> {
+
Mockito.when(parser.parseSingle(S3_SQL)).thenReturn(baseCommand);
+ if (parserFailure) {
+
Mockito.when(parser.parseForEncryption(Mockito.eq(S3_SQL), Mockito.anyMap()))
+ .thenThrow(new
IllegalStateException("audit parsing failed"));
+ }
+ // Otherwise leave replacements empty to exercise
the range assertion.
+ });
+ MockedStatic<AuditLogHelper> audit =
Mockito.mockStatic(AuditLogHelper.class)) {
+ insertTask.when(() ->
InsertTask.makeConnectContext(UserIdentity.ROOT, "test_db")).thenReturn(ctx);
+ StreamingInsertTask task = new StreamingInsertTask(1L, 2L, S3_SQL,
provider, "test_db", properties,
+ Collections.emptyMap(), UserIdentity.ROOT, null);
+ // A retry must not audit files from an earlier attempt after
preparation fails.
+ Deencapsulation.setField(task, "auditSql", "stale audit SQL from a
previous attempt");
+ task.before();
+ Assertions.assertNull(task.getAuditSql());
+ Assertions.assertEquals(2, parsers.constructed().size());
+ Mockito.verify(parsers.constructed().get(0)).parseSingle(S3_SQL);
+
Mockito.verify(parsers.constructed().get(1)).parseForEncryption(Mockito.eq(S3_SQL),
Mockito.anyMap());
+ task.run();
+ Assertions.assertEquals(QueryState.MysqlStateType.OK,
state.getStateType());
+ Mockito.verify(taskCommand).run(ctx,
executors.constructed().get(1));
+ audit.verifyNoInteractions();
+ }
+ }
+
+ @Test
+ public void testRewriteS3UriCaseInsensitive() {
+ Map<String, String> originProperties = new HashMap<>();
+ originProperties.put("URI", ORIGIN_URI);
+ UnboundTVFRelation originTvf = new UnboundTVFRelation(
+ new RelationId(1), "s3", new Properties(originProperties));
+ InsertIntoTableCommand originCommand =
Mockito.mock(InsertIntoTableCommand.class);
+
Mockito.when(originCommand.getParsedPlan()).thenReturn(Optional.of(originTvf));
+
+ S3Offset offset = new S3Offset();
+ offset.setFileLists(RESOLVED_URI);
+ Env env = Mockito.mock(Env.class);
+ Mockito.when(env.isMaster()).thenReturn(false);
+ try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+ mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+ InsertIntoTableCommand rewritten = new S3SourceOffsetProvider()
+ .rewriteTvfParams(originCommand, offset, 1L);
+
+ Map<String, String> rewrittenProperties =
+
rewritten.getAllTVFRelation().get(0).getProperties().getMap();
+ Assertions.assertEquals(1, rewrittenProperties.size());
+ Assertions.assertEquals(RESOLVED_URI,
rewrittenProperties.get("URI"));
+ }
+ }
+
+ private AuditEvent runS3Task(RuntimeException commandFailure) throws
Exception {
+ S3Offset offset = new S3Offset();
+ offset.setFileLists(RESOLVED_URI);
+ return runTask(commandFailure, S3_SQL, new S3SourceOffsetProvider(),
offset, true);
+ }
+
+ private AuditEvent runTask(RuntimeException commandFailure, String sql,
+ SourceOffsetProvider offsetProvider, Offset offset, boolean
expectAudit) throws Exception {
+ Env env = Mockito.mock(Env.class);
+ CatalogMgr catalogMgr = Mockito.mock(CatalogMgr.class);
+ InternalCatalog catalog = Mockito.mock(InternalCatalog.class);
+ WorkloadRuntimeStatusMgr statusMgr =
Mockito.mock(WorkloadRuntimeStatusMgr.class);
+ Mockito.when(env.getCatalogMgr()).thenReturn(catalogMgr);
+ Mockito.when(env.getInternalCatalog()).thenReturn(catalog);
+
Mockito.when(catalogMgr.getCatalog(Mockito.anyString())).thenReturn(catalog);
+ Mockito.when(catalog.getName()).thenReturn("internal");
+ Mockito.when(env.getWorkloadRuntimeStatusMgr()).thenReturn(statusMgr);
+
+ try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+ mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+
+ ConnectContext ctx =
InsertTask.makeConnectContext(UserIdentity.ROOT, "test_db");
+ ctx.getState().setOk();
+ InsertIntoTableCommand command =
Mockito.mock(InsertIntoTableCommand.class);
+ if (expectAudit) {
+ Map<String, String> rewrittenTvfProps = new HashMap<>();
+ rewrittenTvfProps.put("uri", RESOLVED_URI);
+ rewrittenTvfProps.put("s3.secret_key", "private-value");
+ rewrittenTvfProps.put("enclose", "\"");
+ UnboundTVFRelation tvf =
Mockito.mock(UnboundTVFRelation.class);
+ Mockito.when(tvf.getProperties()).thenReturn(new
Properties(rewrittenTvfProps));
+
Mockito.when(command.getAllTVFRelation()).thenReturn(Collections.singletonList(tvf));
+ }
+ AtomicLong commandStartTime = new AtomicLong();
+ Mockito.doAnswer(invocation -> {
+ commandStartTime.set(ctx.getStartTime());
+ if (commandFailure != null) {
+ throw commandFailure;
+ }
+ return null;
+ }).when(command).run(Mockito.eq(ctx),
Mockito.any(StmtExecutor.class));
+
+ LogicalPlanAdapter parsedStmt = new LogicalPlanAdapter(
+ new NereidsParser().parseSingle(sql), new
StatementContext());
+ StmtExecutor executor = Mockito.mock(StmtExecutor.class);
+ Mockito.when(executor.getParsedStmt()).thenReturn(parsedStmt);
+
Mockito.when(executor.getSummaryProfile()).thenReturn(Mockito.mock(SummaryProfile.class));
+ ctx.setExecutor(executor);
+
+ StreamingInsertTask task = new StreamingInsertTask(
+ 1L, 2L, sql, offsetProvider, "test_db", null,
+ Collections.emptyMap(), UserIdentity.ROOT, null);
+ Deencapsulation.setField(task, "ctx", ctx);
+ Deencapsulation.setField(task, "taskCommand", command);
+ Deencapsulation.setField(task, "stmtExecutor", executor);
+ Deencapsulation.setField(task, "runningOffset", offset);
+ if (expectAudit) {
+ Deencapsulation.setField(task, "auditSql",
+ Deencapsulation.invoke(task, "buildAuditSql"));
+ }
+
+ if (commandFailure == null) {
+ task.run();
+ } else {
+ Assertions.assertThrows(JobException.class, task::run);
+ }
+
+ if (!expectAudit) {
+ Mockito.verify(statusMgr, Mockito.never())
+
.submitFinishQueryToAudit(Mockito.any(AuditEvent.class));
+ return null;
+ }
+ ArgumentCaptor<AuditEvent> auditEventCaptor =
ArgumentCaptor.forClass(AuditEvent.class);
+
Mockito.verify(statusMgr).submitFinishQueryToAudit(auditEventCaptor.capture());
+ AuditEvent auditEvent = auditEventCaptor.getValue();
+ Assertions.assertTrue(commandStartTime.get() > 0);
+ Assertions.assertEquals(commandStartTime.get(),
auditEvent.timestamp);
+ return auditEvent;
+ } finally {
+ ConnectContext.remove();
+ }
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]