github-actions[bot] commented on code in PR #66634:
URL: https://github.com/apache/doris/pull/66634#discussion_r4228470515


##########
fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java:
##########
@@ -920,6 +938,7 @@ private void 
resetCloudProgress(Cloud.ResetRLProgressRequest.Builder builder) th
     public void replayModifyProperties(AlterRoutineLoadJobOperationLog log) {

Review Comment:
   [P1] Fail closed when source replay prevents the load clause from being 
applied. A PAUSED cloud Kafka job can accept one ALTER containing `WHERE` (or 
`COLUMNS`) plus Kafka source properties. During follower replay, 
`modifyPropertiesInternal()` invokes `resetCloudProgress()`; a MetaService RPC 
error throws `DdlException`, so execution skips this new call and the catch 
only logs. The job remains PAUSED with its old load definition, and a later 
RESUME journal record can schedule it after failover with different loading 
semantics from the accepted ALTER. Cancel the job on any replay failure before 
the definition is installed, or replay the complete ALTER atomically.



##########
fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadJob.java:
##########
@@ -2042,26 +2048,162 @@ public void gsonPostProcess() throws IOException {
                 ctx.cleanup();
             }
         } catch (Exception e) {
-            // Terminalize this unusable job as CANCELLED. Set endTimestamp 
only if unset (avoid
-            // refreshing on every image load) and keep the existing cancel 
reason if present.
-            state = JobState.CANCELLED;
-            routineLoadTaskInfoList.clear();
-            long failureTimestamp = System.currentTimeMillis();
-            if (endTimestamp == -1) {
-                endTimestamp = failureTimestamp;
-            }
-            if (cancelReason == null) {
-                cancelReason = new ErrorReason(InternalErrorCode.INTERNAL_ERR,
-                        "FE restart deserialize failed at " + 
TimeUtils.longToTimeString(failureTimestamp)
-                                + ": " + e.getMessage());
-            }
+            cancelUnrecoverableJob("FE restart deserialize failed", e);
             LOG.warn("error happens when parsing create routine load stmt: " + 
origStmt.originStmt, e);
         }
         if (userIdentity != null) {
             userIdentity.setIsAnalyzed();
         }
     }
 
