diff --git a/fe/fe-core/src/main/java/org/apache/doris/analysis/Separator.java b/fe/fe-core/src/main/java/org/apache/doris/analysis/Separator.java index 67515eaca5c79f..7da2e092a212ad 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/analysis/Separator.java +++ b/fe/fe-core/src/main/java/org/apache/doris/analysis/Separator.java @@ -20,13 +20,16 @@ import org.apache.doris.common.AnalysisException; import com.google.common.base.Strings; +import com.google.gson.annotations.SerializedName; import java.io.StringWriter; public class Separator { private static final String HEX_STRING = "0123456789ABCDEF"; + @SerializedName("os") private final String oriSeparator; + @SerializedName("s") private String separator; public Separator(String separator) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/load/RoutineLoadDesc.java b/fe/fe-core/src/main/java/org/apache/doris/load/RoutineLoadDesc.java index 2c1ede0d13a352..28a429960d2b71 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/load/RoutineLoadDesc.java +++ b/fe/fe-core/src/main/java/org/apache/doris/load/RoutineLoadDesc.java @@ -26,19 +26,29 @@ import org.apache.doris.load.loadv2.LoadTask; import com.google.common.base.Strings; +import com.google.gson.annotations.SerializedName; import java.util.List; public class RoutineLoadDesc { + @SerializedName("cs") private final Separator columnSeparator; + @SerializedName("ld") private final Separator lineDelimiter; + @SerializedName("cols") private final List columnsInfo; + @SerializedName("pf") private final Expr precedingFilter; + @SerializedName("f") private final Expr filter; + @SerializedName("dc") private final Expr deleteCondition; + @SerializedName("mt") private LoadTask.MergeType mergeType; // nullable + @SerializedName("pn") private final PartitionNamesInfo partitionNamesInfo; + @SerializedName("sc") private final String sequenceColName; public RoutineLoadDesc(Separator columnSeparator, Separator lineDelimiter, List columnsInfo, diff --git a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadJob.java b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadJob.java index 9873368f405114..c01edd1ddcd28b 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadJob.java +++ b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadJob.java @@ -116,7 +116,6 @@ public abstract class RoutineLoadJob extends AbstractTxnStateChangeCallback implements Writable, LoadTaskInfo, GsonPostProcessable { private static final Logger LOG = LogManager.getLogger(RoutineLoadJob.class); - public static final long DEFAULT_MAX_ERROR_NUM = 0; public static final double DEFAULT_MAX_FILTER_RATIO = 1.0; @@ -182,11 +181,17 @@ public boolean isFinalState() { // this code is used to verify be task request protected long authCode; // protected RoutineLoadDesc routineLoadDesc; // optional + @SerializedName("pni") protected PartitionNamesInfo partitionNamesInfo; // optional + @SerializedName("cds") protected ImportColumnDescs columnDescs; // optional + @SerializedName("pf") protected Expr precedingFilter; // optional + @SerializedName("we") protected Expr whereExpr; // optional + @SerializedName("cs") protected Separator columnSeparator; // optional + @SerializedName("lidel") protected Separator lineDelimiter; @SerializedName("dtcn") protected int desireTaskConcurrentNum; // optional @@ -201,6 +206,7 @@ public boolean isFinalState() { @SerializedName("men") protected long maxErrorNum = DEFAULT_MAX_ERROR_NUM; // optional protected double maxFilterRatio = DEFAULT_MAX_FILTER_RATIO; + @SerializedName("eml") protected long execMemLimit = DEFAULT_EXEC_MEM_LIMIT; protected int sendBatchParallelism = DEFAULT_SEND_BATCH_PARALLELISM; protected boolean loadToSingleTablet = DEFAULT_LOAD_TO_SINGLE_TABLET; @@ -230,8 +236,10 @@ public boolean isFinalState() { protected TPartialUpdateNewRowPolicy partialUpdateNewKeyPolicy = TPartialUpdateNewRowPolicy.APPEND; protected TUniqueKeyUpdateMode uniqueKeyUpdateMode = TUniqueKeyUpdateMode.UPSERT; + @SerializedName("sc") protected String sequenceCol; + @SerializedName("mosn") protected boolean memtableOnSinkNode = false; protected int currentTaskConcurrentNum; @@ -258,8 +266,7 @@ public boolean isFinalState() { // The tasks belong to this job protected List routineLoadTaskInfoList = Lists.newArrayList(); - // this is the origin stmt of CreateRoutineLoadStmt, we use it to persist the RoutineLoadJob, - // because we can not serialize the Expressions contained in job. + // Keep the original CREATE statement for downgrade compatibility and legacy image migration. @SerializedName("ostmt") protected OriginStatement origStmt; // User who submit this job. Maybe null for the old version job(before v1.1) @@ -270,7 +277,9 @@ public boolean isFinalState() { protected String comment = ""; protected ReentrantReadWriteLock lock = new ReentrantReadWriteLock(true); - protected LoadTask.MergeType mergeType = LoadTask.MergeType.APPEND; // default is all data is load no delete + @SerializedName("mt") + protected LoadTask.MergeType mergeType; + @SerializedName("dc") protected Expr deleteCondition; // TODO(ml): error sample @@ -316,6 +325,7 @@ public RoutineLoadJob(Long id, String name, this.tableId = tableId; this.authCode = 0; this.userIdentity = userIdentity; + this.mergeType = LoadTask.MergeType.APPEND; if (ConnectContext.get() != null) { SessionVariable var = ConnectContext.get().getSessionVariable(); @@ -345,6 +355,7 @@ public RoutineLoadJob(Long id, String name, this.authCode = 0; this.userIdentity = userIdentity; this.isMultiTable = true; + this.mergeType = LoadTask.MergeType.APPEND; if (ConnectContext.get() != null) { SessionVariable var = ConnectContext.get().getSessionVariable(); @@ -1967,87 +1978,137 @@ public void gsonPostProcess() throws IOException { if (tableId == 0) { isMultiTable = true; } - // Process UNIQUE_KEY_UPDATE_MODE first to ensure correct backward compatibility - // with PARTIAL_COLUMNS (HashMap iteration order is not guaranteed) - if (jobProperties.containsKey(CreateRoutineLoadInfo.UNIQUE_KEY_UPDATE_MODE)) { - String modeValue = jobProperties.get(CreateRoutineLoadInfo.UNIQUE_KEY_UPDATE_MODE); - TUniqueKeyUpdateMode mode = CreateRoutineLoadInfo.parseUniqueKeyUpdateMode(modeValue); - if (mode != null) { - uniqueKeyUpdateMode = mode; - isPartialUpdate = (uniqueKeyUpdateMode == TUniqueKeyUpdateMode.UPDATE_FIXED_COLUMNS); - } else { - uniqueKeyUpdateMode = TUniqueKeyUpdateMode.UPSERT; - } + // Legacy images did not persist mergeType. New images always contain it, including jobs + // without any load clause, so its absence is sufficient to identify the one-time fallback. + boolean isOldImage = mergeType == null; + if (isOldImage) { + mergeType = LoadTask.MergeType.APPEND; + // Legacy images did not persist this create-time session option. Preserve their historical + // post-restart behavior instead of inheriting the image-loading thread's ConnectContext. + memtableOnSinkNode = false; } - // Process remaining properties - jobProperties.forEach((k, v) -> { - if (k.equals(CreateRoutineLoadInfo.PARTIAL_COLUMNS)) { - // Backward compatibility: only use partial_columns if unique_key_update_mode is not set - // unique_key_update_mode takes precedence - if (uniqueKeyUpdateMode == TUniqueKeyUpdateMode.UPSERT) { - isPartialUpdate = Boolean.parseBoolean(v); - if (isPartialUpdate) { - uniqueKeyUpdateMode = TUniqueKeyUpdateMode.UPDATE_FIXED_COLUMNS; - } - } - } else if (k.equals(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY)) { - if ("ERROR".equalsIgnoreCase(v)) { - partialUpdateNewKeyPolicy = TPartialUpdateNewRowPolicy.ERROR; - } else { - partialUpdateNewKeyPolicy = TPartialUpdateNewRowPolicy.APPEND; - } - } - }); try { - ConnectContext ctx = new ConnectContext(); - ctx.setDatabase(Env.getCurrentEnv().getInternalCatalog().getDb(dbId).get().getName()); - StatementContext statementContext = new StatementContext(); - statementContext.setConnectContext(ctx); - ctx.setStatementContext(statementContext); - ctx.setEnv(Env.getCurrentEnv()); - ctx.setCurrentUserIdentity(UserIdentity.ADMIN); - ctx.getState().reset(); - try { - ctx.setThreadLocalInfo(); - NereidsParser nereidsParser = new NereidsParser(); - CreateRoutineLoadCommand command = (CreateRoutineLoadCommand) nereidsParser.parseSingle( - origStmt.originStmt); - CreateRoutineLoadInfo createRoutineLoadInfo = command.getCreateRoutineLoadInfo(); - // If tableId is set, resolve the current table name by ID so that - // table rename / SWAP TABLE won't cause replay to fail with stale name in origStmt. - if (!isMultiTable && tableId != 0) { - try { - Database db = Env.getCurrentEnv().getInternalCatalog().getDb(dbId).orElse(null); - if (db != null) { - db.getTable(tableId).ifPresent( - table -> createRoutineLoadInfo.setTableName(table.getName())); - } - } catch (Exception ignored) { - // fall through; let validate() surface the real error - } - } - createRoutineLoadInfo.validate(ctx); - setRoutineLoadDesc(createRoutineLoadInfo.getRoutineLoadDesc()); - } finally { - ctx.cleanup(); + hydrateJobProperties(); + if (isOldImage) { + restoreLegacyDefinition(); } } catch (Exception e) { this.state = JobState.CANCELLED; - LOG.warn("error happens when parsing create routine load stmt: " + origStmt.originStmt, e); + LOG.warn("error happens when restoring routine load job", e); } if (userIdentity != null) { userIdentity.setIsAnalyzed(); } } + private void hydrateJobProperties() throws UserException { + if (jobProperties.containsKey(CreateRoutineLoadInfo.MAX_FILTER_RATIO_PROPERTY)) { + maxFilterRatio = Double.parseDouble( + jobProperties.get(CreateRoutineLoadInfo.MAX_FILTER_RATIO_PROPERTY)); + } + if (jobProperties.containsKey(CreateRoutineLoadInfo.SEND_BATCH_PARALLELISM)) { + sendBatchParallelism = Integer.parseInt( + jobProperties.get(CreateRoutineLoadInfo.SEND_BATCH_PARALLELISM)); + } + if (jobProperties.containsKey(CreateRoutineLoadInfo.LOAD_TO_SINGLE_TABLET)) { + loadToSingleTablet = Boolean.parseBoolean( + jobProperties.get(CreateRoutineLoadInfo.LOAD_TO_SINGLE_TABLET)); + } + + boolean hasUniqueKeyUpdateMode = jobProperties.containsKey(CreateRoutineLoadInfo.UNIQUE_KEY_UPDATE_MODE); + if (hasUniqueKeyUpdateMode) { + TUniqueKeyUpdateMode mode = CreateRoutineLoadInfo.parseUniqueKeyUpdateMode( + jobProperties.get(CreateRoutineLoadInfo.UNIQUE_KEY_UPDATE_MODE)); + uniqueKeyUpdateMode = mode == null ? TUniqueKeyUpdateMode.UPSERT : mode; + isPartialUpdate = uniqueKeyUpdateMode == TUniqueKeyUpdateMode.UPDATE_FIXED_COLUMNS; + } + if (!hasUniqueKeyUpdateMode && jobProperties.containsKey(CreateRoutineLoadInfo.PARTIAL_COLUMNS)) { + isPartialUpdate = Boolean.parseBoolean(jobProperties.get(CreateRoutineLoadInfo.PARTIAL_COLUMNS)); + if (isPartialUpdate) { + uniqueKeyUpdateMode = TUniqueKeyUpdateMode.UPDATE_FIXED_COLUMNS; + } + } + if (jobProperties.containsKey(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY)) { + partialUpdateNewKeyPolicy = "ERROR".equalsIgnoreCase( + jobProperties.get(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY)) + ? TPartialUpdateNewRowPolicy.ERROR : TPartialUpdateNewRowPolicy.APPEND; + } + + if (jobProperties.containsKey(CsvFileFormatProperties.PROP_ENCLOSE)) { + enclose = parseEnclose(jobProperties.get(CsvFileFormatProperties.PROP_ENCLOSE)); + } + if (jobProperties.containsKey(CsvFileFormatProperties.PROP_ESCAPE)) { + escape = parseEscape(jobProperties.get(CsvFileFormatProperties.PROP_ESCAPE)); + } + if (jobProperties.containsKey(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL)) { + emptyFieldAsNull = Boolean.parseBoolean( + jobProperties.get(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL)); + } + } + + private void restoreLegacyDefinition() throws UserException { + ConnectContext ctx = new ConnectContext(); + ctx.setDatabase(Env.getCurrentEnv().getInternalCatalog().getDb(dbId).get().getName()); + StatementContext statementContext = new StatementContext(); + statementContext.setConnectContext(ctx); + ctx.setStatementContext(statementContext); + ctx.setEnv(Env.getCurrentEnv()); + ctx.setCurrentUserIdentity(UserIdentity.ADMIN); + ctx.getState().reset(); + try { + ctx.setThreadLocalInfo(); + NereidsParser nereidsParser = new NereidsParser(); + CreateRoutineLoadCommand command = (CreateRoutineLoadCommand) nereidsParser.parseSingle( + origStmt.originStmt); + CreateRoutineLoadInfo createRoutineLoadInfo = command.getCreateRoutineLoadInfo(); + // Resolve the current table name by ID so table rename or SWAP TABLE does not leave the + // legacy CREATE statement pointing at a stale table name. + if (!isMultiTable && tableId != 0) { + try { + Database db = Env.getCurrentEnv().getInternalCatalog().getDb(dbId).orElse(null); + if (db != null) { + db.getTable(tableId).ifPresent( + table -> createRoutineLoadInfo.setTableName(table.getName())); + } + } catch (Exception ignored) { + // Let validate() below surface the original catalog error. + } + } + createRoutineLoadInfo.validate(ctx); + setRoutineLoadDesc(createRoutineLoadInfo.getRoutineLoadDesc()); + execMemLimit = createRoutineLoadInfo.getExecMemLimit(); + } finally { + ctx.cleanup(); + } + } + public abstract void modifyProperties(AlterRoutineLoadCommand command) throws UserException; public abstract void replayModifyProperties(AlterRoutineLoadJobOperationLog log); public abstract NereidsRoutineLoadTaskInfo toNereidsRoutineLoadTaskInfo() throws UserException; - // for ALTER ROUTINE LOAD + protected void validateCommonJobProperties(Map jobProperties) throws UserException { + validateCsvFormatProperties(jobProperties); + if (jobProperties.containsKey(CreateRoutineLoadInfo.UNIQUE_KEY_UPDATE_MODE)) { + TUniqueKeyUpdateMode newMode = CreateRoutineLoadInfo.parseAndValidateUniqueKeyUpdateMode( + jobProperties.get(CreateRoutineLoadInfo.UNIQUE_KEY_UPDATE_MODE)); + if (newMode == TUniqueKeyUpdateMode.UPDATE_FLEXIBLE_COLUMNS) { + validateFlexiblePartialUpdateForAlter(); + } + } + if (jobProperties.containsKey(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY)) { + String policy = jobProperties.get(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY); + if (!"APPEND".equalsIgnoreCase(policy) && !"ERROR".equalsIgnoreCase(policy)) { + throw new AnalysisException(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY + + " should be one of {'APPEND', 'ERROR'}, but found " + policy); + } + } + } + + // for ALTER ROUTINE LOAD. Validate all common properties before changing any common runtime state. protected void modifyCommonJobProperties(Map jobProperties) throws UserException { + validateCommonJobProperties(jobProperties); if (jobProperties.containsKey(CreateRoutineLoadInfo.DESIRED_CONCURRENT_NUMBER_PROPERTY)) { this.desireTaskConcurrentNum = Integer.parseInt( jobProperties.remove(CreateRoutineLoadInfo.DESIRED_CONCURRENT_NUMBER_PROPERTY)); @@ -2081,12 +2142,7 @@ protected void modifyCommonJobProperties(Map jobProperties) thro if (jobProperties.containsKey(CreateRoutineLoadInfo.UNIQUE_KEY_UPDATE_MODE)) { String modeStr = jobProperties.remove(CreateRoutineLoadInfo.UNIQUE_KEY_UPDATE_MODE); - TUniqueKeyUpdateMode newMode = CreateRoutineLoadInfo.parseAndValidateUniqueKeyUpdateMode(modeStr); - // Validate flexible partial update constraints when changing to UPDATE_FLEXIBLE_COLUMNS - if (newMode == TUniqueKeyUpdateMode.UPDATE_FLEXIBLE_COLUMNS) { - validateFlexiblePartialUpdateForAlter(); - } - this.uniqueKeyUpdateMode = newMode; + this.uniqueKeyUpdateMode = CreateRoutineLoadInfo.parseAndValidateUniqueKeyUpdateMode(modeStr); this.isPartialUpdate = (uniqueKeyUpdateMode == TUniqueKeyUpdateMode.UPDATE_FIXED_COLUMNS); this.jobProperties.put(CreateRoutineLoadInfo.UNIQUE_KEY_UPDATE_MODE, uniqueKeyUpdateMode.name()); this.jobProperties.put(CreateRoutineLoadInfo.PARTIAL_COLUMNS, String.valueOf(isPartialUpdate)); @@ -2103,6 +2159,55 @@ protected void modifyCommonJobProperties(Map jobProperties) thro this.jobProperties.put(CreateRoutineLoadInfo.PARTIAL_COLUMNS, String.valueOf(isPartialUpdate)); this.jobProperties.put(CreateRoutineLoadInfo.UNIQUE_KEY_UPDATE_MODE, uniqueKeyUpdateMode.name()); } + + if (jobProperties.containsKey(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY)) { + String policy = jobProperties.remove(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY); + partialUpdateNewKeyPolicy = "ERROR".equalsIgnoreCase(policy) + ? TPartialUpdateNewRowPolicy.ERROR : TPartialUpdateNewRowPolicy.APPEND; + this.jobProperties.put(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY, + partialUpdateNewKeyPolicy.name()); + } + + if (jobProperties.containsKey(CsvFileFormatProperties.PROP_ENCLOSE)) { + String value = jobProperties.remove(CsvFileFormatProperties.PROP_ENCLOSE); + enclose = parseEnclose(value); + this.jobProperties.put(CsvFileFormatProperties.PROP_ENCLOSE, value); + } + if (jobProperties.containsKey(CsvFileFormatProperties.PROP_ESCAPE)) { + String value = jobProperties.remove(CsvFileFormatProperties.PROP_ESCAPE); + escape = parseEscape(value); + this.jobProperties.put(CsvFileFormatProperties.PROP_ESCAPE, value); + } + if (jobProperties.containsKey(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL)) { + String value = jobProperties.remove(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL); + emptyFieldAsNull = Boolean.parseBoolean(value); + this.jobProperties.put(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL, value); + } + } + + private static void validateCsvFormatProperties(Map jobProperties) { + if (!jobProperties.containsKey(CsvFileFormatProperties.PROP_ENCLOSE) + && !jobProperties.containsKey(CsvFileFormatProperties.PROP_ESCAPE) + && !jobProperties.containsKey(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL)) { + return; + } + Map csvProperties = Maps.newHashMap(); + for (String property : new String[] {CsvFileFormatProperties.PROP_ENCLOSE, + CsvFileFormatProperties.PROP_ESCAPE, CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL}) { + if (jobProperties.containsKey(property)) { + csvProperties.put(property, jobProperties.get(property)); + } + } + CsvFileFormatProperties properties = new CsvFileFormatProperties(FileFormatProperties.FORMAT_CSV); + properties.analyzeFileFormatProperties(csvProperties, false); + } + + private static byte parseEnclose(String value) { + return Strings.isNullOrEmpty(value) ? 0 : (byte) value.charAt(0); + } + + private static byte parseEscape(String value) { + return Strings.isNullOrEmpty(value) ? 0 : value.getBytes()[0]; } /** diff --git a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadManager.java b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadManager.java index df5615016bfa54..9bc9b93abd8c98 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadManager.java +++ b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadManager.java @@ -943,7 +943,6 @@ public void alterRoutineLoadJob(AlterRoutineLoadCommand command) throws UserExce + command.getDataSourceProperties().getDataSourceType()); } job.modifyProperties(command); - job.setRoutineLoadDesc(command.getRoutineLoadDesc()); } public void replayAlterRoutineLoadJob(AlterRoutineLoadJobOperationLog log) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java index 885021440351d7..be124f7c72e6c5 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java +++ b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java @@ -60,7 +60,6 @@ import org.apache.doris.rpc.RpcException; import org.apache.doris.service.FrontendOptions; import org.apache.doris.thrift.TFileCompressType; -import org.apache.doris.thrift.TPartialUpdateNewRowPolicy; import org.apache.doris.transaction.TransactionState; import org.apache.doris.transaction.TransactionStatus; @@ -74,7 +73,6 @@ import com.google.gson.annotations.SerializedName; import org.apache.commons.collections4.CollectionUtils; import org.apache.commons.collections4.MapUtils; -import org.apache.commons.lang3.BooleanUtils; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -778,6 +776,7 @@ public Map getCustomProperties() { @Override public void modifyProperties(AlterRoutineLoadCommand command) throws UserException { Map jobProperties = command.getAnalyzedJobProperties(); + validateCommonJobProperties(jobProperties); KafkaDataSourceProperties dataSourceProperties = (KafkaDataSourceProperties) command.getDataSourceProperties(); if (null != dataSourceProperties) { // if the partition offset is set by timestamp, convert it to real offset @@ -791,9 +790,10 @@ public void modifyProperties(AlterRoutineLoadCommand command) throws UserExcepti } modifyPropertiesInternal(jobProperties, dataSourceProperties); + setRoutineLoadDesc(command.getRoutineLoadDesc()); AlterRoutineLoadJobOperationLog log = new AlterRoutineLoadJobOperationLog(this.id, - jobProperties, dataSourceProperties); + jobProperties, dataSourceProperties, command.getRoutineLoadDesc()); Env.getCurrentEnv().getEditLog().logAlterRoutineLoadJob(log); } finally { writeUnlock(); @@ -881,17 +881,6 @@ private void modifyPropertiesInternal(Map jobProperties, Map copiedJobProperties = Maps.newHashMap(jobProperties); modifyCommonJobProperties(copiedJobProperties); this.jobProperties.putAll(copiedJobProperties); - if (jobProperties.containsKey(CreateRoutineLoadInfo.PARTIAL_COLUMNS)) { - this.isPartialUpdate = BooleanUtils.toBoolean(jobProperties.get(CreateRoutineLoadInfo.PARTIAL_COLUMNS)); - } - if (jobProperties.containsKey(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY)) { - String policy = jobProperties.get(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY); - if ("ERROR".equalsIgnoreCase(policy)) { - this.partialUpdateNewKeyPolicy = TPartialUpdateNewRowPolicy.ERROR; - } else { - this.partialUpdateNewKeyPolicy = TPartialUpdateNewRowPolicy.APPEND; - } - } } LOG.info("modify the properties of kafka routine load job: {}, jobProperties: {}, datasource properties: {}", this.id, jobProperties, dataSourceProperties); @@ -920,6 +909,7 @@ private void resetCloudProgress(Cloud.ResetRLProgressRequest.Builder builder) th public void replayModifyProperties(AlterRoutineLoadJobOperationLog log) { try { modifyPropertiesInternal(log.getJobProperties(), (KafkaDataSourceProperties) log.getDataSourceProperties()); + setRoutineLoadDesc(log.getRoutineLoadDesc()); } catch (UserException e) { // should not happen LOG.error("failed to replay modify kafka routine load job: {}", id, e); diff --git a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kinesis/KinesisRoutineLoadJob.java b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kinesis/KinesisRoutineLoadJob.java index 7cebc3f5165b49..9c5bac1c7895a2 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kinesis/KinesisRoutineLoadJob.java +++ b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kinesis/KinesisRoutineLoadJob.java @@ -51,7 +51,6 @@ import org.apache.doris.persist.AlterRoutineLoadJobOperationLog; import org.apache.doris.qe.ConnectContext; import org.apache.doris.thrift.TFileCompressType; -import org.apache.doris.thrift.TPartialUpdateNewRowPolicy; import org.apache.doris.transaction.TransactionState; import org.apache.doris.transaction.TransactionStatus; @@ -65,7 +64,6 @@ import com.google.gson.annotations.SerializedName; import org.apache.commons.collections4.CollectionUtils; import org.apache.commons.collections4.MapUtils; -import org.apache.commons.lang3.BooleanUtils; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -687,9 +685,10 @@ public void modifyProperties(AlterRoutineLoadCommand command) throws UserExcepti } modifyPropertiesInternal(jobProperties, dataSourceProperties); + setRoutineLoadDesc(command.getRoutineLoadDesc()); AlterRoutineLoadJobOperationLog log = new AlterRoutineLoadJobOperationLog(this.id, - jobProperties, dataSourceProperties); + jobProperties, dataSourceProperties, command.getRoutineLoadDesc()); Env.getCurrentEnv().getEditLog().logAlterRoutineLoadJob(log); } finally { writeUnlock(); @@ -699,6 +698,7 @@ public void modifyProperties(AlterRoutineLoadCommand command) throws UserExcepti private void modifyPropertiesInternal(Map jobProperties, KinesisDataSourceProperties dataSourceProperties) throws UserException { + validateCommonJobProperties(jobProperties); if (dataSourceProperties != null) { List> shardPositions = Lists.newArrayList(); Map customKinesisProperties = Maps.newHashMap(); @@ -762,17 +762,6 @@ private void modifyPropertiesInternal(Map jobProperties, Map copiedJobProperties = Maps.newHashMap(jobProperties); modifyCommonJobProperties(copiedJobProperties); this.jobProperties.putAll(copiedJobProperties); - if (jobProperties.containsKey(CreateRoutineLoadInfo.PARTIAL_COLUMNS)) { - this.isPartialUpdate = BooleanUtils.toBoolean(jobProperties.get(CreateRoutineLoadInfo.PARTIAL_COLUMNS)); - } - if (jobProperties.containsKey(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY)) { - String policy = jobProperties.get(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY); - if ("ERROR".equalsIgnoreCase(policy)) { - this.partialUpdateNewKeyPolicy = TPartialUpdateNewRowPolicy.ERROR; - } else { - this.partialUpdateNewKeyPolicy = TPartialUpdateNewRowPolicy.APPEND; - } - } } LOG.info("modify the properties of kinesis routine load job: {}, jobProperties: {}, datasource properties: {}", this.id, jobProperties, dataSourceProperties); @@ -783,6 +772,7 @@ public void replayModifyProperties(AlterRoutineLoadJobOperationLog log) { try { modifyPropertiesInternal(log.getJobProperties(), (KinesisDataSourceProperties) log.getDataSourceProperties()); + setRoutineLoadDesc(log.getRoutineLoadDesc()); } catch (UserException e) { LOG.error("failed to replay modify kinesis routine load job: {}", id, e); } diff --git a/fe/fe-core/src/main/java/org/apache/doris/persist/AlterRoutineLoadJobOperationLog.java b/fe/fe-core/src/main/java/org/apache/doris/persist/AlterRoutineLoadJobOperationLog.java index 4729882f7927fb..9d8064943c566a 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/persist/AlterRoutineLoadJobOperationLog.java +++ b/fe/fe-core/src/main/java/org/apache/doris/persist/AlterRoutineLoadJobOperationLog.java @@ -19,6 +19,7 @@ import org.apache.doris.common.io.Text; import org.apache.doris.common.io.Writable; +import org.apache.doris.load.RoutineLoadDesc; import org.apache.doris.load.routineload.AbstractDataSourceProperties; import org.apache.doris.persist.gson.GsonUtils; @@ -37,12 +38,20 @@ public class AlterRoutineLoadJobOperationLog implements Writable { private Map jobProperties; @SerializedName(value = "dataSourceProperties") private AbstractDataSourceProperties dataSourceProperties; + @SerializedName(value = "routineLoadDesc") + private RoutineLoadDesc routineLoadDesc; public AlterRoutineLoadJobOperationLog(long jobId, Map jobProperties, AbstractDataSourceProperties dataSourceProperties) { + this(jobId, jobProperties, dataSourceProperties, null); + } + + public AlterRoutineLoadJobOperationLog(long jobId, Map jobProperties, + AbstractDataSourceProperties dataSourceProperties, RoutineLoadDesc routineLoadDesc) { this.jobId = jobId; this.jobProperties = jobProperties; this.dataSourceProperties = dataSourceProperties; + this.routineLoadDesc = routineLoadDesc; } public long getJobId() { @@ -57,6 +66,10 @@ public AbstractDataSourceProperties getDataSourceProperties() { return dataSourceProperties; } + public RoutineLoadDesc getRoutineLoadDesc() { + return routineLoadDesc; + } + public static AlterRoutineLoadJobOperationLog read(DataInput in) throws IOException { String json = Text.readString(in); return GsonUtils.GSON.fromJson(json, AlterRoutineLoadJobOperationLog.class); diff --git a/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java b/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java index 7f0c8588372403..ca0748775e0f28 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java @@ -32,6 +32,7 @@ import org.apache.doris.common.jmockit.Deencapsulation; import org.apache.doris.datasource.InternalCatalog; import org.apache.doris.datasource.kafka.KafkaUtil; +import org.apache.doris.datasource.property.fileformat.CsvFileFormatProperties; import org.apache.doris.load.RoutineLoadDesc; import org.apache.doris.load.loadv2.LoadTask; import org.apache.doris.load.routineload.kafka.KafkaConfiguration; @@ -40,10 +41,13 @@ import org.apache.doris.load.routineload.kafka.KafkaRoutineLoadJob; import org.apache.doris.load.routineload.kafka.KafkaTaskInfo; import org.apache.doris.mysql.privilege.MockedAuth; +import org.apache.doris.nereids.trees.plans.commands.AlterRoutineLoadCommand; import org.apache.doris.nereids.trees.plans.commands.info.CreateRoutineLoadInfo; import org.apache.doris.nereids.trees.plans.commands.info.LabelNameInfo; import org.apache.doris.nereids.trees.plans.commands.load.LoadProperty; import org.apache.doris.nereids.trees.plans.commands.load.LoadSeparator; +import org.apache.doris.persist.AlterRoutineLoadJobOperationLog; +import org.apache.doris.persist.EditLog; import org.apache.doris.qe.ConnectContext; import org.apache.doris.thrift.TResourceInfo; import org.apache.doris.thrift.TRoutineLoadTask; @@ -58,9 +62,14 @@ import org.junit.Assert; import org.junit.Before; import org.junit.Test; +import org.mockito.ArgumentCaptor; import org.mockito.MockedStatic; import org.mockito.Mockito; +import java.io.ByteArrayInputStream; +import java.io.ByteArrayOutputStream; +import java.io.DataInputStream; +import java.io.DataOutputStream; import java.util.ArrayList; import java.util.Arrays; import java.util.HashMap; @@ -272,6 +281,84 @@ public void testUpdateProgressWarnsWhenReadCommittedTaskHasZeroRowsAndLag() thro Assert.assertTrue(otherMsg.contains("some records may be in uncommitted transactions")); } + @Test + public void testAlterPersistsLoadDescAndCsvPropertiesForReplay() throws Exception { + KafkaRoutineLoadJob leader = createPausedJob(); + KafkaRoutineLoadJob follower = createPausedJob(); + RoutineLoadDesc originalDesc = new RoutineLoadDesc(new Separator("|", "|"), null, null, + null, null, null, null, LoadTask.MergeType.APPEND, "original_sequence"); + leader.setRoutineLoadDesc(originalDesc); + follower.setRoutineLoadDesc(originalDesc); + + Map jobProperties = Maps.newHashMap(); + jobProperties.put(CsvFileFormatProperties.PROP_ENCLOSE, "\""); + jobProperties.put(CsvFileFormatProperties.PROP_ESCAPE, "\\"); + jobProperties.put(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL, "true"); + RoutineLoadDesc delta = new RoutineLoadDesc(null, new Separator("\n", "\\n"), null, + null, null, null, null, LoadTask.MergeType.APPEND, null); + AlterRoutineLoadCommand command = Mockito.mock(AlterRoutineLoadCommand.class); + Mockito.when(command.getAnalyzedJobProperties()).thenReturn(jobProperties); + Mockito.when(command.getDataSourceProperties()).thenReturn(null); + Mockito.when(command.getRoutineLoadDesc()).thenReturn(delta); + + Env env = Mockito.mock(Env.class); + EditLog editLog = Mockito.mock(EditLog.class); + AlterRoutineLoadJobOperationLog alterLog; + try (MockedStatic envStatic = Mockito.mockStatic(Env.class)) { + envStatic.when(Env::getCurrentEnv).thenReturn(env); + Mockito.when(env.getEditLog()).thenReturn(editLog); + + leader.modifyProperties(command); + + ArgumentCaptor logCaptor = + ArgumentCaptor.forClass(AlterRoutineLoadJobOperationLog.class); + Mockito.verify(editLog).logAlterRoutineLoadJob(logCaptor.capture()); + alterLog = logCaptor.getValue(); + } + + Assert.assertSame(delta, alterLog.getRoutineLoadDesc()); + Assert.assertEquals(jobProperties, alterLog.getJobProperties()); + assertAlterState(leader); + + follower.replayModifyProperties(alterLog); + assertAlterState(follower); + + assertAlterState(imageRoundTrip(leader)); + assertAlterState(imageRoundTrip(follower)); + } + + private static KafkaRoutineLoadJob createPausedJob() { + KafkaRoutineLoadJob job = new KafkaRoutineLoadJob(1L, "job1", 1L, + 1L, "127.0.0.1:9020", "topic1", UserIdentity.ADMIN); + Deencapsulation.setField(job, "state", RoutineLoadJob.JobState.PAUSED); + return job; + } + + private static void assertAlterState(RoutineLoadJob job) { + Assert.assertEquals("|", job.getColumnSeparator().getSeparator()); + Assert.assertEquals("\n", job.getLineDelimiter().getSeparator()); + Assert.assertEquals("original_sequence", job.getSequenceCol()); + Assert.assertEquals((byte) '"', job.getEnclose()); + Assert.assertEquals((byte) '\\', job.getEscape()); + Assert.assertTrue(job.getEmptyFieldAsNull()); + Assert.assertEquals(Boolean.TRUE, Deencapsulation.getField(job, "emptyFieldAsNull")); + + Map persistedJobProperties = Deencapsulation.getField(job, "jobProperties"); + Assert.assertEquals("\"", persistedJobProperties.get(CsvFileFormatProperties.PROP_ENCLOSE)); + Assert.assertEquals("\\", persistedJobProperties.get(CsvFileFormatProperties.PROP_ESCAPE)); + Assert.assertEquals("true", persistedJobProperties.get(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL)); + } + + private static RoutineLoadJob imageRoundTrip(RoutineLoadJob routineLoadJob) throws Exception { + ByteArrayOutputStream bytes = new ByteArrayOutputStream(); + try (DataOutputStream out = new DataOutputStream(bytes)) { + routineLoadJob.write(out); + } + try (DataInputStream in = new DataInputStream(new ByteArrayInputStream(bytes.toByteArray()))) { + return RoutineLoadJob.read(in); + } + } + @Test public void testDisplayCustomPropertiesMasksKafkaSecrets() { KafkaRoutineLoadJob routineLoadJob = new KafkaRoutineLoadJob(1L, "kafka_routine_load_job", 1L, diff --git a/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KinesisRoutineLoadJobTest.java b/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KinesisRoutineLoadJobTest.java index 65aebd0084e729..ea64e7e2871b0c 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KinesisRoutineLoadJobTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KinesisRoutineLoadJobTest.java @@ -17,21 +17,39 @@ package org.apache.doris.load.routineload; +import org.apache.doris.analysis.Separator; import org.apache.doris.analysis.UserIdentity; +import org.apache.doris.catalog.Env; import org.apache.doris.common.Config; +import org.apache.doris.common.io.Text; import org.apache.doris.common.jmockit.Deencapsulation; +import org.apache.doris.datasource.property.fileformat.CsvFileFormatProperties; +import org.apache.doris.load.RoutineLoadDesc; +import org.apache.doris.load.loadv2.LoadTask; import org.apache.doris.load.routineload.kinesis.KinesisConfiguration; import org.apache.doris.load.routineload.kinesis.KinesisDataSourceProperties; import org.apache.doris.load.routineload.kinesis.KinesisProgress; import org.apache.doris.load.routineload.kinesis.KinesisRoutineLoadJob; import org.apache.doris.load.routineload.kinesis.KinesisTaskInfo; +import org.apache.doris.nereids.exceptions.AnalysisException; +import org.apache.doris.nereids.trees.plans.commands.AlterRoutineLoadCommand; +import org.apache.doris.persist.AlterRoutineLoadJobOperationLog; +import org.apache.doris.persist.EditLog; import com.google.common.collect.Lists; import com.google.common.collect.Maps; import com.google.gson.Gson; +import com.google.gson.JsonParser; import org.junit.Assert; import org.junit.Test; - +import org.mockito.ArgumentCaptor; +import org.mockito.MockedStatic; +import org.mockito.Mockito; + +import java.io.ByteArrayInputStream; +import java.io.ByteArrayOutputStream; +import java.io.DataInputStream; +import java.io.DataOutputStream; import java.util.HashMap; import java.util.HashSet; import java.util.List; @@ -229,6 +247,58 @@ public void testModifyPropertiesShouldReplaceCustomShardsWhenExplicitShardsProvi Assert.assertEquals("202", progress.getSequenceNumberByShard("shard-2")); } + @Test + public void testAlterReplayKeepsDeltaAndCsvCachesInCheckpointParity() throws Exception { + KinesisRoutineLoadJob leader = createPausedJobWithInitialLoadDesc(); + KinesisRoutineLoadJob replay = createPausedJobWithInitialLoadDesc(); + + Map jobProperties = Maps.newHashMap(); + jobProperties.put(CsvFileFormatProperties.PROP_ENCLOSE, "\""); + jobProperties.put(CsvFileFormatProperties.PROP_ESCAPE, "\\"); + jobProperties.put(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL, "true"); + RoutineLoadDesc delta = new RoutineLoadDesc(null, new Separator("\n", "\\n"), + null, null, null, null, null, LoadTask.MergeType.APPEND, "sequence_col"); + AlterRoutineLoadCommand command = Mockito.mock(AlterRoutineLoadCommand.class); + Mockito.when(command.getAnalyzedJobProperties()).thenReturn(jobProperties); + Mockito.when(command.getDataSourceProperties()).thenReturn(null); + Mockito.when(command.getRoutineLoadDesc()).thenReturn(delta); + + Env env = Mockito.mock(Env.class); + EditLog editLog = Mockito.mock(EditLog.class); + ArgumentCaptor logCaptor = + ArgumentCaptor.forClass(AlterRoutineLoadJobOperationLog.class); + try (MockedStatic envStatic = Mockito.mockStatic(Env.class)) { + envStatic.when(Env::getCurrentEnv).thenReturn(env); + Mockito.when(env.getEditLog()).thenReturn(editLog); + leader.modifyProperties(command); + Mockito.verify(editLog).logAlterRoutineLoadJob(logCaptor.capture()); + } + + AlterRoutineLoadJobOperationLog log = logCaptor.getValue(); + replay.replayModifyProperties(log); + + Assert.assertSame(delta, log.getRoutineLoadDesc()); + assertAlterResult(leader); + assertAlterResult(replay); + Assert.assertEquals(JsonParser.parseString(checkpointJson(leader)), + JsonParser.parseString(checkpointJson(replay))); + } + + @Test + public void testAlterValidatesCsvBeforeDataSourceMutation() { + KinesisRoutineLoadJob job = createPausedJobWithInitialLoadDesc(); + Map jobProperties = Maps.newHashMap(); + jobProperties.put(CsvFileFormatProperties.PROP_ENCLOSE, "invalid"); + KinesisDataSourceProperties dataSourceProperties = Mockito.mock(KinesisDataSourceProperties.class); + AlterRoutineLoadCommand command = Mockito.mock(AlterRoutineLoadCommand.class); + Mockito.when(command.getAnalyzedJobProperties()).thenReturn(jobProperties); + Mockito.when(command.getDataSourceProperties()).thenReturn(dataSourceProperties); + + Assert.assertThrows(AnalysisException.class, () -> job.modifyProperties(command)); + Assert.assertEquals("stream-1", job.getStream()); + Mockito.verifyNoInteractions(dataSourceProperties); + } + @Test public void testShardRefreshShouldMoveRetiredParentToClosedUntilConsumed() throws Exception { KinesisRoutineLoadJob routineLoadJob = @@ -339,6 +409,35 @@ public void testDisplayCustomPropertiesMasksKinesisSecrets() { Assert.assertEquals("role_arn_value", showCreateCustomProperties.get("property.aws.role_arn")); } + private KinesisRoutineLoadJob createPausedJobWithInitialLoadDesc() { + KinesisRoutineLoadJob job = new KinesisRoutineLoadJob(1L, "kinesis_routine_load_job", 1L, + 1L, "ap-southeast-1", "stream-1", UserIdentity.ADMIN); + Deencapsulation.setField(job, "state", RoutineLoadJob.JobState.PAUSED); + Deencapsulation.setField(job, "createTimestamp", 123L); + job.setRoutineLoadDesc(new RoutineLoadDesc(new Separator("|", "|"), null, + null, null, null, null, null, LoadTask.MergeType.APPEND, null)); + return job; + } + + private void assertAlterResult(KinesisRoutineLoadJob job) { + Assert.assertEquals("|", job.getColumnSeparator().getSeparator()); + Assert.assertEquals("\n", job.getLineDelimiter().getSeparator()); + Assert.assertEquals("sequence_col", job.getSequenceCol()); + Assert.assertEquals((byte) '"', job.getEnclose()); + Assert.assertEquals((byte) '\\', job.getEscape()); + Assert.assertTrue(job.getEmptyFieldAsNull()); + } + + private String checkpointJson(RoutineLoadJob job) throws Exception { + ByteArrayOutputStream bytes = new ByteArrayOutputStream(); + try (DataOutputStream out = new DataOutputStream(bytes)) { + job.write(out); + } + try (DataInputStream in = new DataInputStream(new ByteArrayInputStream(bytes.toByteArray()))) { + return Text.readString(in); + } + } + private Set collectAssignedShards(KinesisRoutineLoadJob routineLoadJob) { List routineLoadTaskInfoList = Deencapsulation.getField(routineLoadJob, "routineLoadTaskInfoList"); diff --git a/fe/fe-core/src/test/java/org/apache/doris/load/routineload/RoutineLoadJobPersistenceTest.java b/fe/fe-core/src/test/java/org/apache/doris/load/routineload/RoutineLoadJobPersistenceTest.java new file mode 100644 index 00000000000000..5c221334e22254 --- /dev/null +++ b/fe/fe-core/src/test/java/org/apache/doris/load/routineload/RoutineLoadJobPersistenceTest.java @@ -0,0 +1,437 @@ +// 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.load.routineload; + +import org.apache.doris.analysis.BinaryPredicate; +import org.apache.doris.analysis.Expr; +import org.apache.doris.analysis.ImportColumnDesc; +import org.apache.doris.analysis.IntLiteral; +import org.apache.doris.analysis.Separator; +import org.apache.doris.analysis.SlotRef; +import org.apache.doris.analysis.UserIdentity; +import org.apache.doris.catalog.Database; +import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.OlapTable; +import org.apache.doris.catalog.Table; +import org.apache.doris.catalog.info.PartitionNamesInfo; +import org.apache.doris.common.io.Text; +import org.apache.doris.common.jmockit.Deencapsulation; +import org.apache.doris.datasource.CatalogMgr; +import org.apache.doris.datasource.InternalCatalog; +import org.apache.doris.datasource.property.fileformat.CsvFileFormatProperties; +import org.apache.doris.load.RoutineLoadDesc; +import org.apache.doris.load.loadv2.LoadTask; +import org.apache.doris.load.routineload.kafka.KafkaConfiguration; +import org.apache.doris.load.routineload.kafka.KafkaRoutineLoadJob; +import org.apache.doris.load.routineload.kinesis.KinesisRoutineLoadJob; +import org.apache.doris.nereids.load.NereidsRoutineLoadTaskInfo; +import org.apache.doris.nereids.trees.plans.commands.info.CreateRoutineLoadInfo; +import org.apache.doris.qe.OriginStatement; +import org.apache.doris.thrift.TPartialUpdateNewRowPolicy; +import org.apache.doris.thrift.TUniqueKeyUpdateMode; + +import com.google.common.collect.ImmutableMap; +import com.google.common.collect.Lists; +import com.google.common.collect.Maps; +import com.google.gson.JsonObject; +import com.google.gson.JsonParser; +import org.junit.Assert; +import org.junit.Test; +import org.mockito.MockedStatic; +import org.mockito.Mockito; + +import java.io.ByteArrayInputStream; +import java.io.ByteArrayOutputStream; +import java.io.DataInputStream; +import java.io.DataOutputStream; +import java.io.IOException; +import java.io.InputStream; +import java.nio.charset.StandardCharsets; +import java.util.Base64; +import java.util.List; +import java.util.Map; +import java.util.Optional; + +public class RoutineLoadJobPersistenceTest { + private static final String LEGACY_IMAGE = + "/upgrade/routine-load/a8928245/routine-load-kafka-image.b64"; + + @Test + public void testDirectStateImageRoundTripDoesNotParseOrigStmt() throws Exception { + KafkaRoutineLoadJob job = new KafkaRoutineLoadJob(1001L, "direct_job", 1002L, + 1003L, "127.0.0.1:9092", "direct_topic", UserIdentity.ADMIN); + job.state = RoutineLoadJob.JobState.PAUSED; + job.origStmt = new OriginStatement("this is deliberately not valid SQL", 0); + + Separator columnSeparator = analyzedSeparator("\\x01"); + Separator lineDelimiter = analyzedSeparator("\\n"); + List columns = Lists.newArrayList( + new ImportColumnDesc("source_col"), + new ImportColumnDesc("mapped_col", new IntLiteral(7L))); + Expr precedingFilter = predicate(BinaryPredicate.Operator.GT, "source_col", 1L); + Expr whereExpr = predicate(BinaryPredicate.Operator.LE, "mapped_col", 10L); + Expr deleteCondition = predicate(BinaryPredicate.Operator.EQ, "delete_flag", 1L); + PartitionNamesInfo partitions = new PartitionNamesInfo(false, Lists.newArrayList("p1", "p2")); + job.setRoutineLoadDesc(new RoutineLoadDesc(columnSeparator, lineDelimiter, columns, + precedingFilter, whereExpr, partitions, deleteCondition, LoadTask.MergeType.MERGE, "seq_col")); + + job.desireTaskConcurrentNum = 5; + job.maxErrorNum = 17L; + job.maxBatchIntervalS = 23L; + job.maxBatchRows = 300001L; + job.maxBatchSizeBytes = 104857601L; + job.execMemLimit = 345678901L; + job.maxFilterRatio = 0.99; + job.sendBatchParallelism = 99; + job.loadToSingleTablet = false; + job.memtableOnSinkNode = true; + + Map jobProperties = Maps.newHashMap(); + jobProperties.put(CreateRoutineLoadInfo.MAX_FILTER_RATIO_PROPERTY, "0.25"); + jobProperties.put(CreateRoutineLoadInfo.SEND_BATCH_PARALLELISM, "4"); + jobProperties.put(CreateRoutineLoadInfo.LOAD_TO_SINGLE_TABLET, "true"); + jobProperties.put(CreateRoutineLoadInfo.UNIQUE_KEY_UPDATE_MODE, "UPDATE_FIXED_COLUMNS"); + jobProperties.put(CreateRoutineLoadInfo.PARTIAL_COLUMNS, "true"); + jobProperties.put(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY, "ERROR"); + jobProperties.put(CsvFileFormatProperties.PROP_ENCLOSE, "\""); + jobProperties.put(CsvFileFormatProperties.PROP_ESCAPE, "\\"); + jobProperties.put(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL, "true"); + job.jobProperties = jobProperties; + + JsonObject json = imageJson(job); + Assert.assertTrue(json.has("ostmt")); + Assert.assertEquals(LoadTask.MergeType.MERGE.name(), json.get("mt").getAsString()); + for (String key : Lists.newArrayList( + "pni", "cds", "pf", "we", "cs", "lidel", "sc", "mt", "dc", "eml", "mosn")) { + Assert.assertTrue("missing direct-state key " + key, json.has(key)); + } + Assert.assertFalse(json.has("ld")); + Assert.assertEquals("\\x01", json.getAsJsonObject("cs").get("os").getAsString()); + Assert.assertEquals("\u0001", json.getAsJsonObject("cs").get("s").getAsString()); + Assert.assertEquals("\\n", json.getAsJsonObject("lidel").get("os").getAsString()); + Assert.assertEquals("\n", json.getAsJsonObject("lidel").get("s").getAsString()); + Assert.assertEquals(2, json.getAsJsonObject("cds").getAsJsonArray("des").size()); + + RoutineLoadJob restored; + try (MockedStatic envStatic = Mockito.mockStatic(Env.class)) { + restored = imageRoundTrip(job); + envStatic.verifyNoInteractions(); + } + + Assert.assertEquals(RoutineLoadJob.JobState.PAUSED, restored.getState()); + Assert.assertEquals(Lists.newArrayList("p1", "p2"), restored.getPartitionNamesInfo().getPartitionNames()); + Assert.assertEquals(2, restored.columnDescs.descs.size()); + Assert.assertEquals("source_col", restored.columnDescs.descs.get(0).getColumnName()); + Assert.assertEquals("mapped_col", restored.columnDescs.descs.get(1).getColumnName()); + Assert.assertNotNull(restored.columnDescs.descs.get(1).getExpr()); + Assert.assertNotNull(restored.getPrecedingFilter()); + Assert.assertNotNull(restored.getWhereExpr()); + Assert.assertEquals("\\x01", restored.getColumnSeparator().getOriSeparator()); + Assert.assertEquals("\u0001", restored.getColumnSeparator().getSeparator()); + Assert.assertEquals("\\n", restored.getLineDelimiter().getOriSeparator()); + Assert.assertEquals("\n", restored.getLineDelimiter().getSeparator()); + Assert.assertEquals("seq_col", restored.getSequenceCol()); + Assert.assertEquals(LoadTask.MergeType.MERGE, restored.getMergeType()); + Assert.assertNotNull(restored.getDeleteCondition()); + Assert.assertEquals(345678901L, restored.getMemLimit()); + Assert.assertTrue(restored.isMemtableOnSinkNode()); + Assert.assertEquals(5, restored.desireTaskConcurrentNum); + Assert.assertEquals(17L, restored.maxErrorNum); + Assert.assertEquals(23L, restored.getMaxBatchIntervalS()); + Assert.assertEquals(300001L, restored.getMaxBatchRows()); + Assert.assertEquals(104857601L, restored.getMaxBatchSizeBytes()); + + NereidsRoutineLoadTaskInfo taskInfo = restored.toNereidsRoutineLoadTaskInfo(); + Assert.assertEquals(345678901L, taskInfo.getMemLimit()); + Assert.assertEquals(0.25, taskInfo.getMaxFilterRatio(), 0.0); + Assert.assertEquals(4, taskInfo.getSendBatchParallelism()); + Assert.assertTrue(taskInfo.isLoadToSingleTablet()); + Assert.assertEquals(TUniqueKeyUpdateMode.UPDATE_FIXED_COLUMNS, taskInfo.getUniqueKeyUpdateMode()); + Assert.assertTrue(taskInfo.isFixedPartialUpdate()); + Assert.assertEquals(TPartialUpdateNewRowPolicy.ERROR, taskInfo.getPartialUpdateNewRowPolicy()); + Assert.assertEquals((byte) '"', taskInfo.getEnclose()); + Assert.assertEquals((byte) '\\', taskInfo.getEscape()); + Assert.assertTrue(taskInfo.getEmptyFieldAsNull()); + Assert.assertTrue(taskInfo.isMemtableOnSinkNode()); + Assert.assertEquals(LoadTask.MergeType.MERGE, taskInfo.getMergeType()); + Assert.assertNotNull(taskInfo.getDeleteCondition()); + Assert.assertEquals("seq_col", taskInfo.getSequenceCol()); + Assert.assertEquals(Lists.newArrayList("p1", "p2"), + taskInfo.getPartitionNamesInfo().getPartitionNames()); + Assert.assertEquals(2, taskInfo.getColumnExprDescs().descs.size()); + Assert.assertNotNull(taskInfo.getPrecedingFilter()); + Assert.assertNotNull(taskInfo.getWhereExpr()); + Assert.assertEquals("\u0001", taskInfo.getColumnSeparator().getSeparator()); + Assert.assertEquals("\n", taskInfo.getLineDelimiter().getSeparator()); + } + + @Test + public void testDirectStateImageWithNoLoadClausesDoesNotFallback() throws Exception { + KafkaRoutineLoadJob job = new KafkaRoutineLoadJob(2001L, "empty_job", 2002L, + 2003L, "127.0.0.1:9092", "empty_topic", UserIdentity.ADMIN); + job.state = RoutineLoadJob.JobState.PAUSED; + job.origStmt = new OriginStatement("also not valid SQL", 0); + + JsonObject json = imageJson(job); + Assert.assertTrue(json.has("ostmt")); + Assert.assertEquals(LoadTask.MergeType.APPEND.name(), json.get("mt").getAsString()); + for (String key : Lists.newArrayList("pni", "cds", "pf", "we", "cs", "lidel", "sc", "dc")) { + Assert.assertFalse("unexpected nullable direct-state key " + key, json.has(key)); + } + + RoutineLoadJob restored; + try (MockedStatic envStatic = Mockito.mockStatic(Env.class)) { + restored = imageRoundTrip(job); + envStatic.verifyNoInteractions(); + } + + Assert.assertEquals(RoutineLoadJob.JobState.PAUSED, restored.getState()); + Assert.assertNull(restored.getPartitionNamesInfo()); + Assert.assertNull(restored.columnDescs); + Assert.assertNull(restored.getPrecedingFilter()); + Assert.assertNull(restored.getWhereExpr()); + Assert.assertNull(restored.getColumnSeparator()); + Assert.assertNull(restored.getLineDelimiter()); + Assert.assertNull(restored.getSequenceCol()); + Assert.assertNull(restored.getDeleteCondition()); + Assert.assertEquals(LoadTask.MergeType.APPEND, restored.getMergeType()); + } + + @Test + public void testLegacyImageMigratesOnce() throws Exception { + Env env = Mockito.mock(Env.class); + CatalogMgr catalogMgr = Mockito.mock(CatalogMgr.class); + InternalCatalog catalog = Mockito.mock(InternalCatalog.class); + Database database = Mockito.mock(Database.class); + OlapTable table = Mockito.mock(OlapTable.class); + Mockito.when(env.getInternalCatalog()).thenReturn(catalog); + Mockito.when(env.getCatalogMgr()).thenReturn(catalogMgr); + Mockito.when(catalogMgr.getCatalog(Mockito.anyString())).thenReturn(catalog); + Mockito.when(catalog.getDb(8001L)).thenReturn(Optional.of(database)); + Mockito.when(catalog.getDb("legacy_db")).thenReturn(Optional.of(database)); + Mockito.when(catalog.getDbOrAnalysisException("legacy_db")).thenReturn(database); + Mockito.when(database.getName()).thenReturn("legacy_db"); + Mockito.when(database.getTable(9001L)).thenReturn(Optional.of((Table) table)); + Mockito.when(database.getTableOrAnalysisException("current_table")).thenReturn(table); + Mockito.when(table.getName()).thenReturn("current_table"); + Mockito.when(table.getType()).thenReturn(Table.TableType.OLAP); + Mockito.when(table.getEnableUniqueKeyMergeOnWrite()).thenReturn(true); + + byte[] legacyImage = loadBase64Fixture(LEGACY_IMAGE); + JsonObject legacyJson = imageJson(legacyImage); + Assert.assertFalse(legacyJson.has("mt")); + Assert.assertTrue(legacyJson.has("ostmt")); + + RoutineLoadJob migrated; + try (MockedStatic envStatic = Mockito.mockStatic(Env.class)) { + envStatic.when(Env::getCurrentEnv).thenReturn(env); + envStatic.when(Env::getCurrentInternalCatalog).thenReturn(catalog); + migrated = readImage(legacyImage); + } + + Assert.assertEquals(RoutineLoadJob.JobState.PAUSED, migrated.getState()); + Assert.assertEquals("|", migrated.getColumnSeparator().getOriSeparator()); + Assert.assertEquals("|", migrated.getColumnSeparator().getSeparator()); + Assert.assertNull(migrated.getSequenceCol()); + Assert.assertEquals(33554432L, migrated.getMemLimit()); + Assert.assertEquals(0.25, migrated.getMaxFilterRatio(), 0.0); + Assert.assertEquals(3, migrated.getSendBatchParallelism()); + Assert.assertTrue(migrated.isLoadToSingleTablet()); + Assert.assertEquals(TUniqueKeyUpdateMode.UPDATE_FIXED_COLUMNS, migrated.getUniqueKeyUpdateMode()); + Assert.assertTrue(migrated.isFixedPartialUpdate()); + Assert.assertEquals(TPartialUpdateNewRowPolicy.ERROR, migrated.partialUpdateNewKeyPolicy); + Assert.assertEquals((byte) '"', migrated.getEnclose()); + Assert.assertEquals((byte) '\\', migrated.getEscape()); + Assert.assertTrue(migrated.getEmptyFieldAsNull()); + Assert.assertFalse(migrated.isMemtableOnSinkNode()); + + JsonObject migratedJson = imageJson(migrated); + Assert.assertTrue(migratedJson.has("ostmt")); + Assert.assertEquals(LoadTask.MergeType.APPEND.name(), migratedJson.get("mt").getAsString()); + Assert.assertTrue(migratedJson.has("cs")); + Assert.assertTrue(migratedJson.has("eml")); + Assert.assertTrue(migratedJson.has("mosn")); + migrated.origStmt = new OriginStatement("invalid after successful migration", 0); + + RoutineLoadJob restoredAgain; + try (MockedStatic envStatic = Mockito.mockStatic(Env.class)) { + restoredAgain = imageRoundTrip(migrated); + envStatic.verifyNoInteractions(); + } + Assert.assertEquals(RoutineLoadJob.JobState.PAUSED, restoredAgain.getState()); + Assert.assertEquals("|", restoredAgain.getColumnSeparator().getSeparator()); + Assert.assertEquals(33554432L, restoredAgain.getMemLimit()); + Assert.assertFalse(restoredAgain.isMemtableOnSinkNode()); + } + + @Test + public void testKafkaDerivedStateIsRebuiltFromDurableProperties() throws Exception { + KafkaRoutineLoadJob job = new KafkaRoutineLoadJob(3001L, "kafka_derived", 3002L, + 3003L, "127.0.0.1:9092", "derived_topic", UserIdentity.ADMIN); + job.origStmt = new OriginStatement("invalid SQL must stay unused", 0); + Map customProperties = Maps.newHashMap(); + customProperties.put("client.id", "durable-client"); + customProperties.put(KafkaConfiguration.KAFKA_ORIGIN_DEFAULT_OFFSETS.getName(), "OFFSET_BEGINNING"); + Deencapsulation.setField(job, "customProperties", customProperties); + Deencapsulation.setField(job, "customKafkaPartitions", Lists.newArrayList(9)); + Deencapsulation.setField(job, "currentKafkaPartitions", Lists.newArrayList(1, 2)); + Deencapsulation.setField(job, "convertedCustomProperties", + Maps.newHashMap(ImmutableMap.of("stale", "value"))); + Deencapsulation.setField(job, "cachedPartitionWithLatestOffsets", + Maps.newHashMap(ImmutableMap.of(1, 100L))); + Deencapsulation.setField(job, "newCurrentKafkaPartition", Lists.newArrayList(3)); + Deencapsulation.setField(job, "kafkaDefaultOffSet", "OFFSET_END"); + + JsonObject json = imageJson(job); + Assert.assertEquals("127.0.0.1:9092", json.get("bl").getAsString()); + Assert.assertEquals("derived_topic", json.get("tp").getAsString()); + Assert.assertEquals("durable-client", json.getAsJsonObject("prop").get("client.id").getAsString()); + Assert.assertEquals(1, json.getAsJsonArray("cskp").size()); + assertNoJavaFieldNames(json, "currentKafkaPartitions", "convertedCustomProperties", + "cachedPartitionWithLatestOffsets", "newCurrentKafkaPartition", "kafkaDefaultOffSet"); + + KafkaRoutineLoadJob restored = (KafkaRoutineLoadJob) imageRoundTrip(job); + Assert.assertEquals("127.0.0.1:9092", restored.getBrokerList()); + Assert.assertEquals("derived_topic", restored.getTopic()); + Assert.assertEquals(Lists.newArrayList(9), Deencapsulation.getField(restored, "customKafkaPartitions")); + Assert.assertTrue(((List) Deencapsulation.getField(restored, "currentKafkaPartitions")).isEmpty()); + Assert.assertTrue(restored.getConvertedCustomProperties().isEmpty()); + Assert.assertTrue(((Map) Deencapsulation.getField( + restored, "cachedPartitionWithLatestOffsets")).isEmpty()); + Assert.assertEquals("", Deencapsulation.getField(restored, "kafkaDefaultOffSet")); + + Env env = Mockito.mock(Env.class); + try (MockedStatic envStatic = Mockito.mockStatic(Env.class)) { + envStatic.when(Env::getCurrentEnv).thenReturn(env); + restored.prepare(); + } + Assert.assertEquals("durable-client", restored.getConvertedCustomProperties().get("client.id")); + Assert.assertFalse(restored.getConvertedCustomProperties().containsKey( + KafkaConfiguration.KAFKA_ORIGIN_DEFAULT_OFFSETS.getName())); + Assert.assertEquals("OFFSET_BEGINNING", Deencapsulation.getField(restored, "kafkaDefaultOffSet")); + } + + @Test + public void testKinesisDerivedStateIsRebuiltFromDurableProperties() throws Exception { + KinesisRoutineLoadJob job = new KinesisRoutineLoadJob(4001L, "kinesis_derived", 4002L, + 4003L, "us-east-1", "derived_stream", UserIdentity.ADMIN); + job.origStmt = new OriginStatement("invalid SQL must stay unused", 0); + Deencapsulation.setField(job, "endpoint", "https://kinesis.example.test"); + Map customProperties = Maps.newHashMap(); + customProperties.put("client.setting", "durable-value"); + customProperties.put("kinesis_default_pos", "TRIM_HORIZON"); + Deencapsulation.setField(job, "customProperties", customProperties); + Deencapsulation.setField(job, "customKinesisShards", Lists.newArrayList("custom-shard")); + Deencapsulation.setField(job, "openKinesisShards", Lists.newArrayList("open-shard")); + Deencapsulation.setField(job, "closedKinesisShards", Lists.newArrayList("closed-shard")); + Deencapsulation.setField(job, "convertedCustomProperties", + Maps.newHashMap(ImmutableMap.of("stale", "value"))); + Deencapsulation.setField(job, "cachedShardWithMillsBehindLatest", + Maps.newHashMap(ImmutableMap.of("open-shard", 99L))); + Deencapsulation.setField(job, "newCurrentKinesisShards", Lists.newArrayList("new-shard")); + Deencapsulation.setField(job, "kinesisDefaultPosition", "LATEST"); + + JsonObject json = imageJson(job); + Assert.assertEquals("us-east-1", json.get("rg").getAsString()); + Assert.assertEquals("derived_stream", json.get("stm").getAsString()); + Assert.assertEquals("https://kinesis.example.test", json.get("ep").getAsString()); + Assert.assertEquals("durable-value", + json.getAsJsonObject("prop").get("client.setting").getAsString()); + Assert.assertEquals("custom-shard", json.getAsJsonArray("csks").get(0).getAsString()); + Assert.assertEquals("open-shard", json.getAsJsonArray("opks").get(0).getAsString()); + Assert.assertEquals("closed-shard", json.getAsJsonArray("clks").get(0).getAsString()); + assertNoJavaFieldNames(json, "convertedCustomProperties", "cachedShardWithMillsBehindLatest", + "newCurrentKinesisShards", "kinesisDefaultPosition"); + + KinesisRoutineLoadJob restored = (KinesisRoutineLoadJob) imageRoundTrip(job); + Assert.assertEquals("us-east-1", restored.getRegion()); + Assert.assertEquals("derived_stream", restored.getStream()); + Assert.assertEquals("https://kinesis.example.test", restored.getEndpoint()); + Assert.assertEquals(Lists.newArrayList("custom-shard"), + Deencapsulation.getField(restored, "customKinesisShards")); + Assert.assertEquals(Lists.newArrayList("open-shard"), + Deencapsulation.getField(restored, "openKinesisShards")); + Assert.assertEquals(Lists.newArrayList("closed-shard"), + Deencapsulation.getField(restored, "closedKinesisShards")); + Assert.assertTrue(restored.getConvertedCustomProperties().isEmpty()); + Assert.assertTrue(((Map) Deencapsulation.getField( + restored, "cachedShardWithMillsBehindLatest")).isEmpty()); + Assert.assertTrue(((List) Deencapsulation.getField(restored, "newCurrentKinesisShards")).isEmpty()); + Assert.assertEquals("", Deencapsulation.getField(restored, "kinesisDefaultPosition")); + + restored.prepare(); + Assert.assertEquals("durable-value", restored.getConvertedCustomProperties().get("client.setting")); + Assert.assertEquals("TRIM_HORIZON", + restored.getConvertedCustomProperties().get("kinesis_default_pos")); + Assert.assertEquals("TRIM_HORIZON", Deencapsulation.getField(restored, "kinesisDefaultPosition")); + } + + private static Separator analyzedSeparator(String value) throws Exception { + Separator separator = new Separator(value); + separator.analyze(); + return separator; + } + + private static Expr predicate(BinaryPredicate.Operator operator, String column, long value) { + return new BinaryPredicate(operator, new SlotRef(null, column), new IntLiteral(value)); + } + + private static JsonObject imageJson(RoutineLoadJob job) throws IOException { + return imageJson(writeImage(job)); + } + + private static JsonObject imageJson(byte[] image) throws IOException { + try (DataInputStream in = new DataInputStream(new ByteArrayInputStream(image))) { + return JsonParser.parseString(Text.readString(in)).getAsJsonObject(); + } + } + + private static RoutineLoadJob imageRoundTrip(RoutineLoadJob job) throws IOException { + return readImage(writeImage(job)); + } + + private static byte[] writeImage(RoutineLoadJob job) throws IOException { + ByteArrayOutputStream bytes = new ByteArrayOutputStream(); + try (DataOutputStream out = new DataOutputStream(bytes)) { + job.write(out); + } + return bytes.toByteArray(); + } + + private static RoutineLoadJob readImage(byte[] image) throws IOException { + try (DataInputStream in = new DataInputStream(new ByteArrayInputStream(image))) { + return RoutineLoadJob.read(in); + } + } + + private static byte[] loadBase64Fixture(String resource) throws IOException { + try (InputStream in = RoutineLoadJobPersistenceTest.class.getResourceAsStream(resource)) { + if (in == null) { + throw new IOException("missing fixture " + resource); + } + String base64 = new String(in.readAllBytes(), StandardCharsets.UTF_8).trim(); + return Base64.getDecoder().decode(base64); + } + } + + private static void assertNoJavaFieldNames(JsonObject json, String... fieldNames) { + for (String fieldName : fieldNames) { + Assert.assertFalse("derived field leaked into image: " + fieldName, json.has(fieldName)); + } + } +} diff --git a/fe/fe-core/src/test/java/org/apache/doris/persist/AlterRoutineLoadOperationLogTest.java b/fe/fe-core/src/test/java/org/apache/doris/persist/AlterRoutineLoadOperationLogTest.java index 8a1550d48f5d13..058cf42f949d54 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/persist/AlterRoutineLoadOperationLogTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/persist/AlterRoutineLoadOperationLogTest.java @@ -17,35 +17,43 @@ package org.apache.doris.persist; +import org.apache.doris.analysis.BinaryPredicate; +import org.apache.doris.analysis.ImportColumnDesc; +import org.apache.doris.analysis.IntLiteral; +import org.apache.doris.analysis.Separator; +import org.apache.doris.analysis.StringLiteral; +import org.apache.doris.catalog.info.PartitionNamesInfo; import org.apache.doris.common.UserException; import org.apache.doris.common.util.TimeUtils; +import org.apache.doris.load.RoutineLoadDesc; +import org.apache.doris.load.loadv2.LoadTask; import org.apache.doris.load.routineload.kafka.KafkaConfiguration; import org.apache.doris.load.routineload.kafka.KafkaDataSourceProperties; import org.apache.doris.nereids.trees.plans.commands.info.CreateRoutineLoadInfo; +import org.apache.doris.persist.gson.GsonUtils; +import com.google.common.collect.Lists; import com.google.common.collect.Maps; import org.junit.Assert; import org.junit.Test; +import java.io.ByteArrayInputStream; +import java.io.ByteArrayOutputStream; import java.io.DataInputStream; import java.io.DataOutputStream; -import java.io.File; -import java.io.FileInputStream; -import java.io.FileOutputStream; import java.io.IOException; +import java.io.InputStream; +import java.nio.charset.StandardCharsets; +import java.util.Base64; +import java.util.List; import java.util.Map; public class AlterRoutineLoadOperationLogTest { - private static String fileName = "./AlterRoutineLoadOperationLogTest"; + private static final String A8928245_LEGACY_LOG = + "/upgrade/routine-load/a8928245/alter-routine-load-log.b64"; @Test public void testSerializeAlterRoutineLoadOperationLog() throws IOException, UserException { - // 1. Write objects to file - File file = new File(fileName); - file.createNewFile(); - file.deleteOnExit(); - DataOutputStream out = new DataOutputStream(new FileOutputStream(file)); - long jobId = 1000; Map jobProperties = Maps.newHashMap(); jobProperties.put(CreateRoutineLoadInfo.DESIRED_CONCURRENT_NUMBER_PROPERTY, "5"); @@ -60,16 +68,31 @@ public void testSerializeAlterRoutineLoadOperationLog() throws IOException, User routineLoadDataSourceProperties.setTimezone(TimeUtils.DEFAULT_TIME_ZONE); routineLoadDataSourceProperties.analyze(); + Separator columnSeparator = new Separator(",", "\\x2c"); + Separator lineDelimiter = new Separator("\n", "\\n"); + List columns = Lists.newArrayList( + new ImportColumnDesc("source_col"), + new ImportColumnDesc("mapped_col", new StringLiteral("mapped_value"))); + BinaryPredicate precedingFilter = new BinaryPredicate(BinaryPredicate.Operator.GT, + new IntLiteral(3L), new IntLiteral(2L)); + BinaryPredicate where = new BinaryPredicate(BinaryPredicate.Operator.EQ, + new StringLiteral("selected"), new StringLiteral("selected")); + PartitionNamesInfo partitions = new PartitionNamesInfo(true, Lists.newArrayList("p1", "p2")); + BinaryPredicate deleteCondition = new BinaryPredicate(BinaryPredicate.Operator.EQ, + new IntLiteral(1L), new IntLiteral(1L)); + RoutineLoadDesc routineLoadDesc = new RoutineLoadDesc(columnSeparator, lineDelimiter, columns, + precedingFilter, where, partitions, deleteCondition, LoadTask.MergeType.MERGE, "sequence_col"); AlterRoutineLoadJobOperationLog log = new AlterRoutineLoadJobOperationLog(jobId, - jobProperties, routineLoadDataSourceProperties); - log.write(out); - out.flush(); - out.close(); - - // 2. Read objects from file - DataInputStream in = new DataInputStream(new FileInputStream(file)); + jobProperties, routineLoadDataSourceProperties, routineLoadDesc); + ByteArrayOutputStream bytes = new ByteArrayOutputStream(); + try (DataOutputStream out = new DataOutputStream(bytes)) { + log.write(out); + } - AlterRoutineLoadJobOperationLog log2 = AlterRoutineLoadJobOperationLog.read(in); + AlterRoutineLoadJobOperationLog log2; + try (DataInputStream in = new DataInputStream(new ByteArrayInputStream(bytes.toByteArray()))) { + log2 = AlterRoutineLoadJobOperationLog.read(in); + } Assert.assertEquals(1, log2.getJobProperties().size()); Assert.assertEquals("5", log2.getJobProperties().get(CreateRoutineLoadInfo.DESIRED_CONCURRENT_NUMBER_PROPERTY)); KafkaDataSourceProperties kafkaDataSourceProperties = (KafkaDataSourceProperties) log2.getDataSourceProperties(); @@ -81,9 +104,46 @@ public void testSerializeAlterRoutineLoadOperationLog() throws IOException, User kafkaDataSourceProperties.getKafkaPartitionOffsets().get(0)); Assert.assertEquals(routineLoadDataSourceProperties.getKafkaPartitionOffsets().get(1), kafkaDataSourceProperties.getKafkaPartitionOffsets().get(1)); + RoutineLoadDesc restoredDesc = log2.getRoutineLoadDesc(); + Assert.assertEquals(",", restoredDesc.getColumnSeparator().getSeparator()); + Assert.assertEquals("\\x2c", restoredDesc.getColumnSeparator().getOriSeparator()); + Assert.assertEquals("\n", restoredDesc.getLineDelimiter().getSeparator()); + Assert.assertEquals("\\n", restoredDesc.getLineDelimiter().getOriSeparator()); + Assert.assertEquals(2, restoredDesc.getColumnsInfo().size()); + Assert.assertEquals("source_col", restoredDesc.getColumnsInfo().get(0).getColumnName()); + Assert.assertEquals("mapped_col", restoredDesc.getColumnsInfo().get(1).getColumnName()); + Assert.assertNotNull(restoredDesc.getColumnsInfo().get(1).getExpr()); + Assert.assertNotNull(restoredDesc.getPrecedingFilter()); + Assert.assertNotNull(restoredDesc.getFilter()); + Assert.assertTrue(restoredDesc.getPartitionNamesInfo().isTemp()); + Assert.assertEquals(Lists.newArrayList("p1", "p2"), + restoredDesc.getPartitionNamesInfo().getPartitionNames()); + Assert.assertNotNull(restoredDesc.getDeleteCondition()); + Assert.assertEquals(LoadTask.MergeType.MERGE, restoredDesc.getMergeType()); + Assert.assertEquals("sequence_col", restoredDesc.getSequenceColName()); + Assert.assertEquals(GsonUtils.GSON.toJson(routineLoadDesc), GsonUtils.GSON.toJson(restoredDesc)); + } - in.close(); + @Test + public void testDeserializeLegacyLogWithoutRoutineLoadDesc() throws IOException { + byte[] bytes = loadBase64Fixture(A8928245_LEGACY_LOG); + try (DataInputStream in = new DataInputStream(new ByteArrayInputStream(bytes))) { + AlterRoutineLoadJobOperationLog log = AlterRoutineLoadJobOperationLog.read(in); + Assert.assertEquals(7001L, log.getJobId()); + Assert.assertTrue(log.getJobProperties().isEmpty()); + Assert.assertNull(log.getDataSourceProperties()); + Assert.assertNull(log.getRoutineLoadDesc()); + } } + private static byte[] loadBase64Fixture(String resource) throws IOException { + try (InputStream in = AlterRoutineLoadOperationLogTest.class.getResourceAsStream(resource)) { + if (in == null) { + throw new IOException("missing fixture " + resource); + } + String base64 = new String(in.readAllBytes(), StandardCharsets.UTF_8).trim(); + return Base64.getDecoder().decode(base64); + } + } }