Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
Expand Up @@ -24,44 +24,59 @@
*/
public class SqlFieldSchema {

private SamzaSqlFieldType fieldType;
private SqlFieldSchema elementType;
private SqlFieldSchema valueType;
private SqlSchema rowSchema;
private final SamzaSqlFieldType fieldType;
private final SqlFieldSchema elementType;
private final SqlFieldSchema valueType;
private final SqlSchema rowSchema;
// A field is considered nullable when the field could have a null value. Please note that nullable field
// needs to be explicitly set while writing and is expected to be set while reading. A non-nullable field
// cannot have a null value.
private final Boolean isNullable;
// A field is considered optional when the field has a default value. Such a field need not be set while writing
// but is expected to be set while reading.
// Please note that nullable field is also optional field if a default value is set but the value for
// nullable non-optional field need to be explicitly set.
private final Boolean isOptional;

private SqlFieldSchema(SamzaSqlFieldType fieldType, SqlFieldSchema elementType, SqlFieldSchema valueType, SqlSchema rowSchema) {
private SqlFieldSchema(SamzaSqlFieldType fieldType, SqlFieldSchema elementType, SqlFieldSchema valueType,
SqlSchema rowSchema, boolean isNullable, boolean isOptional) {
this.fieldType = fieldType;
this.elementType = elementType;
this.valueType = valueType;
this.rowSchema = rowSchema;
this.isNullable = isNullable;
this.isOptional = isOptional;
}

/**
* Create a primitive fi
* Create a primitive field schema.
* @param typeName
* @return
*/
public static SqlFieldSchema createPrimitiveSchema(SamzaSqlFieldType typeName) {
return new SqlFieldSchema(typeName, null, null, null);
public static SqlFieldSchema createPrimitiveSchema(SamzaSqlFieldType typeName, boolean isNullable,
boolean isOptional) {
return new SqlFieldSchema(typeName, null, null, null, isNullable, isOptional);
}

public static SqlFieldSchema createArraySchema(SqlFieldSchema elementType) {
return new SqlFieldSchema(SamzaSqlFieldType.ARRAY, elementType, null, null);
public static SqlFieldSchema createArraySchema(SqlFieldSchema elementType, boolean isNullable,
boolean isOptional) {
return new SqlFieldSchema(SamzaSqlFieldType.ARRAY, elementType, null, null, isNullable, isOptional);
}

public static SqlFieldSchema createMapSchema(SqlFieldSchema valueType) {
return new SqlFieldSchema(SamzaSqlFieldType.MAP, null, valueType, null);
public static SqlFieldSchema createMapSchema(SqlFieldSchema valueType, boolean isNullable, boolean isOptional) {
return new SqlFieldSchema(SamzaSqlFieldType.MAP, null, valueType, null, isNullable, isOptional);
}

public static SqlFieldSchema createRowFieldSchema(SqlSchema rowSchema) {
return new SqlFieldSchema(SamzaSqlFieldType.ROW, null, null, rowSchema);
public static SqlFieldSchema createRowFieldSchema(SqlSchema rowSchema, boolean isNullable, boolean isOptional) {
return new SqlFieldSchema(SamzaSqlFieldType.ROW, null, null, rowSchema, isNullable, isOptional);
}

/**
* @return whether the field is a primitive field type or not.
*/
public boolean isPrimitiveField() {
return fieldType != SamzaSqlFieldType.ARRAY && fieldType != SamzaSqlFieldType.MAP && fieldType != SamzaSqlFieldType.ROW;
return fieldType != SamzaSqlFieldType.ARRAY && fieldType != SamzaSqlFieldType.MAP &&
fieldType != SamzaSqlFieldType.ROW;
}

/**
Expand Down Expand Up @@ -93,5 +108,17 @@ public SqlSchema getRowSchema() {
return rowSchema;
}

/**
* Get if the field type is nullable.
*/
public boolean isNullable() {
return isNullable;
}

/**
* Get if the field type is optional.
*/
public boolean isOptional() {
return isOptional;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -293,9 +293,9 @@ public List<SqlFunction> listFunctions(ExecutionContext context) {
*/
List<SqlFunction> udfs = new ArrayList<>();
udfs.add(new SamzaSqlUdfDisplayInfo("RegexMatch", "Matches the string to the regex",
Arrays.asList(SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.STRING),
SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.STRING)),
SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.BOOLEAN)));
Arrays.asList(SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.STRING, false, false),
SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.STRING, false, false)),
SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.BOOLEAN, false, false)));

