Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
d22630b
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Oct 3, 2019
2086ad6
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Sep 13, 2019
a3424b4
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Sep 13, 2019
a2d9d69
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Sep 13, 2019
63185dc
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Sep 13, 2019
9d2e1b0
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Sep 14, 2019
d89ca2d
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Sep 16, 2019
a41cfab
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Sep 18, 2019
425031b
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Oct 1, 2019
dceb097
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Oct 3, 2019
7d314a5
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Oct 3, 2019
1d8512f
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Oct 3, 2019
cd83220
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Oct 3, 2019
863011e
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Oct 3, 2019
aa0e4d9
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Oct 3, 2019
4ffb086
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Oct 3, 2019
b48c683
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Oct 3, 2019
af1c18f
SAMZA-2320: Samza-sql: Refactor validation to cover more cases and ma…
atoomula Oct 3, 2019
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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;

Expand All @@ -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.
Expand All @@ -44,37 +45,47 @@ public AvroTypeFactoryImpl() {
}

public SqlSchema createType(Schema schema) {
validateTopLevelAvroType(schema);
return convertSchema(schema.getFields(), true);
}

/**
* 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 fields 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;
}

private 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());
}

private SqlSchema convertSchema(List<Schema.Field> fields) {

private SqlSchema convertSchema(List<Schema.Field> fields, boolean isTopLevelField) {
SqlSchemaBuilder schemaBuilder = SqlSchemaBuilder.builder();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Where is isTopLevel being used? It's not used in this function or in isOptional?

@atoomula atoomula Sep 13, 2019 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added the reasoning for isTopLevel in javdoc for isOptional()

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();
}

private SqlFieldSchema convertField(Schema fieldSchema) {
return convertField(fieldSchema, false, false);
}

private 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);
Expand All @@ -98,7 +109,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!
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -81,13 +81,13 @@ public SqlIOConfig(String systemName, String streamName, List<String> 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 = "";
}
Expand Down
Loading