+    // Terminalize a job whose persisted definition can not be restored as 
CANCELLED. Set endTimestamp
+    // only if unset (avoid refreshing on every image load) and keep the 
existing cancel reason if present.
+    private void cancelUnrecoverableJob(String reason, Exception e) {
+        state = JobState.CANCELLED;
+        routineLoadTaskInfoList.clear();
+        long failureTimestamp = System.currentTimeMillis();
+        if (endTimestamp == -1) {
+            endTimestamp = failureTimestamp;
+        }
+        if (cancelReason == null) {
+            cancelReason = new ErrorReason(InternalErrorCode.INTERNAL_ERR,
+                    reason + " at " + 
TimeUtils.longToTimeString(failureTimestamp) + ": " + e.getMessage());
+        }
+    }
+
+    protected void replayLoadDefinition(OriginStatement alterStatement, Long 
sqlMode,
+            Map<String, String> alterSessionVariables) {
+        if (alterStatement == null) {
+            return;
+        }
+        try {
+            Database database = 
Env.getCurrentEnv().getInternalCatalog().getDbOrMetaException(dbId);
+            ConnectContext ctx = createLoadDefinitionContext(database, 
sqlMode);
+            // Old journals have no snapshot and retain the existing recovery 
context.
+            
ctx.getSessionVariable().setAffectQueryResultInPlanSessionVariables(alterSessionVariables);
+            try {
+                ctx.setThreadLocalInfo();
+                AlterRoutineLoadCommand command =
+                        (AlterRoutineLoadCommand) 
parsePersistedStatement(alterStatement);
+                if (command.hasLoadProperty()) {
+                    RoutineLoadDesc loadDesc = 
mergeLoadDesc(command.analyzeLoadProperties(ctx, this));
+                    OriginStatement loadDefinitionStmt = 
buildLoadDefinitionStatement(loadDesc);
+                    applyLoadDefinition(loadDesc, loadDefinitionStmt,
+                            
ctx.getSessionVariable().getAffectQueryResultInPlanVariables(),
+                            ctx.getSessionVariable().getSqlMode());
+                }
+            } finally {
+                ctx.cleanup();
+            }
+        } catch (Exception e) {
+            // Like gsonPostProcess(), an ALTER that this FE can not analyze 
any more (for example, after an
+            // upgrade changed the analysis rules) must not stop journal 
replay. The job keeps its previous
+            // definition and is cancelled, so it never loads data with a 
definition other than the ALTERed one.
+            cancelUnrecoverableJob("FE replay alter routine load failed", e);
+            
Env.getCurrentGlobalTransactionMgr().getCallbackFactory().removeCallback(id);
+            LOG.warn("error happens when replaying alter routine load stmt of 
job {}, cancel it", id, e);
+        }
+    }
+
+    protected void replayLoadDefinition(OriginStatement alterStatement, Long 
sqlMode) {
+        replayLoadDefinition(alterStatement, sqlMode, null);
+    }
+
+    protected void replayLoadDefinition(OriginStatement alterStatement) {
+        replayLoadDefinition(alterStatement, null);
+    }
+
+    private ConnectContext createLoadDefinitionContext(Database database, Long 
sqlMode) {
+        ConnectContext ctx = new ConnectContext();
+        ctx.setDatabase(database.getName());
+        StatementContext statementContext = new StatementContext();
+        statementContext.setConnectContext(ctx);
+        ctx.setStatementContext(statementContext);
+        ctx.setEnv(Env.getCurrentEnv());
+        ctx.setCurrentUserIdentity(UserIdentity.ADMIN);
+        
ctx.getSessionVariable().setAffectQueryResultInPlanSessionVariables(sessionVariables);
+        if (sqlMode != null) {
+            ctx.getSessionVariable().setSqlMode(sqlMode);
+        } else if (sessionVariables.containsKey(SessionVariable.SQL_MODE)) {
+            
ctx.getSessionVariable().setSqlMode(Long.parseLong(sessionVariables.get(SessionVariable.SQL_MODE)));
+        }
+        ctx.getState().reset();
+        return ctx;
+    }
+
+    /**
+     * Overlay the load clauses of an ALTER on the current load definition 
without modifying this job.
+     * A clause replaces the current one only when the ALTER specifies it, the 
same as setRoutineLoadDesc().
+     */
+    protected RoutineLoadDesc mergeLoadDesc(RoutineLoadDesc alterLoadDesc) {
+        List<ImportColumnDesc> columns = alterLoadDesc.getColumnsInfo();
+        if (columns == null && columnDescs != null) {
+            columns = Lists.newArrayList(columnDescs.descs);
+        }
+        return new RoutineLoadDesc(
+                alterLoadDesc.getColumnSeparator() != null ? 
alterLoadDesc.getColumnSeparator() : columnSeparator,
+                alterLoadDesc.getLineDelimiter() != null ? 
alterLoadDesc.getLineDelimiter() : lineDelimiter,
+                columns,
+                alterLoadDesc.getPrecedingFilter() != null ? 
alterLoadDesc.getPrecedingFilter() : precedingFilter,
+                alterLoadDesc.getFilter() != null ? alterLoadDesc.getFilter() 
: whereExpr,
+                alterLoadDesc.getPartitionNamesInfo() != null
+                        ? alterLoadDesc.getPartitionNamesInfo() : 
partitionNamesInfo,
+                alterLoadDesc.getDeleteCondition() != null ? 
alterLoadDesc.getDeleteCondition() : deleteCondition,
+                alterLoadDesc.getMergeType(),
+                alterLoadDesc.hasSequenceCol() ? 
alterLoadDesc.getSequenceColName() : sequenceCol);
+    }
+
+    /**
+     * Build the CREATE statement persisted as origStmt for the given 
effective load definition, and check
+     * that it can be parsed back. This does not modify the job, so callers 
must finish every step that may
+     * fail before applying the result with applyLoadDefinition().
+     */
+    protected OriginStatement buildLoadDefinitionStatement(RoutineLoadDesc 
loadDesc) throws UserException {
+        StringBuilder sql = new StringBuilder("CREATE ROUTINE LOAD ")
+                .append(SqlUtils.getIdentSql(name));
+        if (!isMultiTable) {
+            sql.append(" ON ").append(SqlUtils.getIdentSql(getTableName()));
+        }
+        sql.append(" WITH ").append(loadDesc.getMergeType().name());
+        String loadClauseSql = loadDesc.toSql();
+        if (!loadClauseSql.isEmpty()) {

Review Comment:
   [P2] Validate the merged CREATE definition before persisting it. An APPEND 
job can accept `ALTER ROUTINE LOAD FOR j DELETE ON flag = 1`: ALTER analysis 
translates the clause but does not enforce the CREATE merge/delete check, and 
this line only checks that the generated `WITH APPEND ... DELETE ON` parses. On 
the next image load, `CreateRoutineLoadInfo.validate()` rejects that same 
statement and `gsonPostProcess()` marks the previously accepted job CANCELLED. 
Apply the CREATE semantic checks before mutating or journaling the ALTER.



##########
fe/fe-core/src/main/java/org/apache/doris/load/RoutineLoadDesc.java:
##########
@@ -97,6 +138,59 @@ public boolean hasSequenceCol() {
         return !Strings.isNullOrEmpty(sequenceColName);
     }
 
+    /**
+     * Convert the effective load clauses to SQL so they can be persisted in 
RoutineLoadJob.origStmt.
+     */
+    public String toSql() {
+        List<String> clauses = new ArrayList<>();
+        // Routine Load SQL does not currently expose a line-delimiter clause.
+        if (columnSeparator != null) {
+            // oriSeparator is already the encoded spelling consumed by 
Separator.convertSeparator().
+            // Escaping its backslashes again would turn \t and \x01 into 
literal backslash sequences.
+            String separator = columnSeparator.getOriSeparator();
+            // Keep the encoded spelling intact, including escaped or doubled 
quotes. A single quote
+            // in the spelling does not require double quoting when the SQL 
lexer already accepts it.
+            Pattern singleQuotedSeparator = 
SqlModeHelper.hasNoBackSlashEscapes()
+                    ? SINGLE_QUOTED_SEPARATOR_NO_BACKSLASH_ESCAPES : 
SINGLE_QUOTED_SEPARATOR;
+            String quote = singleQuotedSeparator.matcher(separator).matches() 
? "'" : "\"";
+            clauses.add("COLUMNS TERMINATED BY " + quote + separator + quote);
+        }
+        if (columnsInfo != null) {
+            clauses.add("COLUMNS(" + columnsInfo.stream()
+                    .map(this::columnToSql)
+                    .collect(Collectors.joining(", ")) + ")");
+        }
+        if (precedingFilter != null) {
+            clauses.add("PRECEDING FILTER " + precedingFilter.accept(
+                    PERSISTED_EXPR_TO_SQL_VISITOR, ToSqlParams.WITHOUT_TABLE));
+        }
+        if (filter != null) {
+            clauses.add("WHERE " + 
filter.accept(PERSISTED_EXPR_TO_SQL_VISITOR, ToSqlParams.WITHOUT_TABLE));
+        }
+        if (partitionNamesInfo != null) {
+            String prefix = partitionNamesInfo.isTemp() ? "TEMPORARY 
PARTITION(" : "PARTITION(";
+            clauses.add(prefix + 
partitionNamesInfo.getPartitionNames().stream()
+                    .map(SqlUtils::getIdentSql)
+                    .collect(Collectors.joining(", ")) + ")");
+        }
+        if (deleteCondition != null) {
+            clauses.add("DELETE ON " + deleteCondition.accept(
+                    PERSISTED_EXPR_TO_SQL_VISITOR, ToSqlParams.WITHOUT_TABLE));
+        }
+        if (hasSequenceCol()) {
+            clauses.add("ORDER BY " + SqlUtils.getIdentSql(sequenceColName));
+        }
+        return String.join(", ", clauses);
+    }
+
+    private String columnToSql(ImportColumnDesc columnDesc) {
+        String sql = SqlUtils.getIdentSql(columnDesc.getColumnName());
+        if (columnDesc.getExpr() != null) {

Review Comment:
   [P2] Quote lambda parameter names when persisting inherited COLUMNS 
expressions. A valid mapping such as `COLUMNS(mapped = array_map(`x-y` -> `x-y` 
+ 1, arr))` is translated into a legacy lambda whose parameter name is `x-y`. 
The SQL visitor used here emits `array_map(x-y -> x-y + 1, arr)` without 
backticks, so an unrelated load-clause ALTER fails 
`buildLoadDefinitionStatement()` parsing and cannot be applied. Preserve 
identifier quoting for the lambda declaration and references, then cover this 
mapping in an ALTER test.



##########
fe/fe-core/src/main/java/org/apache/doris/load/RoutineLoadDesc.java:
##########
@@ -97,6 +138,59 @@ public boolean hasSequenceCol() {
         return !Strings.isNullOrEmpty(sequenceColName);
     }
 
+    /**
+     * Convert the effective load clauses to SQL so they can be persisted in 
RoutineLoadJob.origStmt.
+     */
+    public String toSql() {
+        List<String> clauses = new ArrayList<>();
+        // Routine Load SQL does not currently expose a line-delimiter clause.
+        if (columnSeparator != null) {
+            // oriSeparator is already the encoded spelling consumed by 
Separator.convertSeparator().
+            // Escaping its backslashes again would turn \t and \x01 into 
literal backslash sequences.
+            String separator = columnSeparator.getOriSeparator();
+            // Keep the encoded spelling intact, including escaped or doubled 
quotes. A single quote
+            // in the spelling does not require double quoting when the SQL 
lexer already accepts it.
+            Pattern singleQuotedSeparator = 
SqlModeHelper.hasNoBackSlashEscapes()
+                    ? SINGLE_QUOTED_SEPARATOR_NO_BACKSLASH_ESCAPES : 
SINGLE_QUOTED_SEPARATOR;
+            String quote = singleQuotedSeparator.matcher(separator).matches() 
? "'" : "\"";
+            clauses.add("COLUMNS TERMINATED BY " + quote + separator + quote);
+        }
+        if (columnsInfo != null) {
+            clauses.add("COLUMNS(" + columnsInfo.stream()
+                    .map(this::columnToSql)
+                    .collect(Collectors.joining(", ")) + ")");
+        }
+        if (precedingFilter != null) {
+            clauses.add("PRECEDING FILTER " + precedingFilter.accept(
+                    PERSISTED_EXPR_TO_SQL_VISITOR, ToSqlParams.WITHOUT_TABLE));
+        }
+        if (filter != null) {

Review Comment:
   [P1] Preserve the explicit LIKE escape when serializing the load definition. 
A valid filter such as `WHERE city LIKE 'A!_%' ESCAPE '!'` becomes a 
three-argument legacy `like` expression, but the SQL visitor used here prints 
only the first two arguments. After an unrelated load-clause ALTER, the 
generated CREATE persists `LIKE 'A!_%'` without `ESCAPE`; image recovery then 
reverses which rows such as `A_one` and `A!one` match. Emit the third argument 
as `ESCAPE ...` and cover the restored predicate with discriminating rows.



##########
fe/fe-core/src/main/java/org/apache/doris/load/RoutineLoadDesc.java:
##########
@@ -97,6 +138,59 @@ public boolean hasSequenceCol() {
         return !Strings.isNullOrEmpty(sequenceColName);
     }
 
+    /**
+     * Convert the effective load clauses to SQL so they can be persisted in 
RoutineLoadJob.origStmt.
+     */
+    public String toSql() {
+        List<String> clauses = new ArrayList<>();
+        // Routine Load SQL does not currently expose a line-delimiter clause.
+        if (columnSeparator != null) {
+            // oriSeparator is already the encoded spelling consumed by 
Separator.convertSeparator().
+            // Escaping its backslashes again would turn \t and \x01 into 
literal backslash sequences.
+            String separator = columnSeparator.getOriSeparator();
+            // Keep the encoded spelling intact, including escaped or doubled 
quotes. A single quote
+            // in the spelling does not require double quoting when the SQL 
lexer already accepts it.
+            Pattern singleQuotedSeparator = 
SqlModeHelper.hasNoBackSlashEscapes()
+                    ? SINGLE_QUOTED_SEPARATOR_NO_BACKSLASH_ESCAPES : 
SINGLE_QUOTED_SEPARATOR;
+            String quote = singleQuotedSeparator.matcher(separator).matches() 
? "'" : "\"";
+            clauses.add("COLUMNS TERMINATED BY " + quote + separator + quote);
+        }
+        if (columnsInfo != null) {
+            clauses.add("COLUMNS(" + columnsInfo.stream()

Review Comment:
   [P2] Keep inherited map literals parseable when rebuilding the load 
definition. For a valid mapping such as `COLUMNS(mapped = element_at({'a': 1}, 
'a'))`, this path converts the Nereids map literal to a legacy `MapLiteral`, 
whose SQL visitor emits `MAP{'a': 1}`. The generated CREATE is then rejected by 
`buildLoadDefinitionStatement()`, because the Nereids grammar accepts map 
literals only as `{key:value}`. An unrelated load-clause ALTER therefore fails; 
render the Nereids spelling and cover a map literal through ALTER and recovery.



##########
fe/fe-core/src/main/java/org/apache/doris/persist/AlterRoutineLoadJobOperationLog.java:
##########
@@ -37,12 +38,37 @@ public class AlterRoutineLoadJobOperationLog  implements 
Writable {
     private Map<String, String> jobProperties;
     @SerializedName(value = "dataSourceProperties")
     private AbstractDataSourceProperties dataSourceProperties;

Review Comment:
   [P1] Fence this journal extension during rolling FE upgrades. The new load 
delta exists only in these optional JSON fields. A not-yet-upgraded electable 
follower ignores them in its old `replayModifyProperties()`, so after a new 
leader accepts `ALTER ROUTINE LOAD ... COLUMNS ...`, that follower can become 
master with the previous mapping and later RESUME the job with different 
semantics. This PR does not bump the metadata version or gate the ALTER on 
electable FE capability. Reject such ALTERs until all electable FEs can replay 
the fields, or use a journal compatibility fence that prevents old readers from 
promotion.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to