return udfs;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@
import org.apache.avro.Schema;
import org.apache.calcite.rel.type.RelDataTypeSystem;
import org.apache.calcite.sql.type.SqlTypeFactoryImpl;
import org.apache.commons.lang3.Validate;
import org.apache.samza.SamzaException;
import org.apache.samza.sql.schema.SamzaSqlFieldType;
import org.apache.samza.sql.schema.SqlFieldSchema;
Expand All @@ -33,8 +32,8 @@
import org.slf4j.LoggerFactory;

/**
* Factory that creates the Calcite relational types from the Avro Schema. This is used by the
* AvroRelConverter to convert the Avro schema to calcite relational schema.
* 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 {

Expand All @@ -60,62 +59,68 @@ private SqlSchema convertSchema(List<Schema.Field> fields) {

SqlSchemaBuilder schemaBuilder = SqlSchemaBuilder.builder();
for (Schema.Field field : fields) {
SqlFieldSchema fieldSchema = convertField(field.schema());
boolean isOptional = (field.defaultValue() != null);
SqlFieldSchema fieldSchema = convertField(field.schema(), false, isOptional);
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());
return SqlFieldSchema.createArraySchema(elementSchema);
return SqlFieldSchema.createArraySchema(elementSchema, isNullable, isOptional);
case BOOLEAN:
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.BOOLEAN);
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.BOOLEAN, isNullable, isOptional);
case DOUBLE:
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.DOUBLE);
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.DOUBLE, isNullable, isOptional);
case FLOAT:
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.FLOAT);
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.FLOAT, isNullable, isOptional);
case ENUM:
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.STRING);
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.STRING, isNullable, isOptional);
case UNION:
return getSqlTypeFromUnionTypes(fieldSchema.getTypes());
return getSqlTypeFromUnionTypes(fieldSchema.getTypes(), isNullable, isOptional);
case FIXED:
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.BYTES);
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.BYTES, isNullable, isOptional);
case STRING:
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.STRING);
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.STRING, isNullable, isOptional);
case BYTES:
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.BYTES);
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.BYTES, isNullable, isOptional);
case INT:
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.INT32);
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.INT32, isNullable, isOptional);
case LONG:
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.INT64);
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.INT64, isNullable, isOptional);
case RECORD:
SqlSchema rowSchema = convertSchema(fieldSchema.getFields());
return SqlFieldSchema.createRowFieldSchema(rowSchema);
return SqlFieldSchema.createRowFieldSchema(rowSchema, isNullable, isOptional);
case MAP:
SqlFieldSchema valueType = convertField(fieldSchema.getValueType());
return SqlFieldSchema.createMapSchema(valueType);
// Can the value type be nullable and have default values ? Guess not!
SqlFieldSchema valueType = convertField(fieldSchema.getValueType(), false, false);
return SqlFieldSchema.createMapSchema(valueType, isNullable, isOptional);
default:
String msg = String.format("Field Type %s is not supported", fieldSchema.getType());
LOG.error(msg);
throw new SamzaException(msg);
}
}

private SqlFieldSchema getSqlTypeFromUnionTypes(List<Schema> types) {
private SqlFieldSchema getSqlTypeFromUnionTypes(List<Schema> 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) {
if (types.get(0).getType() == Schema.Type.NULL) {
return convertField(types.get(1));
return convertField(types.get(1), true, isOptional);
} else if ((types.get(1).getType() == Schema.Type.NULL)) {
return convertField(types.get(0));
return convertField(types.get(0), true, isOptional);
}
}

return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.ANY);
return SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.ANY, isNullable, isOptional);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -152,25 +152,29 @@ public RelRoot plan(String query) {
}
}

public static RelDataType getSourceRelSchema(RelSchemaProvider relSchemaProvider,
RelSchemaConverter relSchemaConverter) {
// If the source part is the last one, then fetch the schema corresponding to the stream and register.
public static SqlSchema getSourceSqlSchema(RelSchemaProvider relSchemaProvider) {
SqlSchema sqlSchema = relSchemaProvider.getSqlSchema();

List<String> fieldNames = new ArrayList<>();
List<SqlFieldSchema> fieldTypes = new ArrayList<>();
if (!sqlSchema.containsField(SamzaSqlRelMessage.KEY_NAME)) {
fieldNames.add(SamzaSqlRelMessage.KEY_NAME);
fieldTypes.add(SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.ANY));
// Key is a nullable and optional field. It is defaulted to null in SamzaSqlRelMessage.
fieldTypes.add(SqlFieldSchema.createPrimitiveSchema(SamzaSqlFieldType.ANY, true, true));
}

fieldNames.addAll(
sqlSchema.getFields().stream().map(SqlSchema.SqlField::getFieldName).collect(Collectors.toList()));
fieldTypes.addAll(
sqlSchema.getFields().stream().map(SqlSchema.SqlField::getFieldSchema).collect(Collectors.toList()));

SqlSchema newSchema = new SqlSchema(fieldNames, fieldTypes);
return relSchemaConverter.convertToRelSchema(newSchema);
return new SqlSchema(fieldNames, fieldTypes);
}

public static RelDataType getSourceRelSchema(RelSchemaProvider relSchemaProvider,
RelSchemaConverter relSchemaConverter) {
// If the source part is the last one, then fetch the schema corresponding to the stream and register.
return relSchemaConverter.convertToRelSchema(getSourceSqlSchema(relSchemaProvider));
}

private static Table createTableFromRelSchema(RelDataType relationalSchema) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import org.apache.calcite.rel.type.RelDataTypeSystem;
import org.apache.calcite.rel.type.RelRecordType;
import org.apache.calcite.sql.type.ArraySqlType;
import org.apache.calcite.sql.type.MapSqlType;
import org.apache.calcite.sql.type.SqlTypeFactoryImpl;
import org.apache.calcite.sql.type.SqlTypeName;
import org.apache.samza.SamzaException;
Expand Down Expand Up @@ -73,30 +74,29 @@ private RelDataType getRelDataType(SqlFieldSchema fieldSchema) {
switch (fieldSchema.getFieldType()) {
case ARRAY:
RelDataType elementType = getRelDataType(fieldSchema.getElementSchema());
return new ArraySqlType(elementType, true);
return new ArraySqlType(elementType, fieldSchema.isNullable());
case BOOLEAN:
return createTypeWithNullability(createSqlType(SqlTypeName.BOOLEAN), true);
return createTypeWithNullability(createSqlType(SqlTypeName.BOOLEAN), fieldSchema.isNullable());
case DOUBLE:
return createTypeWithNullability(createSqlType(SqlTypeName.DOUBLE), true);
return createTypeWithNullability(createSqlType(SqlTypeName.DOUBLE), fieldSchema.isNullable());
case FLOAT:
return createTypeWithNullability(createSqlType(SqlTypeName.FLOAT), true);
return createTypeWithNullability(createSqlType(SqlTypeName.FLOAT), fieldSchema.isNullable());
case STRING:
return createTypeWithNullability(createSqlType(SqlTypeName.VARCHAR), true);
return createTypeWithNullability(createSqlType(SqlTypeName.VARCHAR), fieldSchema.isNullable());
case BYTES:
return createTypeWithNullability(createSqlType(SqlTypeName.VARBINARY), true);
return createTypeWithNullability(createSqlType(SqlTypeName.VARBINARY), fieldSchema.isNullable());
case INT16:
case INT32:
return createTypeWithNullability(createSqlType(SqlTypeName.INTEGER), true);
return createTypeWithNullability(createSqlType(SqlTypeName.INTEGER), fieldSchema.isNullable());
case INT64:
return createTypeWithNullability(createSqlType(SqlTypeName.BIGINT), true);
return createTypeWithNullability(createSqlType(SqlTypeName.BIGINT), fieldSchema.isNullable());
case ROW:
case ANY:
// TODO Calcite execution engine doesn't support record type yet.
return createTypeWithNullability(createSqlType(SqlTypeName.ANY), true);
return createTypeWithNullability(createSqlType(SqlTypeName.ANY), fieldSchema.isNullable());
case MAP:
RelDataType valueType = getRelDataType(fieldSchema.getValueScehma());
return super.createMapType(createTypeWithNullability(createSqlType(SqlTypeName.VARCHAR), true),
createTypeWithNullability(valueType, true));
return new MapSqlType(createSqlType(SqlTypeName.VARCHAR), valueType, fieldSchema.isNullable());
default:
String msg = String.format("Field Type %s is not supported", fieldSchema.getFieldType());
LOG.error(msg);
Expand Down
Loading