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]