This is an automated email from the ASF dual-hosted git repository.
JNSimba pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 225bcd3018b [improvement](streaming) Improve streaming job error
messages (#68241)
225bcd3018b is described below
commit 225bcd3018b01c16e4b2f22dd1a6919b02c87e0d
Author: wudi <[email protected]>
AuthorDate: Tue Sep 29 09:57:17 2026 +0800
[improvement](streaming) Improve streaming job error messages (#68241)
### What problem does this PR solve?
Problem Summary: Streaming Job failures often exposed generic validation
text, misleading HTTP status lines, or insufficient runtime context.
This change makes CREATE/ALTER validation and runtime failures report
the relevant parameter, task, source object, threshold, recovery action,
or remote error. In particular, CDC offset commit failures now preserve
the error returned in the HTTP 200 response body.
---
.../doris/httpv2/rest/StreamingJobAction.java | 2 +-
.../streaming/DataSourceConfigValidator.java | 24 ++++++-
.../insert/streaming/StreamingInsertJob.java | 10 ++-
.../insert/streaming/StreamingJobProperties.java | 10 ++-
.../streaming/StreamingJobSchedulerTask.java | 3 +-
.../insert/streaming/StreamingMultiTblTask.java | 7 +-
.../job/offset/jdbc/JdbcSourceOffsetProvider.java | 20 ++++--
.../apache/doris/job/util/StreamingJobUtils.java | 11 +--
.../trees/plans/commands/AlterJobCommand.java | 2 +-
.../trees/plans/commands/info/CreateJobInfo.java | 3 +-
.../CdcStreamTableValuedFunction.java | 13 ++--
.../StreamingInsertJobCheckDataQualityTest.java | 6 ++
.../streaming/StreamingJobPropertiesTest.java | 3 +-
.../JdbcSourceOffsetProviderErrorHandlingTest.java | 19 ++++-
.../doris/cdcclient/sink/DorisBatchStreamLoad.java | 11 +++
.../cdcclient/sink/DorisBatchStreamLoadTest.java | 84 ++++++++++++++++++++++
.../test_streaming_mysql_job_create_alter.groovy | 2 +-
.../cdc/test_streaming_mysql_job_dup.groovy | 2 +-
.../cdc/test_streaming_postgres_job_dup.groovy | 2 +-
.../test_streaming_insert_job_crud.groovy | 8 +--
.../test_streaming_job_max_retry.groovy | 2 +
21 files changed, 204 insertions(+), 40 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/StreamingJobAction.java
b/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/StreamingJobAction.java
index 42d243b5754..7eac63ddb77 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/StreamingJobAction.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/StreamingJobAction.java
@@ -113,7 +113,7 @@ public class StreamingJobAction extends RestBaseController {
private void checkAuth(HttpServletRequest request) {
String authToken = request.getHeader("token");
if (Strings.isNullOrEmpty(authToken)) {
- throw new UnauthorizedException("Miss token");
+ throw new UnauthorizedException("Missing token");
}
if (!checkClusterToken(authToken)) {
throw new UnauthorizedException("Invalid token");
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/DataSourceConfigValidator.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/DataSourceConfigValidator.java
index f75bf03c56e..f74feb336ed 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/DataSourceConfigValidator.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/DataSourceConfigValidator.java
@@ -127,7 +127,10 @@ public class DataSourceConfigValidator {
}
if (!isValidValue(key, value, dataSourceType)) {
- throw new IllegalArgumentException("Invalid value for key '" +
key + "': " + value);
+ if (DataSourceConfigKeys.SERVER_ID.equals(key)) {
+ validateServerIdConfig(input);
+ }
+ throw new IllegalArgumentException(invalidValueMessage(key,
value));
}
}
@@ -290,6 +293,25 @@ public class DataSourceConfigValidator {
return true;
}
+ public static String invalidValueMessage(String key, String value) {
+ String message = "Invalid value for key '" + key + "': " + value + ".";
+ if (DataSourceConfigKeys.SLOT_NAME.equals(key)
+ || DataSourceConfigKeys.PUBLICATION_NAME.equals(key)) {
+ return message + " Must match [a-z_][a-z0-9_]* (max 63
characters).";
+ }
+ if (DataSourceConfigKeys.SSL_MODE.equals(key)) {
+ return message + " Use disable, require, or verify-ca.";
+ }
+ if (DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE.equals(key)
+ || DataSourceConfigKeys.SNAPSHOT_PARALLELISM.equals(key)) {
+ return message + " Use a positive integer.";
+ }
+ if (DataSourceConfigKeys.SKIP_SNAPSHOT_BACKFILL.equals(key)) {
+ return message + " Use true or false.";
+ }
+ return message;
+ }
+
// Strict boolean: only "true"/"false" (case-insensitive);
Boolean.parseBoolean would
// silently coerce typos like "yes" to false.
public static boolean isValidBoolean(String value) {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
index a8cbcab6ac6..23b2a5ae66e 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
@@ -1236,7 +1236,7 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
private static boolean checkPrivilege(ConnectContext ctx, String sql)
throws AnalysisException {
LogicalPlan logicalPlan = new NereidsParser().parseSingle(sql);
if (!(logicalPlan instanceof InsertIntoTableCommand)) {
- throw new AnalysisException("Only support insert command");
+ throw new AnalysisException("Streaming jobs only support INSERT
statements");
}
LogicalPlan logicalQuery = ((InsertIntoTableCommand)
logicalPlan).getLogicalQuery();
List<String> targetTable =
InsertUtils.getTargetTableQualified(logicalQuery, ctx);
@@ -1562,7 +1562,8 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
&& runningMultiTask.isTimeout(status)) {
String timeoutReason = status == null ? "" :
status.getFailReason();
if (StringUtils.isEmpty(timeoutReason)) {
- timeoutReason = "task failed cause timeout";
+ timeoutReason = "Streaming task " +
runningMultiTask.getTaskId()
+ + " timed out because no progress was reported.";
}
runningMultiTask.onFail(timeoutReason);
// renew streaming task by auto resume
@@ -1660,7 +1661,10 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
getJobId(), ratio, maxFilterRatio,
sampleWindowFilteredRows, sampleWindowScannedRows);
log.error(msg);
FailureReason failureReason = new
FailureReason(InternalErrorCode.TOO_MANY_FAILURE_ROWS_ERR,
- "too many filtered rows exceeded max_filter_ratio " +
maxFilterRatio);
+ String.format(
+ "too many filtered rows: ratio %s exceeds
load.max_filter_ratio %s. "
+ + "Fix the source data or adjust the
limit, then run RESUME JOB.",
+ ratio, maxFilterRatio));
this.setFailureReason(failureReason);
this.updateJobStatus(JobStatus.PAUSED);
throw new JobException(failureReason.getMsg());
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobProperties.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobProperties.java
index 3db3aecaf8b..1bde1bf6205 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobProperties.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobProperties.java
@@ -106,18 +106,22 @@ public class StreamingJobProperties implements
JobProperties {
this.maxIntervalSecond = Util.getLongPropertyOrDefault(
properties.get(StreamingJobProperties.MAX_INTERVAL_SECOND_PROPERTY),
StreamingJobProperties.DEFAULT_MAX_INTERVAL_SECOND,
(v) -> v >= 1,
- StreamingJobProperties.MAX_INTERVAL_SECOND_PROPERTY + " should
> 1");
+ StreamingJobProperties.MAX_INTERVAL_SECOND_PROPERTY + " must
be at least 1 second, but was "
+ +
properties.get(StreamingJobProperties.MAX_INTERVAL_SECOND_PROPERTY));
this.s3BatchFiles = Util.getLongPropertyOrDefault(
properties.get(StreamingJobProperties.S3_MAX_BATCH_FILES_PROPERTY),
StreamingJobProperties.DEFAULT_MAX_S3_BATCH_FILES, (v)
-> v >= 1,
- StreamingJobProperties.S3_MAX_BATCH_FILES_PROPERTY + " should
>=1 ");
+ StreamingJobProperties.S3_MAX_BATCH_FILES_PROPERTY + " must be
at least 1, but was "
+ +
properties.get(StreamingJobProperties.S3_MAX_BATCH_FILES_PROPERTY));
this.s3BatchBytes = Util.getLongPropertyOrDefault(
properties.get(StreamingJobProperties.S3_MAX_BATCH_BYTES_PROPERTY),
StreamingJobProperties.DEFAULT_MAX_S3_BATCH_BYTES, (v)
-> v >= 100 * 1024 * 1024
&& v <= (long) (1024 * 1024 * 1024) * 10,
- StreamingJobProperties.S3_MAX_BATCH_BYTES_PROPERTY + " should
between 100MB and 10GB");
+ StreamingJobProperties.S3_MAX_BATCH_BYTES_PROPERTY
+ + " must be between 100 MB and 10 GB, but was "
+ +
properties.get(StreamingJobProperties.S3_MAX_BATCH_BYTES_PROPERTY));
// validate session variables
try {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobSchedulerTask.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobSchedulerTask.java
index c5dfacaadfa..0a4b7ca09f6 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobSchedulerTask.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobSchedulerTask.java
@@ -113,7 +113,8 @@ public class StreamingJobSchedulerTask extends AbstractTask
{
streamingInsertJob.setFailureReason(new FailureReason(
InternalErrorCode.CANNOT_RESUME_ERR,
"Auto resume failed after " + autoResumeCount
- + " attempts. Last error: " +
failureReason.getMsg()));
+ + " attempts. Last error: " +
failureReason.getMsg()
+ + ". Run RESUME JOB to retry."));
return;
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingMultiTblTask.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingMultiTblTask.java
index e652f7c9e1e..f141762bdf8 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingMultiTblTask.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingMultiTblTask.java
@@ -163,18 +163,21 @@ public class StreamingMultiTblTask extends
AbstractStreamingTask {
log.info("Send write records request successfully,
response: {}", responseObj.getData());
return;
}
+ String errorMessage =
StringUtils.defaultIfBlank(responseObj.getData(), responseObj.getMsg());
+ throw new JobException(StringUtils.defaultIfBlank(
+ errorMessage, "cdc_client failed to start streaming
write"));
} catch (JsonProcessingException e) {
log.warn("Failed to parse write records response: {}",
response);
throw new JobException("Failed to parse write records
response: " + response);
}
- throw new JobException("Failed to send write records request ,
error message: " + response);
} catch (TimeoutException te) {
log.warn("cdc_client RPC timeout api=/api/writeRecords taskId={}
jobId={} backend={}:{} timeout_sec={}",
taskId, getJobId(), backend.getHost(),
backend.getBrpcPort(),
Config.streaming_cdc_heavy_rpc_timeout_sec);
// the request may have been dispatched and still running remotely
noRetry = true;
- throw new JobException("cdc_client RPC timeout: /api/writeRecords
taskId=" + taskId);
+ throw new JobException("cdc_client RPC timeout: /api/writeRecords
jobId="
+ + getJobId() + " taskId=" + taskId);
} catch (ExecutionException | InterruptedException ex) {
if (ex instanceof InterruptedException) {
Thread.currentThread().interrupt();
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java
index 64752144853..16e4d109159 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java
@@ -59,6 +59,7 @@ import lombok.Setter;
import lombok.extern.log4j.Log4j2;
import org.apache.commons.collections4.CollectionUtils;
import org.apache.commons.collections4.MapUtils;
+import org.apache.commons.lang3.StringUtils;
import java.util.ArrayList;
import java.util.Arrays;
@@ -960,9 +961,9 @@ public class JdbcSourceOffsetProvider implements
SourceOffsetProvider {
/**
* Decode a remote response envelope. A failure is returned as {@code
{code:1, data:"<message>"}}
* over HTTP 200, while success carries the typed payload in {@code data}.
Decode the envelope
- * with a lenient {@link JsonNode} data field first so a failure throws
the raw response (which
- * carries the real error in {@code data}) instead of a misleading
type-mismatch from forcing the
- * success type onto an error string. Package-private for unit testing.
+ * with a lenient {@link JsonNode} data field first so a failure surfaces
the error in
+ * {@code data} instead of a misleading type-mismatch from forcing the
success type onto an
+ * error string. Package-private for unit testing.
*/
<T> T parseCdcResponseData(String response, TypeReference<T> dataType)
throws JobException {
ResponseBody<JsonNode> body;
@@ -972,6 +973,13 @@ public class JdbcSourceOffsetProvider implements
SourceOffsetProvider {
throw new JobException(response);
}
if (body.getCode() != RestApiStatusCode.OK.code) {
+ JsonNode data = body.getData();
+ if (data != null && data.isTextual() &&
StringUtils.isNotBlank(data.asText())) {
+ throw new JobException(data.asText());
+ }
+ if (StringUtils.isNotBlank(body.getMsg())) {
+ throw new JobException(body.getMsg());
+ }
throw new JobException(response);
}
try {
@@ -1102,12 +1110,12 @@ public class JdbcSourceOffsetProvider implements
SourceOffsetProvider {
if (responseObj.getCode() == RestApiStatusCode.OK.code) {
log.info("Init {} source reader successfully, response:
{}", getJobId(), responseObj.getData());
return;
- } else {
- throw new JobException("Failed to init source reader,
error: " + responseObj.getData());
}
+ String errorMessage =
StringUtils.defaultIfBlank(responseObj.getData(), responseObj.getMsg());
+ throw new JobException("Failed to init source reader, error: "
+ errorMessage);
} catch (JobException jobex) {
log.warn("Failed to init {} source reader, {}", getJobId(),
response);
- throw new JobException(jobex.getMessage());
+ throw jobex;
} catch (Exception e) {
log.warn("Failed to init {} source reader, {}", getJobId(),
response);
throw new JobException("Failed to init source reader, cause "
+ e.getMessage());
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/util/StreamingJobUtils.java
b/fe/fe-core/src/main/java/org/apache/doris/job/util/StreamingJobUtils.java
index bb5d88f48cf..a379a997589 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/job/util/StreamingJobUtils.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/job/util/StreamingJobUtils.java
@@ -382,9 +382,10 @@ public class StreamingJobUtils {
JdbcClient jdbcClient = getJdbcClient(sourceType, properties);
try {
String database = getRemoteDbName(sourceType, properties);
+ String sourceLocation = sourceType == DataSourceType.POSTGRES ?
"schema" : "database";
List<String> tablesNameList =
jdbcClient.getTablesNameList(database);
if (tablesNameList.isEmpty()) {
- throw new JobException("No tables found in database " +
database);
+ throw new JobException("No source tables found in " +
sourceLocation + " '" + database + "'");
}
Database targetDatabase =
Env.getCurrentEnv().getInternalCatalog().getDbNullable(targetDb);
Preconditions.checkNotNull(targetDatabase, "target database %s
does not exist", targetDb);
@@ -483,8 +484,9 @@ public class StreamingJobUtils {
}
if (!noPrimaryKeyTables.isEmpty()) {
- throw new JobException("The following tables do not have
primary key defined: "
- + String.join(", ", noPrimaryKeyTables));
+ throw new JobException("Source tables require primary keys: "
+ + String.join(", ", noPrimaryKeyTables)
+ + ". Add a primary key or exclude these tables.");
}
return createtblCmds;
} finally {
@@ -499,7 +501,8 @@ public class StreamingJobUtils {
List<Column> columns = jdbcClient.getColumnsFromJdbc(database, table);
columns.forEach(col -> {
Preconditions.checkArgument(!col.getType().isUnsupported(),
- "Unsupported column type, table:[%s], column:[%s]", table,
col.getName());
+ "Unsupported column type for source column '%s.%s.%s'",
+ database, table, col.getName());
if (col.getType().isVarchar()) {
// The length of varchar needs to be multiplied by 3.
int len = col.getType().getLength() * 3;
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterJobCommand.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterJobCommand.java
index d75be51c457..0a529d55424 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterJobCommand.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterJobCommand.java
@@ -338,7 +338,7 @@ public class AlterJobCommand extends AlterCommand
implements ForwardWithSync, Ne
private Pair<List<String>, UnboundTVFRelation> getTargetTableAndTvf(String
sql) throws AnalysisException {
LogicalPlan logicalPlan = new NereidsParser().parseSingle(sql);
if (!(logicalPlan instanceof InsertIntoTableCommand)) {
- throw new AnalysisException("Only support insert command");
+ throw new AnalysisException("Streaming jobs only support INSERT
statements");
}
LogicalPlan logicalQuery = ((InsertIntoTableCommand)
logicalPlan).getLogicalQuery();
List<String> targetTable =
InsertUtils.getTargetTableQualified(logicalQuery, ConnectContext.get());
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/CreateJobInfo.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/CreateJobInfo.java
index 49749918d4c..f70f046fec2 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/CreateJobInfo.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/CreateJobInfo.java
@@ -346,8 +346,7 @@ public class CreateJobInfo {
throw new AnalysisException(e.getMessage());
}
} else {
- throw new AnalysisException("Only " +
logicalPlan.getClass().getName()
- + " is supported to use with streaming job together");
+ throw new AnalysisException("Streaming jobs only support INSERT
statements");
}
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/tablefunction/CdcStreamTableValuedFunction.java
b/fe/fe-core/src/main/java/org/apache/doris/tablefunction/CdcStreamTableValuedFunction.java
index f6df6595cd8..ffcc83c85c9 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/tablefunction/CdcStreamTableValuedFunction.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/tablefunction/CdcStreamTableValuedFunction.java
@@ -154,11 +154,13 @@ public class CdcStreamTableValuedFunction extends
ExternalFileTableValuedFunctio
}
String offset = properties.get(DataSourceConfigKeys.OFFSET);
if (!DataSourceConfigValidator.isValidOffset(offset,
sourceType.name())) {
- throw new AnalysisException("Invalid value for key 'offset': " +
offset);
+ throw new
AnalysisException(DataSourceConfigValidator.invalidValueMessage(
+ DataSourceConfigKeys.OFFSET, offset));
}
String sslMode = properties.get(DataSourceConfigKeys.SSL_MODE);
if (sslMode != null &&
!DataSourceConfigValidator.isValidSslMode(sslMode)) {
- throw new AnalysisException("Invalid value for key 'ssl_mode': " +
sslMode);
+ throw new
AnalysisException(DataSourceConfigValidator.invalidValueMessage(
+ DataSourceConfigKeys.SSL_MODE, sslMode));
}
try {
DataSourceConfigValidator.validateSslVerifyCaPair(properties);
@@ -184,7 +186,7 @@ public class CdcStreamTableValuedFunction extends
ExternalFileTableValuedFunctio
return;
}
if (!DataSourceConfigValidator.isPositiveInt(value)) {
- throw new AnalysisException("Invalid value for key '" + key + "':
" + value);
+ throw new
AnalysisException(DataSourceConfigValidator.invalidValueMessage(key, value));
}
}
@@ -195,7 +197,7 @@ public class CdcStreamTableValuedFunction extends
ExternalFileTableValuedFunctio
return;
}
if (!DataSourceConfigValidator.isValidPgIdentifier(value)) {
- throw new AnalysisException("Invalid value for key '" + key + "':
" + value);
+ throw new
AnalysisException(DataSourceConfigValidator.invalidValueMessage(key, value));
}
}
@@ -206,7 +208,8 @@ public class CdcStreamTableValuedFunction extends
ExternalFileTableValuedFunctio
return;
}
if (!DataSourceConfigValidator.isValidBoolean(value)) {
- throw new AnalysisException("Invalid value for key '" + key + "':
" + value);
+ throw new AnalysisException("Invalid value for key '" + key + "':
" + value
+ + ". Expected true or false.");
}
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobCheckDataQualityTest.java
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobCheckDataQualityTest.java
index c15e22d9224..6a69c374f68 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobCheckDataQualityTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobCheckDataQualityTest.java
@@ -17,10 +17,12 @@
package org.apache.doris.job.extensions.insert.streaming;
+import org.apache.doris.common.InternalErrorCode;
import org.apache.doris.common.jmockit.Deencapsulation;
import org.apache.doris.job.base.JobExecutionConfiguration;
import org.apache.doris.job.base.TimerDefinition;
import org.apache.doris.job.cdc.request.CommitOffsetRequest;
+import org.apache.doris.job.common.FailureReason;
import org.apache.doris.job.common.JobStatus;
import org.apache.doris.job.exception.JobException;
@@ -93,6 +95,10 @@ public class StreamingInsertJobCheckDataQualityTest {
thrown = e;
}
Assertions.assertNotNull(thrown, "expected pause when combined ratio
exceeds threshold");
+
Assertions.assertTrue(thrown.getMessage().contains("load.max_filter_ratio"));
+ Assertions.assertTrue(thrown.getMessage().contains("RESUME JOB"));
+ Assertions.assertEquals(InternalErrorCode.TOO_MANY_FAILURE_ROWS_ERR,
+ new FailureReason(thrown.getMessage()).getCode());
Assertions.assertEquals(JobStatus.PAUSED, job.getJobStatus());
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobPropertiesTest.java
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobPropertiesTest.java
index 10ae94e1cc0..64f4515816e 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobPropertiesTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobPropertiesTest.java
@@ -104,7 +104,8 @@ public class StreamingJobPropertiesTest {
Assertions.assertEquals(StreamingJobProperties.DEFAULT_MAX_INTERVAL_SECOND,
p.getMaxIntervalSecond());
// but validate() should throw
- Assertions.assertThrows(AnalysisException.class, p::validate);
+ AnalysisException exception =
Assertions.assertThrows(AnalysisException.class, p::validate);
+ Assertions.assertTrue(exception.getMessage().contains("max_interval
must be at least 1 second"));
}
@Test
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderErrorHandlingTest.java
b/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderErrorHandlingTest.java
index fe7b2304105..1099aa23a6e 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderErrorHandlingTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderErrorHandlingTest.java
@@ -56,7 +56,19 @@ public class JdbcSourceOffsetProviderErrorHandlingTest {
provider.parseCdcResponseData(response, new
TypeReference<List<SnapshotSplit>>() {});
Assertions.fail("a failed envelope must throw");
} catch (JobException e) {
- Assertions.assertTrue(e.getMessage().contains(realError), "the
real remote error must be surfaced, got: " + e.getMessage());
+ Assertions.assertEquals(realError, e.getMessage());
+ }
+ }
+
+ @Test
+ public void testParseFailureEnvelopeFallsBackToMessageWithoutTextData() {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ String response = "{\"code\":1,\"msg\":\"cdc_client is
unavailable\",\"data\":null}";
+ try {
+ provider.parseCdcResponseData(response, new
TypeReference<List<SnapshotSplit>>() {});
+ Assertions.fail("a failed envelope must throw");
+ } catch (JobException e) {
+ Assertions.assertEquals("cdc_client is unavailable",
e.getMessage());
}
}
@@ -69,7 +81,8 @@ public class JdbcSourceOffsetProviderErrorHandlingTest {
provider.parseCdcResponseData(response, new
TypeReference<Map<String, String>>() {});
Assertions.fail("an incompatible success payload must throw");
} catch (JobException e) {
- Assertions.assertTrue(e.getMessage().contains("not-a-map"), "the
raw response must be surfaced, got: " + e.getMessage());
+ Assertions.assertTrue(e.getMessage().contains("not-a-map"),
+ "the raw response must be surfaced, got: " +
e.getMessage());
}
}
@@ -81,7 +94,7 @@ public class JdbcSourceOffsetProviderErrorHandlingTest {
provider.parseCdcResponseData(response, new
TypeReference<Integer>() {});
Assertions.fail("an unparseable response must throw");
} catch (JobException e) {
- Assertions.assertTrue(e.getMessage().contains("502"), "the raw
response must be surfaced, got: " + e.getMessage());
+ Assertions.assertEquals(response, e.getMessage());
}
}
diff --git
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/sink/DorisBatchStreamLoad.java
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/sink/DorisBatchStreamLoad.java
index dbc68c680a9..370e3c30ee9 100644
---
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/sink/DorisBatchStreamLoad.java
+++
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/sink/DorisBatchStreamLoad.java
@@ -572,6 +572,17 @@ public class DorisBatchStreamLoad implements Serializable {
taskId);
return;
}
+ JsonNode data = root.get("data");
+ JsonNode msg = root.get("msg");
+ if (data != null
+ && data.isTextual()
+ && StringUtils.isNotBlank(data.asText())) {
+ reason = data.asText();
+ } else if (msg != null
+ && msg.isTextual()
+ && StringUtils.isNotBlank(msg.asText())) {
+ reason = msg.asText();
+ }
}
LOG.error(
"commit offset failed with {}, reason {}, to
retry",
diff --git
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/sink/DorisBatchStreamLoadTest.java
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/sink/DorisBatchStreamLoadTest.java
new file mode 100644
index 00000000000..7c271f046ad
--- /dev/null
+++
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/sink/DorisBatchStreamLoadTest.java
@@ -0,0 +1,84 @@
+// 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.cdcclient.sink;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+import org.apache.doris.cdcclient.exception.StreamLoadException;
+
+import org.apache.commons.lang3.exception.ExceptionUtils;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
+
+import java.net.InetSocketAddress;
+import java.nio.charset.StandardCharsets;
+import java.util.Collections;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import com.sun.net.httpserver.HttpServer;
+
+class DorisBatchStreamLoadTest {
+
+ @ParameterizedTest
+ @CsvSource(
+ delimiter = '|',
+ value = {
+ "{\"code\":1,\"msg\":\"Error\",\"data\":\"Offset commit
rejected\"}|Offset commit rejected",
+ "{\"code\":1,\"msg\":\"Offset commit
rejected\",\"data\":null}|Offset commit rejected",
+ "{\"code\":1,\"msg\":\"\",\"data\":null}|HTTP/1.1 200 OK"
+ })
+ void commitOffsetPreservesFailureReasonAfterRetries(String response,
String expectedReason)
+ throws Exception {
+ HttpServer server = HttpServer.create(new
InetSocketAddress("127.0.0.1", 0), 0);
+ AtomicInteger requests = new AtomicInteger();
+ server.createContext(
+ "/api/streaming/commit_offset",
+ exchange -> {
+ try {
+ exchange.getRequestBody().readAllBytes();
+ requests.incrementAndGet();
+ byte[] body =
response.getBytes(StandardCharsets.UTF_8);
+ exchange.sendResponseHeaders(200, body.length);
+ exchange.getResponseBody().write(body);
+ } finally {
+ exchange.close();
+ }
+ });
+ server.start();
+ try {
+ DorisBatchStreamLoad streamLoad = new DorisBatchStreamLoad("1",
"test_db");
+ try {
+ streamLoad.setFrontendAddress("127.0.0.1:" +
server.getAddress().getPort());
+ streamLoad.setToken("test-token");
+ StreamLoadException error = assertThrows(
+ StreamLoadException.class,
+ () -> streamLoad.commitOffset(
+ "2", Collections.emptyList(), 0, new
LoadStatistic(), null));
+ assertEquals(
+ "StreamLoadException: commit offset failed with: " +
expectedReason,
+ ExceptionUtils.getRootCauseMessage(error));
+ assertEquals(4, requests.get());
+ } finally {
+ streamLoad.close();
+ }
+ } finally {
+ server.stop(0);
+ }
+ }
+}
diff --git
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_create_alter.groovy
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_create_alter.groovy
index 2e6443db633..30e3af99afe 100644
---
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_create_alter.groovy
+++
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_create_alter.groovy
@@ -181,7 +181,7 @@ suite("test_streaming_mysql_job_create_alter",
"p0,external,mysql,external_docke
"table.create.properties.replication_num" = "1"
)
"""
- exception "No tables found in database"
+ exception "No source tables found in database"
}
// no match table
diff --git
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_dup.groovy
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_dup.groovy
index 2ddacebde16..ac9c8631280 100644
---
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_dup.groovy
+++
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_dup.groovy
@@ -61,7 +61,7 @@ suite("test_streaming_mysql_job_dup",
"p0,external,mysql,external_docker,externa
"table.create.properties.replication_num" = "1"
)
"""
- exception "The following tables do not have primary key defined:
${table1}"
+ exception "Source tables require primary keys: ${table1}"
}
def jobInfo = sql """
diff --git
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_postgres_job_dup.groovy
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_postgres_job_dup.groovy
index 1bec6cd3a25..cbaab4f153a 100644
---
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_postgres_job_dup.groovy
+++
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_postgres_job_dup.groovy
@@ -64,7 +64,7 @@ suite("test_streaming_postgres_job_dup",
"p0,external,pg,external_docker,externa
"table.create.properties.replication_num" = "1"
)
"""
- exception "The following tables do not have primary key defined:
${table1}"
+ exception "Source tables require primary keys: ${table1}"
}
def jobInfo = sql """
diff --git
a/regression-test/suites/job_p0/streaming_job/test_streaming_insert_job_crud.groovy
b/regression-test/suites/job_p0/streaming_job/test_streaming_insert_job_crud.groovy
index f648767fdaf..8de7f1b2748 100644
---
a/regression-test/suites/job_p0/streaming_job/test_streaming_insert_job_crud.groovy
+++
b/regression-test/suites/job_p0/streaming_job/test_streaming_insert_job_crud.groovy
@@ -112,7 +112,7 @@ suite("test_streaming_insert_job_crud") {
"s3.secret_key" = "${getS3SK()}"
);
"""
- }, "s3.max_batch_files should >=1")
+ }, "s3.max_batch_files must be at least 1")
jobCount = sql """ select count(1) from jobs("type"="insert") where Name =
'${jobNameError}' and ExecuteType='STREAMING' """
assert jobCount.get(0).get(0) == 0
@@ -136,7 +136,7 @@ suite("test_streaming_insert_job_crud") {
"s3.secret_key" = "${getS3SK()}"
);
"""
- }, "s3.max_batch_bytes should between 100MB and 10GB")
+ }, "s3.max_batch_bytes must be between 100 MB and 10 GB")
jobCount = sql """ select count(1) from jobs("type"="insert") where Name =
'${jobNameError}' and ExecuteType='STREAMING' """
assert jobCount.get(0).get(0) == 0
@@ -160,7 +160,7 @@ suite("test_streaming_insert_job_crud") {
"s3.secret_key" = "${getS3SK()}"
);
"""
- }, "max_interval should > 1")
+ }, "max_interval must be at least 1 second")
jobCount = sql """ select count(1) from jobs("type"="insert") where Name =
'${jobNameError}' and ExecuteType='STREAMING' """
assert jobCount.get(0).get(0) == 0
@@ -409,7 +409,7 @@ suite("test_streaming_insert_job_crud") {
"s3.secret_key" = "${getS3SK()}"
);
"""
- exception "s3.max_batch_files should >=1"
+ exception "s3.max_batch_files must be at least 1"
}
// alter session var
diff --git
a/regression-test/suites/job_p0/streaming_job/test_streaming_job_max_retry.groovy
b/regression-test/suites/job_p0/streaming_job/test_streaming_job_max_retry.groovy
index 30fc6eba1c1..b137198e174 100644
---
a/regression-test/suites/job_p0/streaming_job/test_streaming_job_max_retry.groovy
+++
b/regression-test/suites/job_p0/streaming_job/test_streaming_job_max_retry.groovy
@@ -101,6 +101,8 @@ suite("test_streaming_job_max_retry", "nonConcurrent") {
"ErrorMsg should contain CANNOT_RESUME_ERR code, got: " +
errorMsgJson
assert errorMsgJson.contains("Auto resume failed after"),
"ErrorMsg should contain the burn-out message, got: " +
errorMsgJson
+ assert errorMsgJson.contains("RESUME JOB"),
+ "ErrorMsg should explain the manual recovery action, got: " +
errorMsgJson
} finally {
GetDebugPoint().disableDebugPointForAllFEs('StreamingJob.scheduleTask.exception')
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]