From d22630bb4eecdc83914e57e490ee274c4b76fd2a Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Wed, 2 Oct 2019 22:47:16 -0700 Subject: [PATCH 01/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../samza/sql/planner/SamzaSqlValidator.java | 144 ++++++++++++------ 1 file changed, 97 insertions(+), 47 deletions(-) diff --git a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java index 97b5de93c4..81d9715e09 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java @@ -33,11 +33,14 @@ import org.apache.calcite.rel.type.RelDataTypeField; import org.apache.calcite.rel.type.RelRecordType; import org.apache.calcite.sql.type.SqlTypeName; +import org.apache.commons.lang.StringUtils; +import org.apache.commons.lang.Validate; import org.apache.samza.SamzaException; import org.apache.samza.config.Config; import org.apache.samza.sql.data.SamzaSqlRelMessage; import org.apache.samza.sql.dsl.SamzaSqlDslConverter; import org.apache.samza.sql.interfaces.RelSchemaProvider; +import org.apache.samza.sql.interfaces.SamzaRelConverter; import org.apache.samza.sql.interfaces.SamzaSqlJavaTypeFactoryImpl; import org.apache.samza.sql.runner.SamzaSqlApplicationConfig; import org.apache.samza.sql.schema.SqlFieldSchema; @@ -78,18 +81,41 @@ public void validate(List sqlStmts) throws SamzaSqlValidatorException { try { relRoot = planner.plan(qinfo.getSelectQuery()); } catch (SamzaException e) { - throw new SamzaSqlValidatorException("Calcite planning for sql failed.", e); + throw new SamzaSqlValidatorException(String.format("Validation failed for sql stmt:\n%s\n", sql), e); } // Now that we have logical plan, validate different aspects. - validate(relRoot, qinfo, sqlConfig); + String sink = qinfo.getSink(); + validate(relRoot, sqlConfig.getRelSchemaProviders().get(sink), sqlConfig.getSamzaRelConverters().get(sink)); } } - protected void validate(RelRoot relRoot, SamzaSqlQueryParser.QueryInfo qinfo, SamzaSqlApplicationConfig sqlConfig) - throws SamzaSqlValidatorException { - // Validate select fields (including Udf return types) with output schema - validateOutput(relRoot, sqlConfig.getRelSchemaProviders().get(qinfo.getSink())); + /** + * Determine if validation needs to be done on Calcite plan based on the schema provider and schema converter. + * @param relRoot + * @param outputSchemaProvider + * @param ouputRelSchemaConverter + * @return if the validation needs to be skipped + */ + protected boolean skipOutputValidation(RelRoot relRoot, RelSchemaProvider outputSchemaProvider, + SamzaRelConverter ouputRelSchemaConverter) { + return false; + } + + // TODO: Remove this API. This API is introduced to take care of cases where RelSchemaProviders have a complex + // mechanism to determine if a given output field is optional. We will need system specific validators to take + // care of such cases and once that is introduced, we can get rid of the below API. + protected boolean isOptional(RelSchemaProvider outputRelSchemaProvider, String outputFieldName, + RelRecordType projectRecord) { + return false; + } + + private void validate(RelRoot relRoot, RelSchemaProvider outputSchemaProvider, + SamzaRelConverter outputRelSchemaConverter) throws SamzaSqlValidatorException { + if (!skipOutputValidation(relRoot, outputSchemaProvider, outputRelSchemaConverter)) { + // Validate select fields (including Udf return types) with output schema + validateOutput(relRoot, outputSchemaProvider); + } // TODO: // 1. SAMZA-2314: Validate Udf arguments. @@ -97,22 +123,25 @@ protected void validate(RelRoot relRoot, SamzaSqlQueryParser.QueryInfo qinfo, Sa // Eg: LogicalAggregate with sum function is not supported by Samza Sql. } - protected void validateOutput(RelRoot relRoot, RelSchemaProvider relSchemaProvider) throws SamzaSqlValidatorException { - RelRecordType outputRecord = (RelRecordType) QueryPlanner.getSourceRelSchema(relSchemaProvider, + private void validateOutput(RelRoot relRoot, RelSchemaProvider outputRelSchemaProvider) + throws SamzaSqlValidatorException { + LogicalProject project = (LogicalProject) relRoot.rel; + + RelRecordType projetRecord = (RelRecordType) project.getRowType(); + RelRecordType outputRecord = (RelRecordType) QueryPlanner.getSourceRelSchema(outputRelSchemaProvider, new RelSchemaConverter()); + // Get Samza Sql schema along with Calcite schema. The reason is that the Calcite schema does not have a way - // to represent optional fields while Samza Sql schema can represent optional fields. This is the only reason that + // to represent optional fields while Samza Sql schema can represent optional fields. This is the reason that // we use SqlSchema in validating output. - SqlSchema outputSqlSchema = QueryPlanner.getSourceSqlSchema(relSchemaProvider); + SqlSchema outputSqlSchema = QueryPlanner.getSourceSqlSchema(outputRelSchemaProvider); - LogicalProject project = (LogicalProject) relRoot.rel; - RelRecordType projetRecord = (RelRecordType) project.getRowType(); - - validateOutputRecords(outputRecord, outputSqlSchema, projetRecord); + validateOutputRecords(outputSqlSchema, outputRecord, projetRecord, outputRelSchemaProvider); + LOG.info("Samza Sql Validation finished successfully."); } - protected void validateOutputRecords(RelRecordType outputRecord, SqlSchema outputSqlSchema, - RelRecordType projectRecord) + private void validateOutputRecords(SqlSchema outputSqlSchema, RelRecordType outputRecord, + RelRecordType projectRecord, RelSchemaProvider outputRelSchemaProvider) throws SamzaSqlValidatorException { Map outputRecordMap = outputRecord.getFieldList().stream().collect( Collectors.toMap(RelDataTypeField::getName, RelDataTypeField::getType)); @@ -121,6 +150,50 @@ protected void validateOutputRecords(RelRecordType outputRecord, SqlSchema outpu Map projectRecordMap = projectRecord.getFieldList().stream().collect( Collectors.toMap(RelDataTypeField::getName, RelDataTypeField::getType)); + // Ensure that all fields from sql statement exist in the output schema and are of the same type. + for (Map.Entry entry : projectRecordMap.entrySet()) { + String projectedFieldName = entry.getKey(); + RelDataType outputFieldType = outputRecordMap.get(projectedFieldName); + SqlFieldSchema outputSqlFieldSchema = outputFieldSchemaMap.get(projectedFieldName); + + if (outputFieldType == null) { + // If the field names are specified more than once in the select query, calcite appends 'n' as suffix to the + // dup fields based on the order they are specified, where 'n' starts from 0 for the first dup field. + // Take the following example: SELECT id as str, secondaryId as str, tertiaryId as str FROM store.myTable + // Calcite renames the projected fieldNames in select query as str, str0, str1 respectively. + // Samza Sql allows a field name to be specified up to 2 times. Do the validation accordingly. + + // This type of pattern is typically followed when users want to just modify one field in the input table while + // keeping rest of the fields the same. Eg: SELECT myUdf(id) as id, * from store.myTable + if (projectedFieldName.endsWith("0")) { + projectedFieldName = StringUtils.chop(projectedFieldName); + outputFieldType = outputRecordMap.get(projectedFieldName); + outputSqlFieldSchema = outputFieldSchemaMap.get(projectedFieldName); + } + + if (outputFieldType == null) { + // If a field in sql query is not found in the output schema, ignore if it is a Samza Sql special op. + // Otherwise, throw an error. + if (entry.getKey().equals(SamzaSqlRelMessage.OP_NAME)) { + continue; + } + String errMsg = String.format("Field '%s' in select query does not match any field in output schema.", entry.getKey()); + LOG.error(errMsg); + throw new SamzaSqlValidatorException(errMsg); + } + } + + Validate.notNull(outputFieldType); + Validate.notNull(outputSqlFieldSchema); + + if (!compareFieldTypes(outputFieldType, outputSqlFieldSchema, entry.getValue(), outputRelSchemaProvider)) { + String errMsg = String.format("Field '%s' with type '%s' in select query does not match the field type '%s' in" + + " output schema.", entry.getKey(), entry.getValue(), outputFieldType); + LOG.error(errMsg); + throw new SamzaSqlValidatorException(errMsg); + } + } + // Ensure that all non-optional fields in output schema are set in the sql query and are of the // same type. for (Map.Entry entry : outputRecordMap.entrySet()) { @@ -130,47 +203,24 @@ protected void validateOutputRecords(RelRecordType outputRecord, SqlSchema outpu if (projectFieldType == null) { // If an output schema field is not found in the sql query, ignore it if the field is optional. // Otherwise, throw an error. - if (outputSqlFieldSchema.isOptional()) { + if (outputSqlFieldSchema.isOptional() || isOptional(outputRelSchemaProvider, entry.getKey(), projectRecord)) { continue; } - String errMsg = String.format("Field '%s' in output schema does not match any projected fields.", - entry.getKey()); + String errMsg = String.format("Non-optional field '%s' in output schema is missing in projected fields of " + + "select query.", entry.getKey()); LOG.error(errMsg); throw new SamzaSqlValidatorException(errMsg); - } else if (!compareFieldTypes(entry.getValue(), outputSqlFieldSchema, projectFieldType)) { + } else if (!compareFieldTypes(entry.getValue(), outputSqlFieldSchema, projectFieldType, outputRelSchemaProvider)) { String errMsg = String.format("Field '%s' with type '%s' in output schema does not match the field type '%s' in" + " projected fields.", entry.getKey(), entry.getValue(), projectFieldType); LOG.error(errMsg); throw new SamzaSqlValidatorException(errMsg); } } - - // Ensure that all fields from sql statement exist in the output schema and are of the same type. - for (Map.Entry entry : projectRecordMap.entrySet()) { - RelDataType outputFieldType = outputRecordMap.get(entry.getKey()); - SqlFieldSchema outputSqlFieldSchema = outputFieldSchemaMap.get(entry.getKey()); - - if (outputFieldType == null) { - // If a field in sql query is not found in the output schema, ignore if it is a Samza Sql special op. - // Otherwise, throw an error. - if (entry.getKey().equals(SamzaSqlRelMessage.OP_NAME)) { - continue; - } - String errMsg = String.format("Field '%s' in select query does not match any field in output schema.", - entry.getKey()); - LOG.error(errMsg); - throw new SamzaSqlValidatorException(errMsg); - } else if (!compareFieldTypes(outputFieldType, outputSqlFieldSchema, entry.getValue())) { - String errMsg = String.format("Field '%s' with type '%s' in select query does not match the field type '%s' in" - + " output schema.", entry.getKey(), entry.getValue(), outputFieldType); - LOG.error(errMsg); - throw new SamzaSqlValidatorException(errMsg); - } - } } - protected boolean compareFieldTypes(RelDataType outputFieldType, SqlFieldSchema sqlFieldSchema, - RelDataType selectQueryFieldType) { + private boolean compareFieldTypes(RelDataType outputFieldType, SqlFieldSchema sqlFieldSchema, + RelDataType selectQueryFieldType, RelSchemaProvider outputRelSchemaProvider) { RelDataType projectFieldType; // JavaTypes are relevant for Udf argument and return types @@ -206,8 +256,8 @@ protected boolean compareFieldTypes(RelDataType outputFieldType, SqlFieldSchema return projectSqlType == SqlTypeName.FLOAT; case ROW: try { - validateOutputRecords((RelRecordType) outputFieldType, sqlFieldSchema.getRowSchema(), - (RelRecordType) projectFieldType); + validateOutputRecords(sqlFieldSchema.getRowSchema(), (RelRecordType) outputFieldType, + (RelRecordType) projectFieldType, outputRelSchemaProvider); } catch (SamzaSqlValidatorException e) { LOG.error("A field in select query does not match with the output schema.", e); return false; From 2086ad6e19054dfcff14ba187e8b1c201005f1e1 Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Fri, 13 Sep 2019 00:58:40 -0700 Subject: [PATCH 02/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../samza/sql/avro/AvroTypeFactoryImpl.java | 13 +- .../samza/sql/planner/SamzaSqlValidator.java | 120 ++++++++++++------ .../sql/translator/FilterTranslator.java | 21 ++- .../sql/translator/ProjectTranslator.java | 10 +- .../sql/planner/TestSamzaSqlValidator.java | 28 +++- .../org/apache/samza/sql/util/MyTestUdf.java | 5 + .../test/samzasql/TestSamzaSqlEndToEnd.java | 4 +- 7 files changed, 142 insertions(+), 59 deletions(-) diff --git a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java index 6bf0c3c4a1..9de2d978f5 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java @@ -52,15 +52,18 @@ public SqlSchema createType(Schema schema) { throw new SamzaException(msg); } - return convertSchema(schema.getFields()); + return convertSchema(schema.getFields(), true); } - private SqlSchema convertSchema(List fields) { + protected boolean isOptional(Schema.Field field, boolean isTopLevel) { + return field.defaultValue() != null; + } + + private SqlSchema convertSchema(List fields, boolean isTopLevel) { SqlSchemaBuilder schemaBuilder = SqlSchemaBuilder.builder(); for (Schema.Field field : fields) { - boolean isOptional = (field.defaultValue() != null); - SqlFieldSchema fieldSchema = convertField(field.schema(), false, isOptional); + SqlFieldSchema fieldSchema = convertField(field.schema(), false, isOptional(field, isTopLevel)); schemaBuilder.addField(field.name(), fieldSchema); } @@ -98,7 +101,7 @@ private SqlFieldSchema convertField(Schema fieldSchema, boolean isNullable, bool case LONG: return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.INT64, isNullable, isOptional); case RECORD: - SqlSchema rowSchema = convertSchema(fieldSchema.getFields()); + SqlSchema rowSchema = convertSchema(fieldSchema.getFields(), false); return SqlFieldSchema.createRowFieldSchema(rowSchema, isNullable, isOptional); case MAP: // Can the value type be nullable and have default values ? Guess not! diff --git a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java index 97b5de93c4..f81600d933 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java @@ -33,6 +33,8 @@ import org.apache.calcite.rel.type.RelDataTypeField; import org.apache.calcite.rel.type.RelRecordType; import org.apache.calcite.sql.type.SqlTypeName; +import org.apache.commons.lang.StringUtils; +import org.apache.commons.lang.Validate; import org.apache.samza.SamzaException; import org.apache.samza.config.Config; import org.apache.samza.sql.data.SamzaSqlRelMessage; @@ -86,10 +88,12 @@ public void validate(List sqlStmts) throws SamzaSqlValidatorException { } } - protected void validate(RelRoot relRoot, SamzaSqlQueryParser.QueryInfo qinfo, SamzaSqlApplicationConfig sqlConfig) + private void validate(RelRoot relRoot, SamzaSqlQueryParser.QueryInfo qinfo, SamzaSqlApplicationConfig sqlConfig) throws SamzaSqlValidatorException { - // Validate select fields (including Udf return types) with output schema - validateOutput(relRoot, sqlConfig.getRelSchemaProviders().get(qinfo.getSink())); + if (!skipOutputValidation(relRoot, qinfo, sqlConfig)) { + // Validate select fields (including Udf return types) with output schema + validateOutput(relRoot, sqlConfig.getRelSchemaProviders().get(qinfo.getSink())); + } // TODO: // 1. SAMZA-2314: Validate Udf arguments. @@ -97,22 +101,33 @@ protected void validate(RelRoot relRoot, SamzaSqlQueryParser.QueryInfo qinfo, Sa // Eg: LogicalAggregate with sum function is not supported by Samza Sql. } - protected void validateOutput(RelRoot relRoot, RelSchemaProvider relSchemaProvider) throws SamzaSqlValidatorException { - RelRecordType outputRecord = (RelRecordType) QueryPlanner.getSourceRelSchema(relSchemaProvider, + private void validateOutput(RelRoot relRoot, RelSchemaProvider outputRelSchemaProvider) + throws SamzaSqlValidatorException { + LogicalProject project = (LogicalProject) relRoot.rel; + RelRecordType projetRecord = (RelRecordType) project.getRowType(); + RelRecordType outputRecord = (RelRecordType) QueryPlanner.getSourceRelSchema(outputRelSchemaProvider, new RelSchemaConverter()); // Get Samza Sql schema along with Calcite schema. The reason is that the Calcite schema does not have a way // to represent optional fields while Samza Sql schema can represent optional fields. This is the only reason that // we use SqlSchema in validating output. - SqlSchema outputSqlSchema = QueryPlanner.getSourceSqlSchema(relSchemaProvider); + SqlSchema outputSqlSchema = QueryPlanner.getSourceSqlSchema(outputRelSchemaProvider); - LogicalProject project = (LogicalProject) relRoot.rel; - RelRecordType projetRecord = (RelRecordType) project.getRowType(); + validateOutputRecords(outputRecord, outputSqlSchema, projetRecord, outputRelSchemaProvider); + LOG.info("Samza Sql Validation finished successfully."); + } + + protected boolean skipOutputValidation(RelRoot relRoot, SamzaSqlQueryParser.QueryInfo qinfo, + SamzaSqlApplicationConfig sqlConfig) { + return false; + } - validateOutputRecords(outputRecord, outputSqlSchema, projetRecord); + protected boolean isOptional(RelSchemaProvider outputRelSchemaProvider, String outputFieldName, + RelRecordType projectRecord) { + return false; } - protected void validateOutputRecords(RelRecordType outputRecord, SqlSchema outputSqlSchema, - RelRecordType projectRecord) + private void validateOutputRecords(RelRecordType outputRecord, SqlSchema outputSqlSchema, + RelRecordType projectRecord, RelSchemaProvider outputRelSchemaProvider) throws SamzaSqlValidatorException { Map outputRecordMap = outputRecord.getFieldList().stream().collect( Collectors.toMap(RelDataTypeField::getName, RelDataTypeField::getType)); @@ -121,6 +136,47 @@ protected void validateOutputRecords(RelRecordType outputRecord, SqlSchema outpu Map projectRecordMap = projectRecord.getFieldList().stream().collect( Collectors.toMap(RelDataTypeField::getName, RelDataTypeField::getType)); + // Ensure that all fields from sql statement exist in the output schema and are of the same type. + for (Map.Entry entry : projectRecordMap.entrySet()) { + String projectedFieldName = entry.getKey(); + RelDataType outputFieldType = outputRecordMap.get(projectedFieldName); + SqlFieldSchema outputSqlFieldSchema = outputFieldSchemaMap.get(projectedFieldName); + + if (outputFieldType == null) { + // If the field names are specified more than once in the select query, calcite appends 'n' as suffix to the + // dup fields based on the order they are specified, where 'n' starts from 0 for the first dup field. + // Take the following example: SELECT id as str, secondaryId as str, tertiaryId as str FROM store.myTable + // Calcite renames the projected fieldNames in select query as str, str0, str1 respectively. + // Samza Sql allows a field name to be specified up to 2 times. Do the validation accordingly. + if (projectedFieldName.endsWith("0")) { + projectedFieldName = StringUtils.chop(projectedFieldName); + outputFieldType = outputRecordMap.get(projectedFieldName); + outputSqlFieldSchema = outputFieldSchemaMap.get(projectedFieldName); + } + + if (outputFieldType == null) { + // If a field in sql query is not found in the output schema, ignore if it is a Samza Sql special op. + // Otherwise, throw an error. + if (entry.getKey().equals(SamzaSqlRelMessage.OP_NAME)) { + continue; + } + String errMsg = String.format("Field '%s' in select query does not match any field in output schema.", entry.getKey()); + LOG.error(errMsg); + throw new SamzaSqlValidatorException(errMsg); + } + } + + Validate.notNull(outputFieldType); + Validate.notNull(outputSqlFieldSchema); + + if (!compareFieldTypes(outputFieldType, outputSqlFieldSchema, entry.getValue(), outputRelSchemaProvider)) { + String errMsg = String.format("Field '%s' with type '%s' in select query does not match the field type '%s' in" + + " output schema.", entry.getKey(), entry.getValue(), outputFieldType); + LOG.error(errMsg); + throw new SamzaSqlValidatorException(errMsg); + } + } + // Ensure that all non-optional fields in output schema are set in the sql query and are of the // same type. for (Map.Entry entry : outputRecordMap.entrySet()) { @@ -130,47 +186,24 @@ protected void validateOutputRecords(RelRecordType outputRecord, SqlSchema outpu if (projectFieldType == null) { // If an output schema field is not found in the sql query, ignore it if the field is optional. // Otherwise, throw an error. - if (outputSqlFieldSchema.isOptional()) { + if (outputSqlFieldSchema.isOptional() || isOptional(outputRelSchemaProvider, entry.getKey(), projectRecord)) { continue; } - String errMsg = String.format("Field '%s' in output schema does not match any projected fields.", - entry.getKey()); + String errMsg = String.format("Non-optional field '%s' in output schema is missing in projected fields of " + + "select query.", entry.getKey()); LOG.error(errMsg); throw new SamzaSqlValidatorException(errMsg); - } else if (!compareFieldTypes(entry.getValue(), outputSqlFieldSchema, projectFieldType)) { + } else if (!compareFieldTypes(entry.getValue(), outputSqlFieldSchema, projectFieldType, outputRelSchemaProvider)) { String errMsg = String.format("Field '%s' with type '%s' in output schema does not match the field type '%s' in" + " projected fields.", entry.getKey(), entry.getValue(), projectFieldType); LOG.error(errMsg); throw new SamzaSqlValidatorException(errMsg); } } - - // Ensure that all fields from sql statement exist in the output schema and are of the same type. - for (Map.Entry entry : projectRecordMap.entrySet()) { - RelDataType outputFieldType = outputRecordMap.get(entry.getKey()); - SqlFieldSchema outputSqlFieldSchema = outputFieldSchemaMap.get(entry.getKey()); - - if (outputFieldType == null) { - // If a field in sql query is not found in the output schema, ignore if it is a Samza Sql special op. - // Otherwise, throw an error. - if (entry.getKey().equals(SamzaSqlRelMessage.OP_NAME)) { - continue; - } - String errMsg = String.format("Field '%s' in select query does not match any field in output schema.", - entry.getKey()); - LOG.error(errMsg); - throw new SamzaSqlValidatorException(errMsg); - } else if (!compareFieldTypes(outputFieldType, outputSqlFieldSchema, entry.getValue())) { - String errMsg = String.format("Field '%s' with type '%s' in select query does not match the field type '%s' in" - + " output schema.", entry.getKey(), entry.getValue(), outputFieldType); - LOG.error(errMsg); - throw new SamzaSqlValidatorException(errMsg); - } - } } - protected boolean compareFieldTypes(RelDataType outputFieldType, SqlFieldSchema sqlFieldSchema, - RelDataType selectQueryFieldType) { + private boolean compareFieldTypes(RelDataType outputFieldType, SqlFieldSchema sqlFieldSchema, + RelDataType selectQueryFieldType, RelSchemaProvider outputRelSchemaProvider) { RelDataType projectFieldType; // JavaTypes are relevant for Udf argument and return types @@ -207,7 +240,7 @@ protected boolean compareFieldTypes(RelDataType outputFieldType, SqlFieldSchema case ROW: try { validateOutputRecords((RelRecordType) outputFieldType, sqlFieldSchema.getRowSchema(), - (RelRecordType) projectFieldType); + (RelRecordType) projectFieldType, outputRelSchemaProvider); } catch (SamzaSqlValidatorException e) { LOG.error("A field in select query does not match with the output schema.", e); return false; @@ -250,7 +283,10 @@ public static String formatErrorString(String query, Exception e) { Matcher matcher = pattern.matcher(e.getMessage()); String[] queryLines = query.split("\\n"); StringBuilder result = new StringBuilder(); - int startColIdx, endColIdx, startLineIdx, endLineIdx; + int startColIdx; + int endColIdx; + int startLineIdx; + int endLineIdx; try { if (matcher.find()) { diff --git a/samza-sql/src/main/java/org/apache/samza/sql/translator/FilterTranslator.java b/samza-sql/src/main/java/org/apache/samza/sql/translator/FilterTranslator.java index 345ac8534d..29253a5ef6 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/translator/FilterTranslator.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/translator/FilterTranslator.java @@ -19,9 +19,12 @@ package org.apache.samza.sql.translator; +import java.time.Duration; +import java.time.Instant; import java.util.Arrays; import java.util.Collections; import org.apache.calcite.rel.logical.LogicalFilter; +import org.apache.samza.SamzaException; import org.apache.samza.context.ContainerContext; import org.apache.samza.context.Context; import org.apache.samza.metrics.Counter; @@ -42,7 +45,7 @@ */ class FilterTranslator { - private static final Logger log = LoggerFactory.getLogger(FilterTranslator.class); + private static final Logger LOG = LoggerFactory.getLogger(FilterTranslator.class); private final int queryId; FilterTranslator(int queryId) { @@ -95,17 +98,23 @@ public void init(Context context) { public boolean apply(SamzaSqlRelMessage message) { long startProcessing = System.nanoTime(); Object[] result = new Object[1]; - expr.execute(translatorContext.getExecutionContext(), context, translatorContext.getDataContext(), - message.getSamzaSqlRelRecord().getFieldValues().toArray(), result); - if (result.length > 0 && result[0] instanceof Boolean) { + try { + expr.execute(translatorContext.getExecutionContext(), context, translatorContext.getDataContext(), + message.getSamzaSqlRelRecord().getFieldValues().toArray(), result); + } catch (Exception e) { + String errMsg = String.format("Handling the following rel message ran into an error. %s", message); + LOG.error(errMsg, e); + throw new SamzaException(errMsg, e); + } + if (result[0] instanceof Boolean) { boolean retVal = (Boolean) result[0]; - log.debug( + LOG.debug( String.format("return value for input %s is %s", Arrays.asList(message.getSamzaSqlRelRecord().getFieldValues()).toString(), retVal)); updateMetrics(startProcessing, retVal, System.nanoTime()); return retVal; } else { - log.error("return value is not boolean"); + LOG.error("return value is not boolean for rel message: %s", message); return false; } } diff --git a/samza-sql/src/main/java/org/apache/samza/sql/translator/ProjectTranslator.java b/samza-sql/src/main/java/org/apache/samza/sql/translator/ProjectTranslator.java index 6269d55871..bf44815aa9 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/translator/ProjectTranslator.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/translator/ProjectTranslator.java @@ -114,8 +114,14 @@ public SamzaSqlRelMessage apply(SamzaSqlRelMessage message) { long arrivalTime = System.nanoTime(); RelDataType type = project.getRowType(); Object[] output = new Object[type.getFieldCount()]; - expr.execute(translatorContext.getExecutionContext(), context, translatorContext.getDataContext(), - message.getSamzaSqlRelRecord().getFieldValues().toArray(), output); + try { + expr.execute(translatorContext.getExecutionContext(), context, translatorContext.getDataContext(), + message.getSamzaSqlRelRecord().getFieldValues().toArray(), output); + } catch (Exception e) { + String errMsg = String.format("Handling the following rel message ran into an error. %s", message); + LOG.error(errMsg, e); + throw new SamzaException(errMsg, e); + } List names = new ArrayList<>(); for (int index = 0; index < output.length; index++) { names.add(index, project.getNamedProjects().get(index).getValue()); diff --git a/samza-sql/src/test/java/org/apache/samza/sql/planner/TestSamzaSqlValidator.java b/samza-sql/src/test/java/org/apache/samza/sql/planner/TestSamzaSqlValidator.java index 4c9522eb22..ce783be4f4 100644 --- a/samza-sql/src/test/java/org/apache/samza/sql/planner/TestSamzaSqlValidator.java +++ b/samza-sql/src/test/java/org/apache/samza/sql/planner/TestSamzaSqlValidator.java @@ -59,6 +59,30 @@ public void testBasicValidation() throws SamzaSqlValidatorException { new SamzaSqlValidator(samzaConfig).validate(sqlStmts); } + @Test + public void testRepeatedTwiceFieldsValidation() throws SamzaSqlValidatorException { + Map config = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(1); + config.put(SamzaSqlApplicationConfig.CFG_SQL_STMT, + "Insert into testavro.outputTopic select id, true as bool_value, false as bool_value" + + " from testavro.level1.level2.SIMPLE1 as s where s.id = 1"); + Config samzaConfig = SamzaSqlApplicationRunner.computeSamzaConfigs(true, new MapConfig(config)); + + List sqlStmts = fetchSqlFromConfig(config); + new SamzaSqlValidator(samzaConfig).validate(sqlStmts); + } + + @Test (expected = SamzaSqlValidatorException.class) + public void testRepeatedThriceFieldsValidation() throws SamzaSqlValidatorException { + Map config = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(1); + config.put(SamzaSqlApplicationConfig.CFG_SQL_STMT, + "Insert into testavro.outputTopic select id, true as bool_value, false as bool_value, true as bool_value" + + " from testavro.level1.level2.SIMPLE1 as s where s.id = 1"); + Config samzaConfig = SamzaSqlApplicationRunner.computeSamzaConfigs(true, new MapConfig(config)); + + List sqlStmts = fetchSqlFromConfig(config); + new SamzaSqlValidator(samzaConfig).validate(sqlStmts); + } + @Test (expected = SamzaSqlValidatorException.class) public void testNonExistingOutputField() throws SamzaSqlValidatorException { Map config = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(1); @@ -160,7 +184,7 @@ public void testNonDefaultButNullableField() { try { new SamzaSqlValidator(samzaConfig).validate(sqlStmts); } catch (SamzaSqlValidatorException e) { - Assert.assertTrue(e.getMessage().contains("Field 'bool_value' in output schema does not match any projected fields.")); + Assert.assertTrue(e.getMessage().contains("Non-optional field 'bool_value' in output schema is missing")); return; } @@ -181,7 +205,7 @@ public void testNonDefaultOutputField() { try { new SamzaSqlValidator(samzaConfig).validate(sqlStmts); } catch (SamzaSqlValidatorException e) { - Assert.assertTrue(e.getMessage().contains("Field 'id' in output schema does not match")); + Assert.assertTrue(e.getMessage().contains("Non-optional field 'id' in output schema is missing")); return; } diff --git a/samza-sql/src/test/java/org/apache/samza/sql/util/MyTestUdf.java b/samza-sql/src/test/java/org/apache/samza/sql/util/MyTestUdf.java index 983df51000..9716df59e9 100644 --- a/samza-sql/src/test/java/org/apache/samza/sql/util/MyTestUdf.java +++ b/samza-sql/src/test/java/org/apache/samza/sql/util/MyTestUdf.java @@ -42,6 +42,11 @@ public Integer execute(Integer value) { return value * 2; } + @SamzaSqlUdfMethod(params = SamzaSqlFieldType.STRING) + public String execute(String value) { + return ""; + } + @Override public void init(Config udfConfig, Context context) { LOG.info("Init called with {}", udfConfig); diff --git a/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java b/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java index 20820f6825..c0de212ed6 100644 --- a/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java +++ b/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java @@ -762,13 +762,13 @@ public void testEndToEndStreamTableInnerJoinWithPrimaryKey() throws Exception { @Test public void testEndToEndStreamTableInnerJoinWithUdf() throws Exception { - int numMessages = 20; + int numMessages = 10; TestAvroSystemFactory.messages.clear(); Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql = "Insert into testavro.enrichedPageViewTopic " - + "select pv.pageKey as __key__, pv.pageKey as pageKey, coalesce(null, 'N/A') as companyName," + + "select MyTest(pv.pageKey) as __key__, pv.pageKey as pageKey, coalesce(null, 'N/A') as companyName," + " p.name as profileName, p.address as profileAddress " + "from testavro.PROFILE.`$table` as p " + "join testavro.PAGEVIEW as pv " From a3424b428ad32d6e5921243d9a0621d7c0bb9dab Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Fri, 13 Sep 2019 01:06:39 -0700 Subject: [PATCH 03/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../src/test/java/org/apache/samza/sql/util/MyTestUdf.java | 5 ----- .../org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java | 2 +- 2 files changed, 1 insertion(+), 6 deletions(-) diff --git a/samza-sql/src/test/java/org/apache/samza/sql/util/MyTestUdf.java b/samza-sql/src/test/java/org/apache/samza/sql/util/MyTestUdf.java index 9716df59e9..983df51000 100644 --- a/samza-sql/src/test/java/org/apache/samza/sql/util/MyTestUdf.java +++ b/samza-sql/src/test/java/org/apache/samza/sql/util/MyTestUdf.java @@ -42,11 +42,6 @@ public Integer execute(Integer value) { return value * 2; } - @SamzaSqlUdfMethod(params = SamzaSqlFieldType.STRING) - public String execute(String value) { - return ""; - } - @Override public void init(Config udfConfig, Context context) { LOG.info("Init called with {}", udfConfig); diff --git a/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java b/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java index c0de212ed6..ea3181086e 100644 --- a/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java +++ b/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java @@ -768,7 +768,7 @@ public void testEndToEndStreamTableInnerJoinWithUdf() throws Exception { Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql = "Insert into testavro.enrichedPageViewTopic " - + "select MyTest(pv.pageKey) as __key__, pv.pageKey as pageKey, coalesce(null, 'N/A') as companyName," + + "select pv.pageKey as __key__, pv.pageKey as pageKey, coalesce(null, 'N/A') as companyName," + " p.name as profileName, p.address as profileAddress " + "from testavro.PROFILE.`$table` as p " + "join testavro.PAGEVIEW as pv " From a2d9d6921e76c6237c5a42563a529ce5535b83c7 Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Fri, 13 Sep 2019 01:07:15 -0700 Subject: [PATCH 04/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java b/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java index ea3181086e..20820f6825 100644 --- a/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java +++ b/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java @@ -762,7 +762,7 @@ public void testEndToEndStreamTableInnerJoinWithPrimaryKey() throws Exception { @Test public void testEndToEndStreamTableInnerJoinWithUdf() throws Exception { - int numMessages = 10; + int numMessages = 20; TestAvroSystemFactory.messages.clear(); Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); From 63185dca4f5b1bec9f225c8ba49a4cd961e8748a Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Fri, 13 Sep 2019 15:47:10 -0700 Subject: [PATCH 05/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../samza/sql/avro/AvroTypeFactoryImpl.java | 25 +++++++++---------- .../samza/sql/interfaces/SqlIOConfig.java | 8 +++--- .../samza/sql/planner/SamzaSqlValidator.java | 22 +++++++--------- .../sql/planner/TestSamzaSqlValidator.java | 12 +++++++++ 4 files changed, 37 insertions(+), 30 deletions(-) diff --git a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java index 9de2d978f5..6c207c11da 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java @@ -31,6 +31,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; + /** * Factory that creates {@link SqlSchema} from the Avro Schema. This is used by the * {@link AvroRelConverter} to convert Avro schema to Samza Sql schema. @@ -44,37 +45,35 @@ public AvroTypeFactoryImpl() { } public SqlSchema createType(Schema schema) { + validateTopLevelAvroType(schema); + return convertSchema(schema.getFields()); + } + + protected void validateTopLevelAvroType(Schema schema) { Schema.Type type = schema.getType(); if (type != Schema.Type.RECORD) { String msg = - String.format("System supports only RECORD as top level avro type, But the Schema's type is %s", type); + String.format("Samza Sql supports only RECORD as top level avro type, But the Schema's type is %s", type); LOG.error(msg); throw new SamzaException(msg); } - - return convertSchema(schema.getFields(), true); } - protected boolean isOptional(Schema.Field field, boolean isTopLevel) { - return field.defaultValue() != null; - } - - private SqlSchema convertSchema(List fields, boolean isTopLevel) { - + protected SqlSchema convertSchema(List fields) { SqlSchemaBuilder schemaBuilder = SqlSchemaBuilder.builder(); for (Schema.Field field : fields) { - SqlFieldSchema fieldSchema = convertField(field.schema(), false, isOptional(field, isTopLevel)); + SqlFieldSchema fieldSchema = convertField(field.schema(), false, field.defaultValue() != null); schemaBuilder.addField(field.name(), fieldSchema); } return schemaBuilder.build(); } - private SqlFieldSchema convertField(Schema fieldSchema) { + protected SqlFieldSchema convertField(Schema fieldSchema) { return convertField(fieldSchema, false, false); } - private SqlFieldSchema convertField(Schema fieldSchema, boolean isNullable, boolean isOptional) { + protected SqlFieldSchema convertField(Schema fieldSchema, boolean isNullable, boolean isOptional) { switch (fieldSchema.getType()) { case ARRAY: SqlFieldSchema elementSchema = convertField(fieldSchema.getElementType()); @@ -101,7 +100,7 @@ private SqlFieldSchema convertField(Schema fieldSchema, boolean isNullable, bool case LONG: return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.INT64, isNullable, isOptional); case RECORD: - SqlSchema rowSchema = convertSchema(fieldSchema.getFields(), false); + SqlSchema rowSchema = convertSchema(fieldSchema.getFields()); return SqlFieldSchema.createRowFieldSchema(rowSchema, isNullable, isOptional); case MAP: // Can the value type be nullable and have default values ? Guess not! diff --git a/samza-sql/src/main/java/org/apache/samza/sql/interfaces/SqlIOConfig.java b/samza-sql/src/main/java/org/apache/samza/sql/interfaces/SqlIOConfig.java index 0761f47f7b..69d301d6c7 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/interfaces/SqlIOConfig.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/interfaces/SqlIOConfig.java @@ -81,13 +81,13 @@ public SqlIOConfig(String systemName, String streamName, List sourcePart this.streamId = String.format("%s-%s", systemName, streamName); samzaRelConverterName = streamConfigs.get(CFG_SAMZA_REL_CONVERTER); - Validate.notEmpty(samzaRelConverterName, - String.format("%s is not set or empty for system %s", CFG_SAMZA_REL_CONVERTER, systemName)); + Validate.notEmpty(samzaRelConverterName, String.format("System %s is unknown. %s is not set or empty for this" + + " system", systemName, CFG_SAMZA_REL_CONVERTER)); if (isRemoteTable()) { samzaRelTableKeyConverterName = streamConfigs.get(CFG_SAMZA_REL_TABLE_KEY_CONVERTER); - Validate.notEmpty(samzaRelTableKeyConverterName, - String.format("%s is not set or empty for system %s", CFG_SAMZA_REL_CONVERTER, systemName)); + Validate.notEmpty(samzaRelTableKeyConverterName, String.format("System %s is unknown. %s is not set or empty for" + + " this system", systemName, CFG_SAMZA_REL_CONVERTER)); } else { samzaRelTableKeyConverterName = ""; } diff --git a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java index f81600d933..5b5a4f6ca2 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java @@ -80,20 +80,18 @@ public void validate(List sqlStmts) throws SamzaSqlValidatorException { try { relRoot = planner.plan(qinfo.getSelectQuery()); } catch (SamzaException e) { - throw new SamzaSqlValidatorException("Calcite planning for sql failed.", e); + throw new SamzaSqlValidatorException(e); } // Now that we have logical plan, validate different aspects. - validate(relRoot, qinfo, sqlConfig); + validate(relRoot, qinfo.getSink(), sqlConfig); } } - private void validate(RelRoot relRoot, SamzaSqlQueryParser.QueryInfo qinfo, SamzaSqlApplicationConfig sqlConfig) + protected void validate(RelRoot relRoot, String sink, SamzaSqlApplicationConfig sqlConfig) throws SamzaSqlValidatorException { - if (!skipOutputValidation(relRoot, qinfo, sqlConfig)) { - // Validate select fields (including Udf return types) with output schema - validateOutput(relRoot, sqlConfig.getRelSchemaProviders().get(qinfo.getSink())); - } + // Validate select fields (including Udf return types) with output schema + validateOutput(relRoot, sqlConfig.getRelSchemaProviders().get(sink)); // TODO: // 1. SAMZA-2314: Validate Udf arguments. @@ -101,7 +99,7 @@ private void validate(RelRoot relRoot, SamzaSqlQueryParser.QueryInfo qinfo, Samz // Eg: LogicalAggregate with sum function is not supported by Samza Sql. } - private void validateOutput(RelRoot relRoot, RelSchemaProvider outputRelSchemaProvider) + protected void validateOutput(RelRoot relRoot, RelSchemaProvider outputRelSchemaProvider) throws SamzaSqlValidatorException { LogicalProject project = (LogicalProject) relRoot.rel; RelRecordType projetRecord = (RelRecordType) project.getRowType(); @@ -116,11 +114,9 @@ private void validateOutput(RelRoot relRoot, RelSchemaProvider outputRelSchemaPr LOG.info("Samza Sql Validation finished successfully."); } - protected boolean skipOutputValidation(RelRoot relRoot, SamzaSqlQueryParser.QueryInfo qinfo, - SamzaSqlApplicationConfig sqlConfig) { - return false; - } - + // TODO: Remove this API. This API is introduced to take care of cases where RelSchemaProviders have a complex + // mechanism to determine if a given output field is optional. Once the RelSchemaProviders are fixed to properly + // mark the fields as optional, we can remove this API. protected boolean isOptional(RelSchemaProvider outputRelSchemaProvider, String outputFieldName, RelRecordType projectRecord) { return false; diff --git a/samza-sql/src/test/java/org/apache/samza/sql/planner/TestSamzaSqlValidator.java b/samza-sql/src/test/java/org/apache/samza/sql/planner/TestSamzaSqlValidator.java index ce783be4f4..7a67d50181 100644 --- a/samza-sql/src/test/java/org/apache/samza/sql/planner/TestSamzaSqlValidator.java +++ b/samza-sql/src/test/java/org/apache/samza/sql/planner/TestSamzaSqlValidator.java @@ -83,6 +83,18 @@ public void testRepeatedThriceFieldsValidation() throws SamzaSqlValidatorExcepti new SamzaSqlValidator(samzaConfig).validate(sqlStmts); } + @Test (expected = SamzaSqlValidatorException.class) + public void testFieldEndingInZeroValidation() throws SamzaSqlValidatorException { + Map config = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(1); + config.put(SamzaSqlApplicationConfig.CFG_SQL_STMT, + "Insert into testavro.outputTopic select id, true as bool_value, false as non_existing_name0" + + " from testavro.level1.level2.SIMPLE1 as s where s.id = 1"); + Config samzaConfig = SamzaSqlApplicationRunner.computeSamzaConfigs(true, new MapConfig(config)); + + List sqlStmts = fetchSqlFromConfig(config); + new SamzaSqlValidator(samzaConfig).validate(sqlStmts); + } + @Test (expected = SamzaSqlValidatorException.class) public void testNonExistingOutputField() throws SamzaSqlValidatorException { Map config = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(1); From 9d2e1b0aff4e1d36aac21e88b4d96e6afeab3566 Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Sat, 14 Sep 2019 06:31:14 -0700 Subject: [PATCH 06/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../apache/samza/sql/avro/AvroTypeFactoryImpl.java | 3 ++- .../samza/sql/planner/SamzaSqlValidator.java | 14 +++++++------- 2 files changed, 9 insertions(+), 8 deletions(-) diff --git a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java index 6c207c11da..c4299012fe 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java @@ -62,7 +62,8 @@ protected void validateTopLevelAvroType(Schema schema) { protected SqlSchema convertSchema(List fields) { SqlSchemaBuilder schemaBuilder = SqlSchemaBuilder.builder(); for (Schema.Field field : fields) { - SqlFieldSchema fieldSchema = convertField(field.schema(), false, field.defaultValue() != null); + boolean isOptional = (field.defaultValue() != null); + SqlFieldSchema fieldSchema = convertField(field.schema(), false, isOptional); schemaBuilder.addField(field.name(), fieldSchema); } diff --git a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java index 5b5a4f6ca2..87618c3855 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java @@ -116,13 +116,13 @@ protected void validateOutput(RelRoot relRoot, RelSchemaProvider outputRelSchema // TODO: Remove this API. This API is introduced to take care of cases where RelSchemaProviders have a complex // mechanism to determine if a given output field is optional. Once the RelSchemaProviders are fixed to properly - // mark the fields as optional, we can remove this API. + // mark the fields as optional, this API will be removed. protected boolean isOptional(RelSchemaProvider outputRelSchemaProvider, String outputFieldName, RelRecordType projectRecord) { return false; } - private void validateOutputRecords(RelRecordType outputRecord, SqlSchema outputSqlSchema, + protected void validateOutputRecords(RelRecordType outputRecord, SqlSchema outputSqlSchema, RelRecordType projectRecord, RelSchemaProvider outputRelSchemaProvider) throws SamzaSqlValidatorException { Map outputRecordMap = outputRecord.getFieldList().stream().collect( @@ -144,6 +144,9 @@ private void validateOutputRecords(RelRecordType outputRecord, SqlSchema outputS // Take the following example: SELECT id as str, secondaryId as str, tertiaryId as str FROM store.myTable // Calcite renames the projected fieldNames in select query as str, str0, str1 respectively. // Samza Sql allows a field name to be specified up to 2 times. Do the validation accordingly. + + // This type of pattern is typically followed when users want to just modify one field in the input table while + // keeping rest of the fields the same. Eg: SELECT myUdf(id) as id, * from store.myTable if (projectedFieldName.endsWith("0")) { projectedFieldName = StringUtils.chop(projectedFieldName); outputFieldType = outputRecordMap.get(projectedFieldName); @@ -198,7 +201,7 @@ private void validateOutputRecords(RelRecordType outputRecord, SqlSchema outputS } } - private boolean compareFieldTypes(RelDataType outputFieldType, SqlFieldSchema sqlFieldSchema, + protected boolean compareFieldTypes(RelDataType outputFieldType, SqlFieldSchema sqlFieldSchema, RelDataType selectQueryFieldType, RelSchemaProvider outputRelSchemaProvider) { RelDataType projectFieldType; @@ -279,10 +282,7 @@ public static String formatErrorString(String query, Exception e) { Matcher matcher = pattern.matcher(e.getMessage()); String[] queryLines = query.split("\\n"); StringBuilder result = new StringBuilder(); - int startColIdx; - int endColIdx; - int startLineIdx; - int endLineIdx; + int startColIdx, endColIdx, startLineIdx, endLineIdx; try { if (matcher.find()) { From d89ca2d039af7d94e73960537a73e6fd685700ed Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Mon, 16 Sep 2019 07:32:41 -0700 Subject: [PATCH 07/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../java/org/apache/samza/sql/translator/FilterTranslator.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/samza-sql/src/main/java/org/apache/samza/sql/translator/FilterTranslator.java b/samza-sql/src/main/java/org/apache/samza/sql/translator/FilterTranslator.java index 29253a5ef6..be861479e0 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/translator/FilterTranslator.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/translator/FilterTranslator.java @@ -19,8 +19,6 @@ package org.apache.samza.sql.translator; -import java.time.Duration; -import java.time.Instant; import java.util.Arrays; import java.util.Collections; import org.apache.calcite.rel.logical.LogicalFilter; From a41cfabceb00e936f18b5b1e863b4af78fd5134d Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Tue, 17 Sep 2019 18:26:39 -0700 Subject: [PATCH 08/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../samza/sql/avro/AvroTypeFactoryImpl.java | 14 +++----- .../samza/sql/runner/SamzaSqlApplication.java | 1 - .../sql/translator/FilterTranslator.java | 2 +- .../samza/sql/avro/TestAvroRelConversion.java | 6 ++-- .../samza/sql/avro/schemas/ComplexRecord.avsc | 2 +- .../samza/sql/avro/schemas/ComplexRecord.java | 8 ++--- .../sql/planner/TestSamzaSqlValidator.java | 32 +++++++++++-------- .../test/samzasql/TestSamzaSqlEndToEnd.java | 6 ++-- 8 files changed, 36 insertions(+), 35 deletions(-) diff --git a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java index c4299012fe..ee441c1cec 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java @@ -49,7 +49,7 @@ public SqlSchema createType(Schema schema) { return convertSchema(schema.getFields()); } - protected void validateTopLevelAvroType(Schema schema) { + public static void validateTopLevelAvroType(Schema schema) { Schema.Type type = schema.getType(); if (type != Schema.Type.RECORD) { String msg = @@ -59,7 +59,7 @@ protected void validateTopLevelAvroType(Schema schema) { } } - protected SqlSchema convertSchema(List fields) { + public static SqlSchema convertSchema(List fields) { SqlSchemaBuilder schemaBuilder = SqlSchemaBuilder.builder(); for (Schema.Field field : fields) { boolean isOptional = (field.defaultValue() != null); @@ -70,14 +70,10 @@ protected SqlSchema convertSchema(List fields) { return schemaBuilder.build(); } - protected SqlFieldSchema convertField(Schema fieldSchema) { - return convertField(fieldSchema, false, false); - } - - protected SqlFieldSchema convertField(Schema fieldSchema, boolean isNullable, boolean isOptional) { + public static SqlFieldSchema convertField(Schema fieldSchema, boolean isNullable, boolean isOptional) { switch (fieldSchema.getType()) { case ARRAY: - SqlFieldSchema elementSchema = convertField(fieldSchema.getElementType()); + SqlFieldSchema elementSchema = convertField(fieldSchema.getElementType(), false, false); return SqlFieldSchema.createArraySchema(elementSchema, isNullable, isOptional); case BOOLEAN: return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.BOOLEAN, isNullable, isOptional); @@ -114,7 +110,7 @@ protected SqlFieldSchema convertField(Schema fieldSchema, boolean isNullable, bo } } - private SqlFieldSchema getSqlTypeFromUnionTypes(List types, boolean isNullable, boolean isOptional) { + private static SqlFieldSchema getSqlTypeFromUnionTypes(List types, boolean isNullable, boolean isOptional) { // Typically a nullable field's schema is configured as an union of Null and a Type. // This is to check whether the Union is a Nullable field if (types.size() == 2) { diff --git a/samza-sql/src/main/java/org/apache/samza/sql/runner/SamzaSqlApplication.java b/samza-sql/src/main/java/org/apache/samza/sql/runner/SamzaSqlApplication.java index 4304b65337..d13529fec5 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/runner/SamzaSqlApplication.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/runner/SamzaSqlApplication.java @@ -28,7 +28,6 @@ import org.apache.samza.application.StreamApplication; import org.apache.samza.application.descriptors.StreamApplicationDescriptor; import org.apache.samza.context.ApplicationContainerContext; -import org.apache.samza.context.ApplicationTaskContext; import org.apache.samza.context.ApplicationTaskContextFactory; import org.apache.samza.context.ContainerContext; import org.apache.samza.context.ExternalContext; diff --git a/samza-sql/src/main/java/org/apache/samza/sql/translator/FilterTranslator.java b/samza-sql/src/main/java/org/apache/samza/sql/translator/FilterTranslator.java index be861479e0..6515dc209c 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/translator/FilterTranslator.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/translator/FilterTranslator.java @@ -112,7 +112,7 @@ public boolean apply(SamzaSqlRelMessage message) { updateMetrics(startProcessing, retVal, System.nanoTime()); return retVal; } else { - LOG.error("return value is not boolean for rel message: %s", message); + LOG.error("return value is not boolean for rel message: {}", message); return false; } } diff --git a/samza-sql/src/test/java/org/apache/samza/sql/avro/TestAvroRelConversion.java b/samza-sql/src/test/java/org/apache/samza/sql/avro/TestAvroRelConversion.java index b7b3dc3ec8..7bb001f046 100644 --- a/samza-sql/src/test/java/org/apache/samza/sql/avro/TestAvroRelConversion.java +++ b/samza-sql/src/test/java/org/apache/samza/sql/avro/TestAvroRelConversion.java @@ -199,7 +199,7 @@ public void testComplexRecordConversion() throws IOException { record.put("id", id); record.put("bool_value", boolValue); record.put("double_value", doubleValue); - record.put("float_value", floatValue); + record.put("float_value0", floatValue); record.put("string_value", testStrValue); record.put("bytes_value", testBytes); record.put("fixed_value", fixedBytes); @@ -212,7 +212,7 @@ public void testComplexRecordConversion() throws IOException { complexRecord.id = id; complexRecord.bool_value = boolValue; complexRecord.double_value = doubleValue; - complexRecord.float_value = floatValue; + complexRecord.float_value0 = floatValue; complexRecord.string_value = testStrValue; complexRecord.bytes_value = testBytes; complexRecord.fixed_value = fixedBytes; @@ -352,7 +352,7 @@ private void validateAvroSerializedData(byte[] serializedData, Object unionValue Assert.assertEquals(message.getSamzaSqlRelRecord().getField("bool_value").get(), boolValue); Assert.assertEquals(message.getSamzaSqlRelRecord().getField("double_value").get(), doubleValue); Assert.assertEquals(message.getSamzaSqlRelRecord().getField("string_value").get(), new Utf8(testStrValue)); - Assert.assertEquals(message.getSamzaSqlRelRecord().getField("float_value").get(), floatValue); + Assert.assertEquals(message.getSamzaSqlRelRecord().getField("float_value0").get(), doubleValue); Assert.assertEquals(message.getSamzaSqlRelRecord().getField("long_value").get(), longValue); if (unionValue instanceof String) { Assert.assertEquals(message.getSamzaSqlRelRecord().getField("union_value").get(), new Utf8((String) unionValue)); diff --git a/samza-sql/src/test/java/org/apache/samza/sql/avro/schemas/ComplexRecord.avsc b/samza-sql/src/test/java/org/apache/samza/sql/avro/schemas/ComplexRecord.avsc index c307b10f0f..e2e67e2b3b 100644 --- a/samza-sql/src/test/java/org/apache/samza/sql/avro/schemas/ComplexRecord.avsc +++ b/samza-sql/src/test/java/org/apache/samza/sql/avro/schemas/ComplexRecord.avsc @@ -40,7 +40,7 @@ "default":null }, { - "name": "float_value", + "name": "float_value0", "doc": "float Value.", "type": ["null", "float"], "default":null diff --git a/samza-sql/src/test/java/org/apache/samza/sql/avro/schemas/ComplexRecord.java b/samza-sql/src/test/java/org/apache/samza/sql/avro/schemas/ComplexRecord.java index 91a447f641..d21fc0fc47 100644 --- a/samza-sql/src/test/java/org/apache/samza/sql/avro/schemas/ComplexRecord.java +++ b/samza-sql/src/test/java/org/apache/samza/sql/avro/schemas/ComplexRecord.java @@ -26,7 +26,7 @@ @SuppressWarnings("all") public class ComplexRecord extends org.apache.avro.specific.SpecificRecordBase implements org.apache.avro.specific.SpecificRecord { - public static final org.apache.avro.Schema SCHEMA$ = org.apache.avro.Schema.parse("{\"type\":\"record\",\"name\":\"ComplexRecord\",\"namespace\":\"org.apache.samza.sql.avro.schemas\",\"fields\":[{\"name\":\"id\",\"type\":\"int\",\"doc\":\"Record id.\"},{\"name\":\"bool_value\",\"type\":[\"null\",\"boolean\"],\"doc\":\"Boolean Value.\"},{\"name\":\"double_value\",\"type\":[\"null\",\"double\"],\"doc\":\"double Value.\",\"default\":null},{\"name\":\"float_value\",\"type\":[\"null\",\"float\"],\"doc\":\"float Value.\",\"default\":null},{\"name\":\"string_value\",\"type\":[\"null\",\"string\"],\"doc\":\"string Value.\",\"default\":null},{\"name\":\"bytes_value\",\"type\":[\"null\",\"bytes\"],\"doc\":\"bytes Value.\",\"default\":null},{\"name\":\"long_value\",\"type\":[\"null\",\"long\"],\"doc\":\"long Value.\",\"default\":null},{\"name\":\"fixed_value\",\"type\":[\"null\",{\"type\":\"fixed\",\"name\":\"MyFixed\",\"size\":16}],\"doc\":\"fixed Value.\",\"default\":null},{\"name\":\"array_values\",\"type\":[\"null\",{\"type\":\"array\",\"items\":\"string\"}],\"doc\":\"array values in the record.\",\"default\":[]},{\"name\":\"map_values\",\"type\":[\"null\",{\"type\":\"map\",\"values\":\"string\"}],\"doc\":\"map values in the record.\",\"default\":[]},{\"name\":\"enum_value\",\"type\":[\"null\",{\"type\":\"enum\",\"name\":\"TestEnumType\",\"symbols\":[\"foo\",\"bar\"]}],\"doc\":\"enum value.\",\"default\":[]},{\"name\":\"empty_record\",\"type\":[\"null\",{\"type\":\"record\",\"name\":\"emptySubRecord\",\"fields\":[]}],\"default\":null},{\"name\":\"array_records\",\"type\":[\"null\",{\"type\":\"record\",\"name\":\"SubRecord\",\"fields\":[{\"name\":\"id\",\"type\":[\"null\",\"int\"],\"doc\":\"sub record id\"},{\"name\":\"sub_values\",\"type\":{\"type\":\"array\",\"items\":\"string\"},\"doc\":\"Sub record \"}]}],\"doc\":\"array of records.\",\"default\":[]},{\"name\":\"union_value\",\"type\":[\"null\",\"SubRecord\",\"string\"],\"doc\":\"union Value.\",\"default\":null}]}"); + public static final org.apache.avro.Schema SCHEMA$ = org.apache.avro.Schema.parse("{\"type\":\"record\",\"name\":\"ComplexRecord\",\"namespace\":\"org.apache.samza.sql.avro.schemas\",\"fields\":[{\"name\":\"id\",\"type\":\"int\",\"doc\":\"Record id.\"},{\"name\":\"bool_value\",\"type\":[\"null\",\"boolean\"],\"doc\":\"Boolean Value.\"},{\"name\":\"double_value\",\"type\":[\"null\",\"double\"],\"doc\":\"double Value.\",\"default\":null},{\"name\":\"float_value0\",\"type\":[\"null\",\"float\"],\"doc\":\"float Value.\",\"default\":null},{\"name\":\"string_value\",\"type\":[\"null\",\"string\"],\"doc\":\"string Value.\",\"default\":null},{\"name\":\"bytes_value\",\"type\":[\"null\",\"bytes\"],\"doc\":\"bytes Value.\",\"default\":null},{\"name\":\"long_value\",\"type\":[\"null\",\"long\"],\"doc\":\"long Value.\",\"default\":null},{\"name\":\"fixed_value\",\"type\":[\"null\",{\"type\":\"fixed\",\"name\":\"MyFixed\",\"size\":16}],\"doc\":\"fixed Value.\",\"default\":null},{\"name\":\"array_values\",\"type\":[\"null\",{\"type\":\"array\",\"items\":\"string\"}],\"doc\":\"array values in the record.\",\"default\":[]},{\"name\":\"map_values\",\"type\":[\"null\",{\"type\":\"map\",\"values\":\"string\"}],\"doc\":\"map values in the record.\",\"default\":[]},{\"name\":\"enum_value\",\"type\":[\"null\",{\"type\":\"enum\",\"name\":\"TestEnumType\",\"symbols\":[\"foo\",\"bar\"]}],\"doc\":\"enum value.\",\"default\":[]},{\"name\":\"empty_record\",\"type\":[\"null\",{\"type\":\"record\",\"name\":\"emptySubRecord\",\"fields\":[]}],\"default\":null},{\"name\":\"array_records\",\"type\":[\"null\",{\"type\":\"record\",\"name\":\"SubRecord\",\"fields\":[{\"name\":\"id\",\"type\":[\"null\",\"int\"],\"doc\":\"sub record id\"},{\"name\":\"sub_values\",\"type\":{\"type\":\"array\",\"items\":\"string\"},\"doc\":\"Sub record \"}]}],\"doc\":\"array of records.\",\"default\":[]},{\"name\":\"union_value\",\"type\":[\"null\",\"SubRecord\",\"string\"],\"doc\":\"union Value.\",\"default\":null}]}"); /** Record id. */ public java.lang.Integer id; /** Boolean Value. */ @@ -34,7 +34,7 @@ public class ComplexRecord extends org.apache.avro.specific.SpecificRecordBase i /** double Value. */ public java.lang.Double double_value; /** float Value. */ - public java.lang.Float float_value; + public java.lang.Float float_value0; /** string Value. */ public java.lang.CharSequence string_value; /** bytes Value. */ @@ -61,7 +61,7 @@ public java.lang.Object get(int field$) { case 0: return id; case 1: return bool_value; case 2: return double_value; - case 3: return float_value; + case 3: return float_value0; case 4: return string_value; case 5: return bytes_value; case 6: return long_value; @@ -82,7 +82,7 @@ public void put(int field$, java.lang.Object value$) { case 0: id = (java.lang.Integer)value$; break; case 1: bool_value = (java.lang.Boolean)value$; break; case 2: double_value = (java.lang.Double)value$; break; - case 3: float_value = (java.lang.Float)value$; break; + case 3: float_value0 = (java.lang.Float)value$; break; case 4: string_value = (java.lang.CharSequence)value$; break; case 5: bytes_value = (java.nio.ByteBuffer)value$; break; case 6: long_value = (java.lang.Long)value$; break; diff --git a/samza-sql/src/test/java/org/apache/samza/sql/planner/TestSamzaSqlValidator.java b/samza-sql/src/test/java/org/apache/samza/sql/planner/TestSamzaSqlValidator.java index 7a67d50181..38a18a1a88 100644 --- a/samza-sql/src/test/java/org/apache/samza/sql/planner/TestSamzaSqlValidator.java +++ b/samza-sql/src/test/java/org/apache/samza/sql/planner/TestSamzaSqlValidator.java @@ -59,24 +59,27 @@ public void testBasicValidation() throws SamzaSqlValidatorException { new SamzaSqlValidator(samzaConfig).validate(sqlStmts); } + // Samza Sql allows users to replace a field in the input stream. For eg: To always set bool_value to false + // while keeping the values of other fields the same, it could be written the below way. + // SELECT false AS bool_value, c.* FROM testavro.COMPLEX1 AS c @Test public void testRepeatedTwiceFieldsValidation() throws SamzaSqlValidatorException { Map config = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(1); config.put(SamzaSqlApplicationConfig.CFG_SQL_STMT, - "Insert into testavro.outputTopic select id, true as bool_value, false as bool_value" - + " from testavro.level1.level2.SIMPLE1 as s where s.id = 1"); + "Insert into testavro.outputTopic select false as bool_value, c.* from testavro.COMPLEX1 as c"); Config samzaConfig = SamzaSqlApplicationRunner.computeSamzaConfigs(true, new MapConfig(config)); List sqlStmts = fetchSqlFromConfig(config); new SamzaSqlValidator(samzaConfig).validate(sqlStmts); } + // Samza Sql allows a field to be replaced only once and validation will fail if the field is replaced more than + // once. We disallow it to keep things simple. @Test (expected = SamzaSqlValidatorException.class) public void testRepeatedThriceFieldsValidation() throws SamzaSqlValidatorException { Map config = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(1); config.put(SamzaSqlApplicationConfig.CFG_SQL_STMT, - "Insert into testavro.outputTopic select id, true as bool_value, false as bool_value, true as bool_value" - + " from testavro.level1.level2.SIMPLE1 as s where s.id = 1"); + "Insert into testavro.outputTopic select id, bool_value, true as bool_value, c.* from testavro.COMPLEX1 as c"); Config samzaConfig = SamzaSqlApplicationRunner.computeSamzaConfigs(true, new MapConfig(config)); List sqlStmts = fetchSqlFromConfig(config); @@ -84,7 +87,7 @@ public void testRepeatedThriceFieldsValidation() throws SamzaSqlValidatorExcepti } @Test (expected = SamzaSqlValidatorException.class) - public void testFieldEndingInZeroValidation() throws SamzaSqlValidatorException { + public void testIllegitFieldEndingInZeroValidation() throws SamzaSqlValidatorException { Map config = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(1); config.put(SamzaSqlApplicationConfig.CFG_SQL_STMT, "Insert into testavro.outputTopic select id, true as bool_value, false as non_existing_name0" @@ -95,25 +98,28 @@ public void testFieldEndingInZeroValidation() throws SamzaSqlValidatorException new SamzaSqlValidator(samzaConfig).validate(sqlStmts); } - @Test (expected = SamzaSqlValidatorException.class) - public void testNonExistingOutputField() throws SamzaSqlValidatorException { + @Test + public void testLegitFieldEndingInZeroValidation() throws SamzaSqlValidatorException { Map config = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(1); config.put(SamzaSqlApplicationConfig.CFG_SQL_STMT, - "Insert into testavro.outputTopic(id) select id, name as strings_value" - + " from testavro.level1.level2.SIMPLE1 as s where s.id = 1"); + "Insert into testavro.outputTopic" + + " select id, bool_value, float_value0 from testavro.COMPLEX1"); Config samzaConfig = SamzaSqlApplicationRunner.computeSamzaConfigs(true, new MapConfig(config)); List sqlStmts = fetchSqlFromConfig(config); new SamzaSqlValidator(samzaConfig).validate(sqlStmts); } - @Test(expected = SamzaException.class) - public void testNonExistingSelectField() throws SamzaSqlValidatorException { + @Test (expected = SamzaSqlValidatorException.class) + public void testNonExistingOutputField() throws SamzaSqlValidatorException { Map config = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(1); config.put(SamzaSqlApplicationConfig.CFG_SQL_STMT, - "Insert into testavro.outputTopic(id) select non_existing_field, name as string_value" + "Insert into testavro.outputTopic(id) select id, name as strings_value" + " from testavro.level1.level2.SIMPLE1 as s where s.id = 1"); - SamzaSqlApplicationRunner.computeSamzaConfigs(true, new MapConfig(config)); + Config samzaConfig = SamzaSqlApplicationRunner.computeSamzaConfigs(true, new MapConfig(config)); + + List sqlStmts = fetchSqlFromConfig(config); + new SamzaSqlValidator(samzaConfig).validate(sqlStmts); } @Test(expected = SamzaSqlValidatorException.class) diff --git a/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java b/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java index 20820f6825..52e6a317d9 100644 --- a/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java +++ b/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java @@ -427,8 +427,8 @@ public void testEndToEndFlatten() throws Exception { LOG.info(" Class Path : " + RelOptUtil.class.getProtectionDomain().getCodeSource().getLocation().toURI().getPath()); String sql1 = - "Insert into testavro.outputTopic(string_value, id, bool_value, bytes_value, fixed_value, float_value) " - + " select Flatten(array_values) as string_value, id, NOT(id = 5) as bool_value, bytes_value, fixed_value, float_value " + "Insert into testavro.outputTopic(string_value, id, bool_value, bytes_value, fixed_value, float_value0) " + + " select Flatten(array_values) as string_value, id, NOT(id = 5) as bool_value, bytes_value, fixed_value, float_value0 " + " from testavro.COMPLEX1"; List sqlStmts = Collections.singletonList(sql1); staticConfigs.put(SamzaSqlApplicationConfig.CFG_SQL_STMTS_JSON, JsonUtil.toJson(sqlStmts)); @@ -458,7 +458,7 @@ public void testEndToEndComplexRecord() throws SamzaSqlValidatorException { String sql1 = "Insert into testavro.outputTopic" + " select bool_value, map_values['key0'] as string_value, union_value, array_values, map_values, id, bytes_value," - + " fixed_value, float_value from testavro.COMPLEX1"; + + " fixed_value, float_value0 from testavro.COMPLEX1"; List sqlStmts = Collections.singletonList(sql1); staticConfigs.put(SamzaSqlApplicationConfig.CFG_SQL_STMTS_JSON, JsonUtil.toJson(sqlStmts)); From 425031b76813588790c570417e1f3334048bd13f Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Tue, 1 Oct 2019 09:09:01 -0700 Subject: [PATCH 09/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../apache/samza/sql/planner/SamzaSqlValidator.java | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java index 87618c3855..dba9251607 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java @@ -102,15 +102,17 @@ protected void validate(RelRoot relRoot, String sink, SamzaSqlApplicationConfig protected void validateOutput(RelRoot relRoot, RelSchemaProvider outputRelSchemaProvider) throws SamzaSqlValidatorException { LogicalProject project = (LogicalProject) relRoot.rel; + RelRecordType projetRecord = (RelRecordType) project.getRowType(); RelRecordType outputRecord = (RelRecordType) QueryPlanner.getSourceRelSchema(outputRelSchemaProvider, new RelSchemaConverter()); + // Get Samza Sql schema along with Calcite schema. The reason is that the Calcite schema does not have a way - // to represent optional fields while Samza Sql schema can represent optional fields. This is the only reason that + // to represent optional fields while Samza Sql schema can represent optional fields. This is the reason that // we use SqlSchema in validating output. SqlSchema outputSqlSchema = QueryPlanner.getSourceSqlSchema(outputRelSchemaProvider); - validateOutputRecords(outputRecord, outputSqlSchema, projetRecord, outputRelSchemaProvider); + validateOutputRecords(outputSqlSchema, outputRecord, projetRecord, outputRelSchemaProvider); LOG.info("Samza Sql Validation finished successfully."); } @@ -122,7 +124,7 @@ protected boolean isOptional(RelSchemaProvider outputRelSchemaProvider, String o return false; } - protected void validateOutputRecords(RelRecordType outputRecord, SqlSchema outputSqlSchema, + protected void validateOutputRecords(SqlSchema outputSqlSchema, RelRecordType outputRecord, RelRecordType projectRecord, RelSchemaProvider outputRelSchemaProvider) throws SamzaSqlValidatorException { Map outputRecordMap = outputRecord.getFieldList().stream().collect( @@ -238,7 +240,7 @@ protected boolean compareFieldTypes(RelDataType outputFieldType, SqlFieldSchema return projectSqlType == SqlTypeName.FLOAT; case ROW: try { - validateOutputRecords((RelRecordType) outputFieldType, sqlFieldSchema.getRowSchema(), + validateOutputRecords(sqlFieldSchema.getRowSchema(), (RelRecordType) outputFieldType, (RelRecordType) projectFieldType, outputRelSchemaProvider); } catch (SamzaSqlValidatorException e) { LOG.error("A field in select query does not match with the output schema.", e); From dceb097568dbd4df466d3b2a50bb82de8087679d Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Wed, 2 Oct 2019 22:47:16 -0700 Subject: [PATCH 10/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../samza/sql/planner/SamzaSqlValidator.java | 50 ++++++++++++------- 1 file changed, 33 insertions(+), 17 deletions(-) diff --git a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java index dba9251607..81d9715e09 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java @@ -40,6 +40,7 @@ import org.apache.samza.sql.data.SamzaSqlRelMessage; import org.apache.samza.sql.dsl.SamzaSqlDslConverter; import org.apache.samza.sql.interfaces.RelSchemaProvider; +import org.apache.samza.sql.interfaces.SamzaRelConverter; import org.apache.samza.sql.interfaces.SamzaSqlJavaTypeFactoryImpl; import org.apache.samza.sql.runner.SamzaSqlApplicationConfig; import org.apache.samza.sql.schema.SqlFieldSchema; @@ -80,18 +81,41 @@ public void validate(List sqlStmts) throws SamzaSqlValidatorException { try { relRoot = planner.plan(qinfo.getSelectQuery()); } catch (SamzaException e) { - throw new SamzaSqlValidatorException(e); + throw new SamzaSqlValidatorException(String.format("Validation failed for sql stmt:\n%s\n", sql), e); } // Now that we have logical plan, validate different aspects. - validate(relRoot, qinfo.getSink(), sqlConfig); + String sink = qinfo.getSink(); + validate(relRoot, sqlConfig.getRelSchemaProviders().get(sink), sqlConfig.getSamzaRelConverters().get(sink)); } } - protected void validate(RelRoot relRoot, String sink, SamzaSqlApplicationConfig sqlConfig) - throws SamzaSqlValidatorException { - // Validate select fields (including Udf return types) with output schema - validateOutput(relRoot, sqlConfig.getRelSchemaProviders().get(sink)); + /** + * Determine if validation needs to be done on Calcite plan based on the schema provider and schema converter. + * @param relRoot + * @param outputSchemaProvider + * @param ouputRelSchemaConverter + * @return if the validation needs to be skipped + */ + protected boolean skipOutputValidation(RelRoot relRoot, RelSchemaProvider outputSchemaProvider, + SamzaRelConverter ouputRelSchemaConverter) { + return false; + } + + // TODO: Remove this API. This API is introduced to take care of cases where RelSchemaProviders have a complex + // mechanism to determine if a given output field is optional. We will need system specific validators to take + // care of such cases and once that is introduced, we can get rid of the below API. + protected boolean isOptional(RelSchemaProvider outputRelSchemaProvider, String outputFieldName, + RelRecordType projectRecord) { + return false; + } + + private void validate(RelRoot relRoot, RelSchemaProvider outputSchemaProvider, + SamzaRelConverter outputRelSchemaConverter) throws SamzaSqlValidatorException { + if (!skipOutputValidation(relRoot, outputSchemaProvider, outputRelSchemaConverter)) { + // Validate select fields (including Udf return types) with output schema + validateOutput(relRoot, outputSchemaProvider); + } // TODO: // 1. SAMZA-2314: Validate Udf arguments. @@ -99,7 +123,7 @@ protected void validate(RelRoot relRoot, String sink, SamzaSqlApplicationConfig // Eg: LogicalAggregate with sum function is not supported by Samza Sql. } - protected void validateOutput(RelRoot relRoot, RelSchemaProvider outputRelSchemaProvider) + private void validateOutput(RelRoot relRoot, RelSchemaProvider outputRelSchemaProvider) throws SamzaSqlValidatorException { LogicalProject project = (LogicalProject) relRoot.rel; @@ -116,15 +140,7 @@ protected void validateOutput(RelRoot relRoot, RelSchemaProvider outputRelSchema LOG.info("Samza Sql Validation finished successfully."); } - // TODO: Remove this API. This API is introduced to take care of cases where RelSchemaProviders have a complex - // mechanism to determine if a given output field is optional. Once the RelSchemaProviders are fixed to properly - // mark the fields as optional, this API will be removed. - protected boolean isOptional(RelSchemaProvider outputRelSchemaProvider, String outputFieldName, - RelRecordType projectRecord) { - return false; - } - - protected void validateOutputRecords(SqlSchema outputSqlSchema, RelRecordType outputRecord, + private void validateOutputRecords(SqlSchema outputSqlSchema, RelRecordType outputRecord, RelRecordType projectRecord, RelSchemaProvider outputRelSchemaProvider) throws SamzaSqlValidatorException { Map outputRecordMap = outputRecord.getFieldList().stream().collect( @@ -203,7 +219,7 @@ protected void validateOutputRecords(SqlSchema outputSqlSchema, RelRecordType ou } } - protected boolean compareFieldTypes(RelDataType outputFieldType, SqlFieldSchema sqlFieldSchema, + private boolean compareFieldTypes(RelDataType outputFieldType, SqlFieldSchema sqlFieldSchema, RelDataType selectQueryFieldType, RelSchemaProvider outputRelSchemaProvider) { RelDataType projectFieldType; From 7d314a5087ece4723b2b8216c8e0bcbf874cd46b Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Wed, 2 Oct 2019 23:17:46 -0700 Subject: [PATCH 11/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../org/apache/samza/sql/planner/SamzaSqlValidator.java | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java index 81d9715e09..23966328eb 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java @@ -86,18 +86,19 @@ public void validate(List sqlStmts) throws SamzaSqlValidatorException { // Now that we have logical plan, validate different aspects. String sink = qinfo.getSink(); - validate(relRoot, sqlConfig.getRelSchemaProviders().get(sink), sqlConfig.getSamzaRelConverters().get(sink)); + validate(relRoot, sink, sqlConfig.getRelSchemaProviders().get(sink), sqlConfig.getSamzaRelConverters().get(sink)); } } /** * Determine if validation needs to be done on Calcite plan based on the schema provider and schema converter. * @param relRoot + * @param sink * @param outputSchemaProvider * @param ouputRelSchemaConverter * @return if the validation needs to be skipped */ - protected boolean skipOutputValidation(RelRoot relRoot, RelSchemaProvider outputSchemaProvider, + protected boolean skipOutputValidation(RelRoot relRoot, String sink, RelSchemaProvider outputSchemaProvider, SamzaRelConverter ouputRelSchemaConverter) { return false; } @@ -110,9 +111,9 @@ protected boolean isOptional(RelSchemaProvider outputRelSchemaProvider, String o return false; } - private void validate(RelRoot relRoot, RelSchemaProvider outputSchemaProvider, + private void validate(RelRoot relRoot, String sink, RelSchemaProvider outputSchemaProvider, SamzaRelConverter outputRelSchemaConverter) throws SamzaSqlValidatorException { - if (!skipOutputValidation(relRoot, outputSchemaProvider, outputRelSchemaConverter)) { + if (!skipOutputValidation(relRoot, sink, outputSchemaProvider, outputRelSchemaConverter)) { // Validate select fields (including Udf return types) with output schema validateOutput(relRoot, outputSchemaProvider); } From 1d8512fdce6172b866fb1fc68efd5afccbb1f4f7 Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Thu, 3 Oct 2019 00:16:47 -0700 Subject: [PATCH 12/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../samza/sql/avro/AvroTypeFactoryImpl.java | 66 +++++++++++-------- 1 file changed, 39 insertions(+), 27 deletions(-) diff --git a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java index ee441c1cec..2171735f8b 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java @@ -1,21 +1,21 @@ /* -* 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. -*/ + * 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.samza.sql.avro; @@ -33,8 +33,8 @@ /** - * Factory that creates {@link SqlSchema} from the Avro Schema. This is used by the - * {@link AvroRelConverter} to convert Avro schema to Samza Sql schema. + * Factory that creates {@link SqlSchema} from the Avro Schema. + * TODO: This class will be renamed to AvroTypeFactoryImpl and copied to open source. */ public class AvroTypeFactoryImpl extends SqlTypeFactoryImpl { @@ -46,10 +46,23 @@ public AvroTypeFactoryImpl() { public SqlSchema createType(Schema schema) { validateTopLevelAvroType(schema); - return convertSchema(schema.getFields()); + return convertSchema(schema.getFields(), true); + } + + /** + * Given a schema field, determine if it is an optional field. There could be cases where a field (like audit header) + * is considered as optional even if it is marked as required in the schema. The producer could be filling in this + * field and hence need not be specified in the query and hence is optional. Typically, such audit headers are + * the top level fields in the schema. + * @param field schema field + * @param isTopLevelField if it is top level field in the schema + * @return if the field is optional + */ + protected boolean isOptional(Schema.Field field, boolean isTopLevelField) { + return field.defaultValue() != null; } - public static void validateTopLevelAvroType(Schema schema) { + private void validateTopLevelAvroType(Schema schema) { Schema.Type type = schema.getType(); if (type != Schema.Type.RECORD) { String msg = @@ -59,18 +72,17 @@ public static void validateTopLevelAvroType(Schema schema) { } } - public static SqlSchema convertSchema(List fields) { + private SqlSchema convertSchema(List fields, boolean isTopLevelField) { SqlSchemaBuilder schemaBuilder = SqlSchemaBuilder.builder(); for (Schema.Field field : fields) { - boolean isOptional = (field.defaultValue() != null); - SqlFieldSchema fieldSchema = convertField(field.schema(), false, isOptional); + SqlFieldSchema fieldSchema = convertField(field.schema(), false, isOptional(field, isTopLevelField)); schemaBuilder.addField(field.name(), fieldSchema); } return schemaBuilder.build(); } - public static SqlFieldSchema convertField(Schema fieldSchema, boolean isNullable, boolean isOptional) { + private SqlFieldSchema convertField(Schema fieldSchema, boolean isNullable, boolean isOptional) { switch (fieldSchema.getType()) { case ARRAY: SqlFieldSchema elementSchema = convertField(fieldSchema.getElementType(), false, false); @@ -97,7 +109,7 @@ public static SqlFieldSchema convertField(Schema fieldSchema, boolean isNullable case LONG: return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.INT64, isNullable, isOptional); case RECORD: - SqlSchema rowSchema = convertSchema(fieldSchema.getFields()); + SqlSchema rowSchema = convertSchema(fieldSchema.getFields(), false); return SqlFieldSchema.createRowFieldSchema(rowSchema, isNullable, isOptional); case MAP: // Can the value type be nullable and have default values ? Guess not! @@ -110,7 +122,7 @@ public static SqlFieldSchema convertField(Schema fieldSchema, boolean isNullable } } - private static SqlFieldSchema getSqlTypeFromUnionTypes(List types, boolean isNullable, boolean isOptional) { + private SqlFieldSchema getSqlTypeFromUnionTypes(List types, boolean isNullable, boolean isOptional) { // Typically a nullable field's schema is configured as an union of Null and a Type. // This is to check whether the Union is a Nullable field if (types.size() == 2) { From cd832202c5152c3a1a6c757fd9267b03f75bddc0 Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Thu, 3 Oct 2019 00:20:10 -0700 Subject: [PATCH 13/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../java/org/apache/samza/sql/avro/TestAvroRelConversion.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/samza-sql/src/test/java/org/apache/samza/sql/avro/TestAvroRelConversion.java b/samza-sql/src/test/java/org/apache/samza/sql/avro/TestAvroRelConversion.java index 7bb001f046..102ad52c5f 100644 --- a/samza-sql/src/test/java/org/apache/samza/sql/avro/TestAvroRelConversion.java +++ b/samza-sql/src/test/java/org/apache/samza/sql/avro/TestAvroRelConversion.java @@ -352,7 +352,7 @@ private void validateAvroSerializedData(byte[] serializedData, Object unionValue Assert.assertEquals(message.getSamzaSqlRelRecord().getField("bool_value").get(), boolValue); Assert.assertEquals(message.getSamzaSqlRelRecord().getField("double_value").get(), doubleValue); Assert.assertEquals(message.getSamzaSqlRelRecord().getField("string_value").get(), new Utf8(testStrValue)); - Assert.assertEquals(message.getSamzaSqlRelRecord().getField("float_value0").get(), doubleValue); + Assert.assertEquals(message.getSamzaSqlRelRecord().getField("float_value0").get(), floatValue); Assert.assertEquals(message.getSamzaSqlRelRecord().getField("long_value").get(), longValue); if (unionValue instanceof String) { Assert.assertEquals(message.getSamzaSqlRelRecord().getField("union_value").get(), new Utf8((String) unionValue)); From aa0e4d92caba826c3745618f85299d47bb8da4cf Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Thu, 3 Oct 2019 00:37:34 -0700 Subject: [PATCH 14/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java index 2171735f8b..9729615aef 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java @@ -33,8 +33,8 @@ /** - * Factory that creates {@link SqlSchema} from the Avro Schema. - * TODO: This class will be renamed to AvroTypeFactoryImpl and copied to open source. + * Factory that creates {@link SqlSchema} from the Avro Schema. This is used by the + * {@link AvroRelConverter} to convert Avro schema to Samza Sql schema. */ public class AvroTypeFactoryImpl extends SqlTypeFactoryImpl { From 4ffb086cfd26fccb527c7f0bdb7e73a06d1acf14 Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Thu, 3 Oct 2019 02:46:35 -0700 Subject: [PATCH 15/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java b/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java index 52e6a317d9..581fbdab11 100644 --- a/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java +++ b/samza-test/src/test/java/org/apache/samza/test/samzasql/TestSamzaSqlEndToEnd.java @@ -479,7 +479,7 @@ public void testEndToEndWithFloatToStringConversion() throws Exception { TestAvroSystemFactory.messages.clear(); Map staticConfigs = SamzaSqlTestConfig.fetchStaticConfigsWithFactories(numMessages); String sql1 = "Insert into testavro.outputTopic" - + " select 'urn:li:member:' || cast(cast(float_value as int) as varchar) as string_value, id, float_value, " + + " select 'urn:li:member:' || cast(cast(float_value0 as int) as varchar) as string_value, id, float_value0, " + " double_value, true as bool_value from testavro.COMPLEX1"; List sqlStmts = Arrays.asList(sql1); staticConfigs.put(SamzaSqlApplicationConfig.CFG_SQL_STMTS_JSON, JsonUtil.toJson(sqlStmts)); From b48c683bc97a930baa10d1fd294ecaef06c52916 Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Thu, 3 Oct 2019 02:49:16 -0700 Subject: [PATCH 16/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../java/org/apache/samza/sql/system/TestAvroSystemFactory.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/samza-sql/src/test/java/org/apache/samza/sql/system/TestAvroSystemFactory.java b/samza-sql/src/test/java/org/apache/samza/sql/system/TestAvroSystemFactory.java index 7d62b8febf..65e0ad0a65 100644 --- a/samza-sql/src/test/java/org/apache/samza/sql/system/TestAvroSystemFactory.java +++ b/samza-sql/src/test/java/org/apache/samza/sql/system/TestAvroSystemFactory.java @@ -325,7 +325,7 @@ private Object createComplexRecord(int index) { record.put("id", index); record.put("string_value", "Name" + index); record.put("bytes_value", ByteBuffer.wrap(("sample bytes").getBytes())); - record.put("float_value", index + 0.123456f); + record.put("float_value0", index + 0.123456f); record.put("double_value", index + 0.0123456789); MyFixed myFixedVar = new MyFixed(); myFixedVar.bytes(DEFAULT_TRACKING_ID_BYTES); From af1c18fea0fdbe492147db16ede7ab30da5cdd21 Mon Sep 17 00:00:00 2001 From: Aditya Toomula Date: Thu, 3 Oct 2019 12:52:33 -0700 Subject: [PATCH 17/17] SAMZA-2320: Samza-sql: Refactor validation to cover more cases and make it more extensible. --- .../java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java index 9729615aef..153d96aca6 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/avro/AvroTypeFactoryImpl.java @@ -50,9 +50,9 @@ public SqlSchema createType(Schema schema) { } /** - * Given a schema field, determine if it is an optional field. There could be cases where a field (like audit header) + * Given a schema field, determine if it is an optional field. There could be cases where a field * is considered as optional even if it is marked as required in the schema. The producer could be filling in this - * field and hence need not be specified in the query and hence is optional. Typically, such audit headers are + * field and hence need not be specified in the query and hence is optional. Typically, such fields are * the top level fields in the schema. * @param field schema field * @param isTopLevelField if it is top level field in the schema