From 77769aa46dd9a972526f805f6c476c4039121b29 Mon Sep 17 00:00:00 2001 From: Anton Kedin Date: Wed, 22 Nov 2017 16:04:05 -0800 Subject: [PATCH 1/2] [BEAM-3238][SQL] Move Coders to map in BeamRecordSqlType Improve readability. All coders are supposed to be thread safe and currently are backed by the static instances. --- .../sdk/extensions/sql/BeamRecordSqlType.java | 111 ++++++++---------- 1 file changed, 47 insertions(+), 64 deletions(-) diff --git a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/BeamRecordSqlType.java b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/BeamRecordSqlType.java index 982494ad2e5a..1784ec14dbca 100644 --- a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/BeamRecordSqlType.java +++ b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/BeamRecordSqlType.java @@ -17,13 +17,13 @@ */ package org.apache.beam.sdk.extensions.sql; +import com.google.common.collect.ImmutableMap; import java.math.BigDecimal; import java.sql.Types; import java.util.ArrayList; import java.util.Collections; import java.util.Date; import java.util.GregorianCalendar; -import java.util.HashMap; import java.util.List; import java.util.Map; import org.apache.beam.sdk.coders.BigDecimalCoder; @@ -50,26 +50,39 @@ * */ public class BeamRecordSqlType extends BeamRecordType { - private static final Map SQL_TYPE_TO_JAVA_CLASS = new HashMap<>(); - static { - SQL_TYPE_TO_JAVA_CLASS.put(Types.TINYINT, Byte.class); - SQL_TYPE_TO_JAVA_CLASS.put(Types.SMALLINT, Short.class); - SQL_TYPE_TO_JAVA_CLASS.put(Types.INTEGER, Integer.class); - SQL_TYPE_TO_JAVA_CLASS.put(Types.BIGINT, Long.class); - SQL_TYPE_TO_JAVA_CLASS.put(Types.FLOAT, Float.class); - SQL_TYPE_TO_JAVA_CLASS.put(Types.DOUBLE, Double.class); - SQL_TYPE_TO_JAVA_CLASS.put(Types.DECIMAL, BigDecimal.class); - - SQL_TYPE_TO_JAVA_CLASS.put(Types.BOOLEAN, Boolean.class); - - SQL_TYPE_TO_JAVA_CLASS.put(Types.CHAR, String.class); - SQL_TYPE_TO_JAVA_CLASS.put(Types.VARCHAR, String.class); - - SQL_TYPE_TO_JAVA_CLASS.put(Types.TIME, GregorianCalendar.class); - - SQL_TYPE_TO_JAVA_CLASS.put(Types.DATE, Date.class); - SQL_TYPE_TO_JAVA_CLASS.put(Types.TIMESTAMP, Date.class); - } + private static final Map JAVA_CLASSES = ImmutableMap + .builder() + .put(Types.TINYINT, Byte.class) + .put(Types.SMALLINT, Short.class) + .put(Types.INTEGER, Integer.class) + .put(Types.BIGINT, Long.class) + .put(Types.FLOAT, Float.class) + .put(Types.DOUBLE, Double.class) + .put(Types.DECIMAL, BigDecimal.class) + .put(Types.BOOLEAN, Boolean.class) + .put(Types.CHAR, String.class) + .put(Types.VARCHAR, String.class) + .put(Types.TIME, GregorianCalendar.class) + .put(Types.DATE, Date.class) + .put(Types.TIMESTAMP, Date.class) + .build(); + + private static final Map CODERS = ImmutableMap + .builder() + .put(Types.TINYINT, ByteCoder.of()) + .put(Types.SMALLINT, ShortCoder.of()) + .put(Types.INTEGER, BigEndianIntegerCoder.of()) + .put(Types.BIGINT, BigEndianLongCoder.of()) + .put(Types.FLOAT, FloatCoder.of()) + .put(Types.DOUBLE, DoubleCoder.of()) + .put(Types.DECIMAL, BigDecimalCoder.of()) + .put(Types.BOOLEAN, BooleanCoder.of()) + .put(Types.CHAR, StringUtf8Coder.of()) + .put(Types.VARCHAR, StringUtf8Coder.of()) + .put(Types.TIME, TimeCoder.of()) + .put(Types.DATE, DateCoder.of()) + .put(Types.TIMESTAMP, DateCoder.of()) + .build(); public List fieldTypes; @@ -84,54 +97,24 @@ private BeamRecordSqlType(List fieldsName, List fieldTypes } public static BeamRecordSqlType create(List fieldNames, - List fieldTypes) { + List fieldTypes) { if (fieldNames.size() != fieldTypes.size()) { throw new IllegalStateException("the sizes of 'dataType' and 'fieldTypes' must match."); } + List fieldCoders = new ArrayList<>(fieldTypes.size()); + for (int idx = 0; idx < fieldTypes.size(); ++idx) { - switch (fieldTypes.get(idx)) { - case Types.INTEGER: - fieldCoders.add(BigEndianIntegerCoder.of()); - break; - case Types.SMALLINT: - fieldCoders.add(ShortCoder.of()); - break; - case Types.TINYINT: - fieldCoders.add(ByteCoder.of()); - break; - case Types.DOUBLE: - fieldCoders.add(DoubleCoder.of()); - break; - case Types.FLOAT: - fieldCoders.add(FloatCoder.of()); - break; - case Types.DECIMAL: - fieldCoders.add(BigDecimalCoder.of()); - break; - case Types.BIGINT: - fieldCoders.add(BigEndianLongCoder.of()); - break; - case Types.VARCHAR: - case Types.CHAR: - fieldCoders.add(StringUtf8Coder.of()); - break; - case Types.TIME: - fieldCoders.add(TimeCoder.of()); - break; - case Types.DATE: - case Types.TIMESTAMP: - fieldCoders.add(DateCoder.of()); - break; - case Types.BOOLEAN: - fieldCoders.add(BooleanCoder.of()); - break; - - default: - throw new UnsupportedOperationException( - "Data type: " + fieldTypes.get(idx) + " not supported yet!"); + Integer fieldType = fieldTypes.get(idx); + + if (!CODERS.containsKey(fieldType)) { + throw new UnsupportedOperationException( + "Data type: " + fieldType + " not supported yet!"); } + + fieldCoders.add(CODERS.get(fieldType)); } + return new BeamRecordSqlType(fieldNames, fieldTypes, fieldCoders); } @@ -142,7 +125,7 @@ public void validateValueType(int index, Object fieldValue) throws IllegalArgume } int fieldType = fieldTypes.get(index); - Class javaClazz = SQL_TYPE_TO_JAVA_CLASS.get(fieldType); + Class javaClazz = JAVA_CLASSES.get(fieldType); if (javaClazz == null) { throw new IllegalArgumentException("Data type: " + fieldType + " not supported yet!"); } @@ -159,7 +142,7 @@ public List getFieldTypes() { return Collections.unmodifiableList(fieldTypes); } - public Integer getFieldTypeByIndex(int index){ + public Integer getFieldTypeByIndex(int index) { return fieldTypes.get(index); } From 335d4862871ee82fbb19abcc1aaec83ad8ac1310 Mon Sep 17 00:00:00 2001 From: Anton Kedin Date: Wed, 22 Nov 2017 19:17:21 -0800 Subject: [PATCH 2/2] [BEAM-3238][SQL] Add BeamRecordSqlTypeBuilder To improve readability of creating BeamRecordSqlTypes. --- .../sdk/extensions/sql/BeamRecordSqlType.java | 85 ++++++++++++- .../extensions/sql/BeamRecordSqlTypeTest.java | 115 ++++++++++++++++++ ...qlBuiltinFunctionsIntegrationTestBase.java | 56 +++++---- 3 files changed, 229 insertions(+), 27 deletions(-) create mode 100644 sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamRecordSqlTypeTest.java diff --git a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/BeamRecordSqlType.java b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/BeamRecordSqlType.java index 1784ec14dbca..a6b23b6310b9 100644 --- a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/BeamRecordSqlType.java +++ b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/BeamRecordSqlType.java @@ -17,6 +17,7 @@ */ package org.apache.beam.sdk.extensions.sql; +import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableMap; import java.math.BigDecimal; import java.sql.Types; @@ -108,8 +109,8 @@ public static BeamRecordSqlType create(List fieldNames, Integer fieldType = fieldTypes.get(idx); if (!CODERS.containsKey(fieldType)) { - throw new UnsupportedOperationException( - "Data type: " + fieldType + " not supported yet!"); + throw new UnsupportedOperationException( + "Data type: " + fieldType + " not supported yet!"); } fieldCoders.add(CODERS.get(fieldType)); @@ -166,4 +167,84 @@ public String toString() { return "BeamRecordSqlType [fieldNames=" + getFieldNames() + ", fieldTypes=" + fieldTypes + "]"; } + + public static Builder builder() { + return new Builder(); + } + + /** + * Builder class to construct {@link BeamRecordSqlType}. + */ + public static class Builder { + + private ImmutableList.Builder fieldNames; + private ImmutableList.Builder fieldTypes; + + public Builder withField(String fieldName, Integer fieldType) { + fieldNames.add(fieldName); + fieldTypes.add(fieldType); + return this; + } + + public Builder withTinyIntField(String fieldName) { + return withField(fieldName, Types.TINYINT); + } + + public Builder withSmallIntField(String fieldName) { + return withField(fieldName, Types.SMALLINT); + } + + public Builder withIntegerField(String fieldName) { + return withField(fieldName, Types.INTEGER); + } + + public Builder withBigIntField(String fieldName) { + return withField(fieldName, Types.BIGINT); + } + + public Builder withFloatField(String fieldName) { + return withField(fieldName, Types.FLOAT); + } + + public Builder withDoubleField(String fieldName) { + return withField(fieldName, Types.DOUBLE); + } + + public Builder withDecimalField(String fieldName) { + return withField(fieldName, Types.DECIMAL); + } + + public Builder withBooleanField(String fieldName) { + return withField(fieldName, Types.BOOLEAN); + } + + public Builder withCharField(String fieldName) { + return withField(fieldName, Types.CHAR); + } + + public Builder withVarcharField(String fieldName) { + return withField(fieldName, Types.VARCHAR); + } + + public Builder withTimeField(String fieldName) { + return withField(fieldName, Types.TIME); + } + + public Builder withDateField(String fieldName) { + return withField(fieldName, Types.DATE); + } + + public Builder withTimestampField(String fieldName) { + return withField(fieldName, Types.TIMESTAMP); + } + + private Builder() { + this.fieldNames = ImmutableList.builder(); + this.fieldTypes = ImmutableList.builder(); + } + + public BeamRecordSqlType build() { + return create(fieldNames.build(), fieldTypes.build()); + } + } } diff --git a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamRecordSqlTypeTest.java b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamRecordSqlTypeTest.java new file mode 100644 index 000000000000..78ff221e0d0c --- /dev/null +++ b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamRecordSqlTypeTest.java @@ -0,0 +1,115 @@ +/* + * 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.beam.sdk.extensions.sql; + +import static org.junit.Assert.assertEquals; + +import com.google.common.collect.ImmutableList; +import java.sql.Types; +import java.util.List; +import org.junit.Test; + +/** + * Unit tests for {@link BeamRecordSqlType}. + */ +public class BeamRecordSqlTypeTest { + + private static final List TYPES = ImmutableList.of( + Types.TINYINT, + Types.SMALLINT, + Types.INTEGER, + Types.BIGINT, + Types.FLOAT, + Types.DOUBLE, + Types.DECIMAL, + Types.BOOLEAN, + Types.CHAR, + Types.VARCHAR, + Types.TIME, + Types.DATE, + Types.TIMESTAMP); + + private static final List NAMES = ImmutableList.of( + "TINYINT_FIELD", + "SMALLINT_FIELD", + "INTEGER_FIELD", + "BIGINT_FIELD", + "FLOAT_FIELD", + "DOUBLE_FIELD", + "DECIMAL_FIELD", + "BOOLEAN_FIELD", + "CHAR_FIELD", + "VARCHAR_FIELD", + "TIME_FIELD", + "DATE_FIELD", + "TIMESTAMP_FIELD"); + + private static final List MORE_NAMES = ImmutableList.of( + "ANOTHER_TINYINT_FIELD", + "ANOTHER_SMALLINT_FIELD", + "ANOTHER_INTEGER_FIELD", + "ANOTHER_BIGINT_FIELD", + "ANOTHER_FLOAT_FIELD", + "ANOTHER_DOUBLE_FIELD", + "ANOTHER_DECIMAL_FIELD", + "ANOTHER_BOOLEAN_FIELD", + "ANOTHER_CHAR_FIELD", + "ANOTHER_VARCHAR_FIELD", + "ANOTHER_TIME_FIELD", + "ANOTHER_DATE_FIELD", + "ANOTHER_TIMESTAMP_FIELD"); + + @Test + public void testBuildsWithCorrectFields() throws Exception { + BeamRecordSqlType.Builder recordTypeBuilder = BeamRecordSqlType.builder(); + + for (int i = 0; i < TYPES.size(); i++) { + recordTypeBuilder.withField(NAMES.get(i), TYPES.get(i)); + } + + recordTypeBuilder.withTinyIntField(MORE_NAMES.get(0)); + recordTypeBuilder.withSmallIntField(MORE_NAMES.get(1)); + recordTypeBuilder.withIntegerField(MORE_NAMES.get(2)); + recordTypeBuilder.withBigIntField(MORE_NAMES.get(3)); + recordTypeBuilder.withFloatField(MORE_NAMES.get(4)); + recordTypeBuilder.withDoubleField(MORE_NAMES.get(5)); + recordTypeBuilder.withDecimalField(MORE_NAMES.get(6)); + recordTypeBuilder.withBooleanField(MORE_NAMES.get(7)); + recordTypeBuilder.withCharField(MORE_NAMES.get(8)); + recordTypeBuilder.withVarcharField(MORE_NAMES.get(9)); + recordTypeBuilder.withTimeField(MORE_NAMES.get(10)); + recordTypeBuilder.withDateField(MORE_NAMES.get(11)); + recordTypeBuilder.withTimestampField(MORE_NAMES.get(12)); + + BeamRecordSqlType recordSqlType = recordTypeBuilder.build(); + + List expectedNames = ImmutableList.builder() + .addAll(NAMES) + .addAll(MORE_NAMES) + .build(); + + List expectedTypes = ImmutableList.builder() + .addAll(TYPES) + .addAll(TYPES) + .build(); + + assertEquals(expectedNames, recordSqlType.getFieldNames()); + assertEquals(expectedTypes, recordSqlType.getFieldTypes()); + } +} diff --git a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/integrationtest/BeamSqlBuiltinFunctionsIntegrationTestBase.java b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/integrationtest/BeamSqlBuiltinFunctionsIntegrationTestBase.java index 339526928553..5997540099c5 100644 --- a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/integrationtest/BeamSqlBuiltinFunctionsIntegrationTestBase.java +++ b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/integrationtest/BeamSqlBuiltinFunctionsIntegrationTestBase.java @@ -19,13 +19,12 @@ package org.apache.beam.sdk.extensions.sql.integrationtest; import com.google.common.base.Joiner; +import com.google.common.collect.ImmutableMap; import java.math.BigDecimal; import java.sql.Types; import java.text.SimpleDateFormat; import java.util.ArrayList; -import java.util.Arrays; import java.util.Date; -import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.TimeZone; @@ -44,35 +43,42 @@ * Base class for all built-in functions integration tests. */ public class BeamSqlBuiltinFunctionsIntegrationTestBase { - private static final Map JAVA_CLASS_TO_SQL_TYPE = new HashMap<>(); - static { - JAVA_CLASS_TO_SQL_TYPE.put(Byte.class, Types.TINYINT); - JAVA_CLASS_TO_SQL_TYPE.put(Short.class, Types.SMALLINT); - JAVA_CLASS_TO_SQL_TYPE.put(Integer.class, Types.INTEGER); - JAVA_CLASS_TO_SQL_TYPE.put(Long.class, Types.BIGINT); - JAVA_CLASS_TO_SQL_TYPE.put(Float.class, Types.FLOAT); - JAVA_CLASS_TO_SQL_TYPE.put(Double.class, Types.DOUBLE); - JAVA_CLASS_TO_SQL_TYPE.put(BigDecimal.class, Types.DECIMAL); - JAVA_CLASS_TO_SQL_TYPE.put(String.class, Types.VARCHAR); - JAVA_CLASS_TO_SQL_TYPE.put(Date.class, Types.DATE); - JAVA_CLASS_TO_SQL_TYPE.put(Boolean.class, Types.BOOLEAN); - } + private static final Map JAVA_CLASS_TO_SQL_TYPE = ImmutableMap + . builder() + .put(Byte.class, Types.TINYINT) + .put(Short.class, Types.SMALLINT) + .put(Integer.class, Types.INTEGER) + .put(Long.class, Types.BIGINT) + .put(Float.class, Types.FLOAT) + .put(Double.class, Types.DOUBLE) + .put(BigDecimal.class, Types.DECIMAL) + .put(String.class, Types.VARCHAR) + .put(Date.class, Types.DATE) + .put(Boolean.class, Types.BOOLEAN) + .build(); + + private static final BeamRecordSqlType RECORD_SQL_TYPE = BeamRecordSqlType.builder() + .withDateField("ts") + .withTinyIntField("c_tinyint") + .withSmallIntField("c_smallint") + .withIntegerField("c_integer") + .withBigIntField("c_bigint") + .withFloatField("c_float") + .withDoubleField("c_double") + .withDecimalField("c_decimal") + .withTinyIntField("c_tinyint_max") + .withSmallIntField("c_smallint_max") + .withIntegerField("c_integer_max") + .withBigIntField("c_bigint_max") + .build(); @Rule public final TestPipeline pipeline = TestPipeline.create(); protected PCollection getTestPCollection() { - BeamRecordSqlType type = BeamRecordSqlType.create( - Arrays.asList("ts", "c_tinyint", "c_smallint", - "c_integer", "c_bigint", "c_float", "c_double", "c_decimal", - "c_tinyint_max", "c_smallint_max", "c_integer_max", "c_bigint_max"), - Arrays.asList(Types.DATE, Types.TINYINT, Types.SMALLINT, - Types.INTEGER, Types.BIGINT, Types.FLOAT, Types.DOUBLE, Types.DECIMAL, - Types.TINYINT, Types.SMALLINT, Types.INTEGER, Types.BIGINT) - ); try { return MockedBoundedTable - .of(type) + .of(RECORD_SQL_TYPE) .addRows( parseDate("1986-02-15 11:35:26"), (byte) 1, @@ -88,7 +94,7 @@ protected PCollection getTestPCollection() { 9223372036854775807L ) .buildIOReader(pipeline) - .setCoder(type.getRecordCoder()); + .setCoder(RECORD_SQL_TYPE.getRecordCoder()); } catch (Exception e) { throw new RuntimeException(e); }