This is an automated email from the ASF dual-hosted git repository.
dataroaring pushed a commit to branch branch-3.0
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-3.0 by this push:
new 9e19b86accd branch-3.0: [fix](audit) update audit table schema do not
work as expected #51363 (#52436)
9e19b86accd is described below
commit 9e19b86accd7e15a3ae4eb7e0a8957a8fbe8e8ba
Author: morrySnow <[email protected]>
AuthorDate: Wed Jul 2 10:40:13 2025 +0800
branch-3.0: [fix](audit) update audit table schema do not work as expected
#51363 (#52436)
pick park from #51363
---
.../doris/catalog/InternalSchemaInitializer.java | 79 +++++++++++++---------
1 file changed, 48 insertions(+), 31 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/catalog/InternalSchemaInitializer.java
b/fe/fe-core/src/main/java/org/apache/doris/catalog/InternalSchemaInitializer.java
index ca16d498f36..cb1c6082320 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/catalog/InternalSchemaInitializer.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/catalog/InternalSchemaInitializer.java
@@ -17,11 +17,11 @@
package org.apache.doris.catalog;
+import org.apache.doris.analysis.AddColumnsClause;
import org.apache.doris.analysis.AlterClause;
import org.apache.doris.analysis.AlterTableStmt;
import org.apache.doris.analysis.ColumnDef;
import org.apache.doris.analysis.ColumnNullableType;
-import org.apache.doris.analysis.ColumnPosition;
import org.apache.doris.analysis.CreateDbStmt;
import org.apache.doris.analysis.CreateTableStmt;
import org.apache.doris.analysis.DbName;
@@ -33,6 +33,7 @@ import org.apache.doris.analysis.ModifyColumnClause;
import org.apache.doris.analysis.ModifyPartitionClause;
import org.apache.doris.analysis.PartitionDesc;
import org.apache.doris.analysis.RangePartitionDesc;
+import org.apache.doris.analysis.ReorderColumnsClause;
import org.apache.doris.analysis.TableName;
import org.apache.doris.analysis.TypeDef;
import org.apache.doris.common.AnalysisException;
@@ -54,10 +55,12 @@ import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
+import java.util.stream.Collectors;
public class InternalSchemaInitializer extends Thread {
@@ -364,41 +367,55 @@ public class InternalSchemaInitializer extends Thread {
// 4. check and update audit table schema
OlapTable auditTable = (OlapTable) optionalStatsTbl.get();
- List<ColumnDef> expectedSchema = InternalSchema.AUDIT_SCHEMA;
// 5. check if we need to add new columns
+ return alterAuditSchemaIfNeeded(auditTable);
+ }
+
+ private boolean alterAuditSchemaIfNeeded(OlapTable auditTable) {
+ List<ColumnDef> expectedSchema = InternalSchema.AUDIT_SCHEMA;
+ List<String> expectedColumnNames = expectedSchema.stream()
+ .map(ColumnDef::getName)
+ .map(String::toLowerCase)
+ .collect(Collectors.toList());
+ List<Column> currentColumns = auditTable.getBaseSchema();
+ List<String> currentColumnNames = currentColumns.stream()
+ .map(Column::getName)
+ .map(String::toLowerCase)
+ .collect(Collectors.toList());
+ // check if all expected columns are exists and in the right order
+ if (currentColumnNames.size() >= expectedColumnNames.size()
+ && expectedColumnNames.equals(currentColumnNames.subList(0,
expectedColumnNames.size()))) {
+ return true;
+ }
+
List<AlterClause> alterClauses = Lists.newArrayList();
- for (int i = 0; i < expectedSchema.size(); i++) {
- ColumnDef def = expectedSchema.get(i);
- if (auditTable.getColumn(def.getName()) == null) {
- // add column if it doesn't exist
- try {
- ColumnDef columnDef = new ColumnDef(def.getName(),
def.getTypeDef(), def.isAllowNull());
- // find the previous column name to determine the position
- String afterColumn = null;
- if (i > 0) {
- for (int j = i - 1; j >= 0; j--) {
- String prevColName =
expectedSchema.get(j).getName();
- if (auditTable.getColumn(prevColName) != null) {
- afterColumn = prevColName;
- break;
- }
- }
- }
- ColumnPosition position = afterColumn == null ?
ColumnPosition.FIRST :
- new ColumnPosition(afterColumn);
- ModifyColumnClause clause = new
ModifyColumnClause(columnDef, position, null,
- Maps.newHashMap());
- clause.setColumn(columnDef.toColumn());
- alterClauses.add(clause);
- } catch (Exception e) {
- LOG.warn("Failed to create alter clause for column: " +
def.getName(), e);
- return false;
- }
+ // add new columns
+ List<ColumnDef> addColumnsDef = Lists.newArrayList();
+ for (ColumnDef expected : expectedSchema) {
+ if
(!currentColumnNames.contains(expected.getName().toLowerCase())) {
+ addColumnsDef.add(expected);
}
}
-
- // apply schema changes if needed
+ if (!addColumnsDef.isEmpty()) {
+ AddColumnsClause addColumnsClause = new
AddColumnsClause(addColumnsDef, null, Collections.emptyMap());
+ try {
+ addColumnsClause.analyze(null);
+ } catch (Exception e) {
+ LOG.warn("Failed to alter audit table schema", e);
+ return false;
+ }
+ alterClauses.add(addColumnsClause);
+ }
+ // reorder columns
+ List<String> removedColumnNames =
Lists.newArrayList(currentColumnNames);
+ removedColumnNames.removeAll(expectedColumnNames);
+ List<String> newColumnOrders = Lists.newArrayList(expectedColumnNames);
+ newColumnOrders.addAll(removedColumnNames);
+ if (!newColumnOrders.isEmpty()) {
+ ReorderColumnsClause reorderColumnsOp = new
ReorderColumnsClause(newColumnOrders, null, Maps.newHashMap());
+ alterClauses.add(reorderColumnsOp);
+ }
if (!alterClauses.isEmpty()) {
try {
TableName tableName = new
TableName(InternalCatalog.INTERNAL_CATALOG_NAME,
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]