EmmyMiao87 commented on a change in pull request #754: Add persist operations
for routine load job
URL: https://github.com/apache/incubator-doris/pull/754#discussion_r265535009
##########
File path:
fe/src/main/java/org/apache/doris/analysis/CreateRoutineLoadStmt.java
##########
@@ -221,65 +219,62 @@ private void checkLoadProperties(Analyzer analyzer)
throws AnalysisException {
throw new AnalysisException("repeat setting of column
separator");
}
columnSeparator = (ColumnSeparator) parseNode;
- columnSeparator.analyze(analyzer);
- } else if (parseNode instanceof LoadColumnsInfo) {
+ columnSeparator.analyze(null);
+ } else if (parseNode instanceof ImportColumnsStmt) {
// check columns info
- if (columnsInfo != null) {
+ if (importColumnsStmt != null) {
throw new AnalysisException("repeat setting of columns
info");
}
- columnsInfo = (LoadColumnsInfo) parseNode;
- columnsInfo.analyze(analyzer);
- } else if (parseNode instanceof Expr) {
+ importColumnsStmt = (ImportColumnsStmt) parseNode;
+ } else if (parseNode instanceof ImportWhereStmt) {
// check where expr
- if (wherePredicate != null) {
+ if (importWhereStmt != null) {
throw new AnalysisException("repeat setting of where
predicate");
}
- wherePredicate = (Expr) parseNode;
- wherePredicate.analyze(analyzer);
+ importWhereStmt = (ImportWhereStmt) parseNode;
} else if (parseNode instanceof PartitionNames) {
// check partition names
if (partitionNames != null) {
throw new AnalysisException("repeat setting of partition
names");
}
partitionNames = (PartitionNames) parseNode;
- partitionNames.analyze(analyzer);
+ partitionNames.analyze(null);
}
}
- routineLoadDesc = new RoutineLoadDesc(columnSeparator, columnsInfo,
wherePredicate,
+ routineLoadDesc = new RoutineLoadDesc(columnSeparator,
importColumnsStmt, importWhereStmt,
partitionNames.getPartitionNames());
}
- private void checkRoutineLoadProperties() throws AnalysisException {
- if (properties != null) {
- Optional<String> optional = properties.keySet().parallelStream()
- .filter(entity ->
!PROPERTIES_SET.contains(entity)).findFirst();
- if (optional.isPresent()) {
- throw new AnalysisException(optional.get() + " is invalid
property");
- }
-
- // check desired concurrent number
- final String desiredConcurrentNumberString =
properties.get(DESIRED_CONCURRENT_NUMBER_PROPERTY);
- if (desiredConcurrentNumberString != null) {
- desiredConcurrentNum =
getIntegerValueFromString(desiredConcurrentNumberString,
-
DESIRED_CONCURRENT_NUMBER_PROPERTY);
- if (desiredConcurrentNum <= 0) {
- throw new
AnalysisException(DESIRED_CONCURRENT_NUMBER_PROPERTY + " must be greater then
0");
- }
- }
+ private void checkJobProperties() throws AnalysisException {
+ Optional<String> optional =
jobProperties.keySet().parallelStream().filter(
+ entity -> !PROPERTIES_SET.contains(entity)).findFirst();
+ if (optional.isPresent()) {
+ throw new AnalysisException(optional.get() + " is invalid
property");
+ }
- // check max error number
- final String maxErrorNumberString =
properties.get(MAX_ERROR_NUMBER_PROPERTY);
- if (maxErrorNumberString != null) {
- maxErrorNum = getIntegerValueFromString(maxErrorNumberString,
MAX_ERROR_NUMBER_PROPERTY);
- if (maxErrorNum < 0) {
- throw new AnalysisException(MAX_ERROR_NUMBER_PROPERTY + "
must be greater then or equal to 0");
- }
+ desiredConcurrentNum =
getIntegetPropertyOrDefault(DESIRED_CONCURRENT_NUMBER_PROPERTY,
+ "must be greater then 0", desiredConcurrentNum);
+ maxErrorNum = getIntegetPropertyOrDefault(MAX_ERROR_NUMBER_PROPERTY,
+ "must be greater then or equal to 0", maxErrorNum);
+ maxBatchIntervalS =
getIntegetPropertyOrDefault(MAX_BATCH_INTERVAL_SECOND,
+ "must be greater then 0", maxBatchIntervalS);
+ maxBatchRows = getIntegetPropertyOrDefault(MAX_BATCH_ROWS, "must be
greater then 0", maxBatchRows);
+ maxBatchSizeBytes = getIntegetPropertyOrDefault(MAX_BATCH_SIZE, "must
be greater then 0", maxBatchSizeBytes);
+ }
+ private int getIntegetPropertyOrDefault(String propName, String hintMsg,
int defaultVal) throws AnalysisException {
+ final String propVal = jobProperties.get(propName);
+ if (propVal != null) {
+ int intVal = getIntegerValueFromString(propVal, propName);
+ if (intVal <= 0) {
+ throw new AnalysisException(propName + " " + hintMsg);
}
+ return intVal;
}
+ return defaultVal;
}
- private void checkCustomProperties() throws AnalysisException {
+ private void checkLoadSourceProperties() throws AnalysisException {
Review comment:
`LoadSource` or `DataSource` ?
----------------------------------------------------------------
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.
For queries about this service, please contact Infrastructure at:
[email protected]
With regards,
Apache Git Services
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]