[ 
https://issues.apache.org/jira/browse/BEAM-3238?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=16270113#comment-16270113
 ] 

ASF GitHub Bot commented on BEAM-3238:
--------------------------------------

xumingming closed pull request #4168: [BEAM-3238][SQL] Add 
BeamRecordSqlTypeBuilder
URL: https://github.com/apache/beam/pull/4168
 
 
   

This is a PR merged from a forked repository.
As GitHub hides the original diff on merge, it is displayed below for
the sake of provenance:

As this is a foreign pull request (from a fork), the diff is supplied
below (as it won't show otherwise due to GitHub magic):

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 982494ad2e5..a6b23b6310b 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,14 @@
  */
 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;
 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 +51,39 @@
  *
  */
 public class BeamRecordSqlType extends BeamRecordType {
-  private static final Map<Integer, Class> 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<Integer, Class> JAVA_CLASSES = ImmutableMap
+      .<Integer, Class>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<Integer, Coder> CODERS = ImmutableMap
+      .<Integer, Coder>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<Integer> fieldTypes;
 
@@ -84,54 +98,24 @@ private BeamRecordSqlType(List<String> fieldsName, 
List<Integer> fieldTypes
   }
 
   public static BeamRecordSqlType create(List<String> fieldNames,
-      List<Integer> fieldTypes) {
+                                         List<Integer> fieldTypes) {
     if (fieldNames.size() != fieldTypes.size()) {
       throw new IllegalStateException("the sizes of 'dataType' and 
'fieldTypes' must match.");
     }
+
     List<Coder> 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:
+      Integer fieldType = fieldTypes.get(idx);
+
+      if (!CODERS.containsKey(fieldType)) {
         throw new UnsupportedOperationException(
-            "Data type: " + fieldTypes.get(idx) + " not supported yet!");
+            "Data type: " + fieldType + " not supported yet!");
       }
+
+      fieldCoders.add(CODERS.get(fieldType));
     }
+
     return new BeamRecordSqlType(fieldNames, fieldTypes, fieldCoders);
   }
 
@@ -142,7 +126,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 +143,7 @@ public void validateValueType(int index, Object fieldValue) 
throws IllegalArgume
     return Collections.unmodifiableList(fieldTypes);
   }
 
-  public Integer getFieldTypeByIndex(int index){
+  public Integer getFieldTypeByIndex(int index) {
     return fieldTypes.get(index);
   }
 
@@ -183,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<String> fieldNames;
+    private ImmutableList.Builder<Integer> 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 00000000000..78ff221e0d0
--- /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<Integer> 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<String> 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<String> 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<String> expectedNames = ImmutableList.<String>builder()
+        .addAll(NAMES)
+        .addAll(MORE_NAMES)
+        .build();
+
+    List<Integer> expectedTypes = ImmutableList.<Integer>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 33952692855..5997540099c 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<Class, Integer> 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<Class, Integer> JAVA_CLASS_TO_SQL_TYPE = 
ImmutableMap
+      .<Class, Integer> 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<BeamRecord> 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 @@
               9223372036854775807L
           )
           .buildIOReader(pipeline)
-          .setCoder(type.getRecordCoder());
+          .setCoder(RECORD_SQL_TYPE.getRecordCoder());
     } catch (Exception e) {
       throw new RuntimeException(e);
     }


 

----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on GitHub and use the
URL above to go to the specific comment.
 
For queries about this service, please contact Infrastructure at:
[email protected]


> [SQL] Add builder to BeamRecordSqlType
> --------------------------------------
>
>                 Key: BEAM-3238
>                 URL: https://issues.apache.org/jira/browse/BEAM-3238
>             Project: Beam
>          Issue Type: Improvement
>          Components: dsl-sql
>            Reporter: Anton Kedin
>            Assignee: Anton Kedin
>
> Currently it's hard to match field names with types when constructing a 
> BeamRecordSqlType, like 
> [here|https://github.com/apache/beam/blob/39e66e953b0f8e16435acb038cad364acf2b3a57/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/integrationtest/BeamSqlBuiltinFunctionsIntegrationTestBase.java#L64-L71]:
> {code:java}
> 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)
> );
> {code}
> It would be much more readable to have a builder, along these lines:
> {code:java}
> BeamRecordSqlType.builder()
>   .withField("f_int", Types.INTEGER)
>   .withStringField("f_str")
>   .build();
> {code}



--
This message was sent by Atlassian JIRA
(v6.4.14#64029)

Reply via email to