[ 
https://issues.apache.org/jira/browse/BEAM-4082?focusedWorklogId=107383&page=com.atlassian.jira.plugin.system.issuetabpanels:worklog-tabpanel#worklog-107383
 ]

ASF GitHub Bot logged work on BEAM-4082:
----------------------------------------

                Author: ASF GitHub Bot
            Created on: 30/May/18 21:24
            Start Date: 30/May/18 21:24
    Worklog Time Spent: 10m 
      Work Description: kennknowles closed pull request #5497: [BEAM-4082] 
Migrate from builders in RowSqlTypes to just using Schema builder directly
URL: https://github.com/apache/beam/pull/5497
 
 
   

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/core/src/main/java/org/apache/beam/sdk/schemas/Schema.java 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/Schema.java
index 94af2c7a809..69da5645e11 100644
--- a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/Schema.java
+++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/Schema.java
@@ -73,6 +73,16 @@ public Builder addField(Field field) {
       return this;
     }
 
+    public Builder addField(String name, FieldType type) {
+      fields.add(Field.of(name, type));
+      return this;
+    }
+
+    public Builder addNullableField(String name, FieldType type) {
+      fields.add(Field.nullable(name, type));
+      return this;
+    }
+
     public Builder addByteField(String name) {
       fields.add(Field.of(name, TypeName.BYTE.type()));
       return this;
diff --git 
a/sdks/java/core/src/test/java/org/apache/beam/sdk/coders/org/apache/beam/sdk/coders/RowCoderTest.java
 
b/sdks/java/core/src/test/java/org/apache/beam/sdk/coders/org/apache/beam/sdk/coders/RowCoderTest.java
index 2c62715b29c..675b9422eda 100644
--- 
a/sdks/java/core/src/test/java/org/apache/beam/sdk/coders/org/apache/beam/sdk/coders/RowCoderTest.java
+++ 
b/sdks/java/core/src/test/java/org/apache/beam/sdk/coders/org/apache/beam/sdk/coders/RowCoderTest.java
@@ -27,7 +27,6 @@
 import java.math.BigDecimal;
 import org.apache.beam.sdk.coders.RowCoder;
 import org.apache.beam.sdk.schemas.Schema;
-import org.apache.beam.sdk.schemas.Schema.Field;
 import org.apache.beam.sdk.schemas.Schema.FieldType;
 import org.apache.beam.sdk.schemas.Schema.TypeName;
 import org.apache.beam.sdk.values.Row;
@@ -114,7 +113,7 @@ public void testArrayOfArray() throws Exception {
     FieldType arrayType = TypeName.ARRAY.type()
         .withCollectionElementType(TypeName.ARRAY.type()
             .withCollectionElementType(TypeName.INT32.type()));
-    Schema schema = Schema.builder().addField(Field.of("f_array", 
arrayType)).build();
+    Schema schema = Schema.builder().addField("f_array", arrayType).build();
     Row row = Row.withSchema(schema).addArray(
         Lists.newArrayList(1, 2, 3, 4),
         Lists.newArrayList(5, 6, 7, 8),
diff --git 
a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/RowJsonDeserializerTest.java
 
b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/RowJsonDeserializerTest.java
index 56210d3d963..e486328b77f 100644
--- 
a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/RowJsonDeserializerTest.java
+++ 
b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/RowJsonDeserializerTest.java
@@ -356,7 +356,7 @@ public void testParsesNulls() throws Exception {
         Schema
             .builder()
             .addByteField("f_byte")
-            .addField(Schema.Field.nullable("f_string", FieldType.STRING))
+            .addNullableField("f_string", FieldType.STRING)
             .build();
 
     String rowString = "{\n"
@@ -562,7 +562,7 @@ private Schema schemaWithField(String fieldName, TypeName 
fieldType) {
     return
         Schema
             .builder()
-            .addField(Schema.Field.of(fieldName, fieldType.type()))
+            .addField(fieldName, fieldType.type())
             .build();
   }
 
diff --git 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/RowSqlTypes.java
 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/RowSqlTypes.java
deleted file mode 100644
index 72cf5efa5a8..00000000000
--- 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/RowSqlTypes.java
+++ /dev/null
@@ -1,178 +0,0 @@
-/*
- * 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 org.apache.beam.sdk.annotations.Experimental;
-import org.apache.beam.sdk.extensions.sql.impl.utils.CalciteUtils;
-import org.apache.beam.sdk.schemas.Schema;
-import org.apache.beam.sdk.schemas.Schema.Field;
-import org.apache.beam.sdk.schemas.Schema.FieldType;
-import org.apache.beam.sdk.schemas.Schema.TypeName;
-import org.apache.beam.sdk.values.Row;
-import org.apache.calcite.rel.type.RelDataType;
-import org.apache.calcite.sql.type.SqlTypeName;
-
-/**
- * Type builder for {@link Row} with SQL types.
- *
- * <p>Limited SQL types are supported now, visit <a
- * href="https://beam.apache.org/documentation/dsls/sql/#data-types";>data 
types</a> for more
- * details.
- *
- * <p>TODO: We should remove this class in favor of directly using 
Beam.Schema.Builder
- */
-@Experimental
-public class RowSqlTypes {
-  // The list of field type names used in SQL as Beam field types.
-  public static final FieldType TINY_INT = 
CalciteUtils.toFieldType(SqlTypeName.TINYINT);
-  public static final FieldType SMALL_INT = 
CalciteUtils.toFieldType(SqlTypeName.SMALLINT);
-  public static final FieldType INTEGER = 
CalciteUtils.toFieldType(SqlTypeName.INTEGER);
-  public static final FieldType BIG_INT = 
CalciteUtils.toFieldType(SqlTypeName.BIGINT);
-  public static final FieldType FLOAT = 
CalciteUtils.toFieldType(SqlTypeName.FLOAT);
-  public static final FieldType DOUBLE = 
CalciteUtils.toFieldType(SqlTypeName.DOUBLE);
-  public static final FieldType DECIMAL = 
CalciteUtils.toFieldType(SqlTypeName.DECIMAL);
-  public static final FieldType BOOLEAN = 
CalciteUtils.toFieldType(SqlTypeName.BOOLEAN);
-  public static final FieldType CHAR = 
CalciteUtils.toFieldType(SqlTypeName.CHAR);
-  public static final FieldType VARCHAR = 
CalciteUtils.toFieldType(SqlTypeName.VARCHAR);
-  public static final FieldType TIME = 
CalciteUtils.toFieldType(SqlTypeName.TIME);
-  public static final FieldType DATE = 
CalciteUtils.toFieldType(SqlTypeName.DATE);
-  public static final FieldType TIMESTAMP = 
CalciteUtils.toFieldType(SqlTypeName.TIMESTAMP);
-
-  public static Builder builder() {
-    return new Builder();
-  }
-
-  /** Builder class to construct {@link Schema}. */
-  public static class Builder {
-    Schema.Builder builder;
-
-    public Builder withTinyIntField(String fieldName) {
-      builder.addField(Field.of(fieldName, TINY_INT).withNullable(true));
-      return this;
-    }
-
-    public Builder withSmallIntField(String fieldName) {
-      builder.addField(Field.of(fieldName, SMALL_INT).withNullable(true));
-      return this;
-    }
-
-    public Builder withIntegerField(String fieldName) {
-      builder.addField(Field.of(fieldName, INTEGER).withNullable(true));
-      return this;
-    }
-
-    public Builder withBigIntField(String fieldName) {
-      builder.addField(Field.of(fieldName, BIG_INT).withNullable(true));
-      return this;
-    }
-
-    public Builder withFloatField(String fieldName) {
-      builder.addField(Field.of(fieldName, FLOAT).withNullable(true));
-      return this;
-    }
-
-    public Builder withDoubleField(String fieldName) {
-      builder.addField(Field.of(fieldName, DOUBLE).withNullable(true));
-      return this;
-    }
-
-    public Builder withDecimalField(String fieldName) {
-      builder.addField(Field.of(fieldName, DECIMAL).withNullable(true));
-      return this;
-    }
-
-    public Builder withBooleanField(String fieldName) {
-      builder.addField(Field.of(fieldName, BOOLEAN).withNullable(true));
-      return this;
-    }
-
-    public Builder withCharField(String fieldName) {
-      builder.addField(Field.of(fieldName, CHAR).withNullable(true));
-      return this;
-    }
-
-    public Builder withVarcharField(String fieldName) {
-      builder.addField(Field.of(fieldName, VARCHAR).withNullable(true));
-      return this;
-    }
-
-    public Builder withTimeField(String fieldName) {
-      builder.addField(Field.of(fieldName, TIME).withNullable(true));
-      return this;
-    }
-
-    public Builder withDateField(String fieldName) {
-      builder.addField(Field.of(fieldName, DATE).withNullable(true));
-      return this;
-    }
-
-    public Builder withTimestampField(String fieldName) {
-      builder.addField(Field.of(fieldName, TIMESTAMP).withNullable(true));
-      return this;
-    }
-
-    /** Adds an ARRAY field with elements of the given type. */
-    public Builder withArrayField(String fieldName, RelDataType relDataType) {
-      builder.addField(Field.of(fieldName, 
CalciteUtils.toArrayType(relDataType)));
-      return this;
-    }
-
-    /** Adds an ARRAY field with elements of the given type. */
-    public Builder withArrayField(String fieldName, SqlTypeName typeName) {
-      builder.addField(Field.of(fieldName, 
CalciteUtils.toArrayType(typeName)));
-      return this;
-    }
-
-    /** Adds a MAP field with elements of the given key/value type. */
-    public Builder withMapField(
-        String fieldName, RelDataType keyRelDataType, RelDataType 
valueRelDataType) {
-      builder.addField(
-          Field.of(fieldName, CalciteUtils.toMapType(keyRelDataType, 
valueRelDataType)));
-      return this;
-    }
-
-    /** Adds a MAP field with elements of the given key/value type. */
-    public Builder withMapField(
-        String fieldName, SqlTypeName keyTypeName, SqlTypeName valueTypeName) {
-      builder.addField(Field.of(fieldName, CalciteUtils.toMapType(keyTypeName, 
valueTypeName)));
-      return this;
-    }
-
-    /** Adds an ARRAY field with elements of {@code schema}. */
-    public Builder withArrayField(String fieldName, Schema schema) {
-      FieldType collectionElementType = 
FieldType.of(TypeName.ROW).withRowSchema(schema);
-      builder.addField(
-          Field.of(
-              fieldName, 
TypeName.ARRAY.type().withCollectionElementType(collectionElementType)));
-      return this;
-    }
-
-    public Builder withRowField(String fieldName, Schema schema) {
-      builder.addField(Field.nullable(fieldName, FieldType.row(schema)));
-      return this;
-    }
-
-    private Builder() {
-      this.builder = Schema.builder();
-    }
-
-    public Schema build() {
-      return builder.build();
-    }
-  }
-}
diff --git 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/example/BeamSqlExample.java
 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/example/BeamSqlExample.java
index 5af18a0e1ab..ab8bfce24ca 100644
--- 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/example/BeamSqlExample.java
+++ 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/example/BeamSqlExample.java
@@ -20,7 +20,6 @@
 import javax.annotation.Nullable;
 import org.apache.beam.sdk.Pipeline;
 import org.apache.beam.sdk.extensions.sql.BeamSql;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.options.PipelineOptions;
 import org.apache.beam.sdk.options.PipelineOptionsFactory;
 import org.apache.beam.sdk.schemas.Schema;
@@ -53,11 +52,7 @@ public static void main(String[] args) {
 
     //define the input row format
     Schema type =
-        RowSqlTypes.builder()
-            .withIntegerField("c1")
-            .withVarcharField("c2")
-            .withDoubleField("c3")
-            .build();
+        
Schema.builder().addInt32Field("c1").addStringField("c2").addDoubleField("c3").build();
 
     Row row1 = Row.withSchema(type).addValues(1, "row", 1.0).build();
     Row row2 = Row.withSchema(type).addValues(2, "row", 2.0).build();
diff --git 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/utils/CalciteUtils.java
 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/utils/CalciteUtils.java
index 9cd3e0c4cbe..74d3007e15c 100644
--- 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/utils/CalciteUtils.java
+++ 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/utils/CalciteUtils.java
@@ -67,6 +67,21 @@
   private static final BiMap<SqlTypeName, FieldType> 
CALCITE_TO_BEAM_TYPE_MAPPING =
       BEAM_TO_CALCITE_TYPE_MAPPING.inverse();
 
+  // The list of field type names used in SQL as Beam field types.
+  public static final FieldType TINY_INT = toFieldType(SqlTypeName.TINYINT);
+  public static final FieldType SMALL_INT = toFieldType(SqlTypeName.SMALLINT);
+  public static final FieldType INTEGER = toFieldType(SqlTypeName.INTEGER);
+  public static final FieldType BIG_INT = toFieldType(SqlTypeName.BIGINT);
+  public static final FieldType FLOAT = toFieldType(SqlTypeName.FLOAT);
+  public static final FieldType DOUBLE = toFieldType(SqlTypeName.DOUBLE);
+  public static final FieldType DECIMAL = toFieldType(SqlTypeName.DECIMAL);
+  public static final FieldType BOOLEAN = toFieldType(SqlTypeName.BOOLEAN);
+  public static final FieldType CHAR = toFieldType(SqlTypeName.CHAR);
+  public static final FieldType VARCHAR = toFieldType(SqlTypeName.VARCHAR);
+  public static final FieldType TIME = toFieldType(SqlTypeName.TIME);
+  public static final FieldType DATE = toFieldType(SqlTypeName.DATE);
+  public static final FieldType TIMESTAMP = toFieldType(SqlTypeName.TIMESTAMP);
+
   // Since there are multiple Calcite type that correspond to a single Beam 
type, this is the
   // default mapping.
   private static final Map<FieldType, SqlTypeName> 
BEAM_TO_CALCITE_DEFAULT_MAPPING =
@@ -97,7 +112,19 @@ public static SqlTypeName toSqlTypeName(FieldType type) {
   }
 
   public static FieldType toFieldType(SqlTypeName sqlTypeName) {
-    return CALCITE_TO_BEAM_TYPE_MAPPING.get(sqlTypeName).getTypeName().type();
+    switch (sqlTypeName) {
+      case MAP:
+      case MULTISET:
+      case ARRAY:
+      case ROW:
+        throw new IllegalArgumentException(
+            String.format(
+                "%s is a type constructor that takes parameters, not a type,"
+                    + "so it cannot be converted to a %s",
+                sqlTypeName, Schema.FieldType.class.getSimpleName()));
+      default:
+        return 
CALCITE_TO_BEAM_TYPE_MAPPING.get(sqlTypeName).getTypeName().type();
+    }
   }
 
   public static FieldType toFieldType(RelDataType calciteType) {
diff --git 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubJsonTableProvider.java
 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubJsonTableProvider.java
index 01f1370a46c..cb320e6d0b5 100644
--- 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubJsonTableProvider.java
+++ 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubJsonTableProvider.java
@@ -17,8 +17,8 @@
  */
 package org.apache.beam.sdk.extensions.sql.meta.provider.pubsub;
 
-import static org.apache.beam.sdk.extensions.sql.RowSqlTypes.TIMESTAMP;
-import static org.apache.beam.sdk.extensions.sql.RowSqlTypes.VARCHAR;
+import static 
org.apache.beam.sdk.extensions.sql.impl.utils.CalciteUtils.TIMESTAMP;
+import static 
org.apache.beam.sdk.extensions.sql.impl.utils.CalciteUtils.VARCHAR;
 import static 
org.apache.beam.sdk.extensions.sql.meta.provider.pubsub.PubsubMessageToRow.ATTRIBUTES_FIELD;
 import static 
org.apache.beam.sdk.extensions.sql.meta.provider.pubsub.PubsubMessageToRow.PAYLOAD_FIELD;
 import static 
org.apache.beam.sdk.extensions.sql.meta.provider.pubsub.PubsubMessageToRow.TIMESTAMP_FIELD;
diff --git 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubMessageToRow.java
 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubMessageToRow.java
index 4a76c792a39..1bb8dee3ba4 100644
--- 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubMessageToRow.java
+++ 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubMessageToRow.java
@@ -27,7 +27,6 @@
 import javax.annotation.Nullable;
 import org.apache.beam.sdk.annotations.Experimental;
 import org.apache.beam.sdk.annotations.Internal;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.io.gcp.pubsub.PubsubMessage;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.schemas.Schema.TypeName;
@@ -58,7 +57,7 @@
    * <p>Required to have exactly 3 top level fields at the moment:
    *
    * <ul>
-   *   <li>'event_timestamp' of type {@link RowSqlTypes#TIMESTAMP}
+   *   <li>'event_timestamp' of type {@link Schema.FieldType#DATETIME}
    *   <li>'attributes' of type {@link TypeName#MAP MAP&lt;VARCHAR,VARCHAR&gt;}
    *   <li>'payload' of type {@link TypeName#ROW ROW&lt;...&gt;}
    * </ul>
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliTest.java
index d3653b1dbba..d218b73e8b7 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliTest.java
@@ -17,9 +17,9 @@
  */
 package org.apache.beam.sdk.extensions.sql;
 
-import static org.apache.beam.sdk.extensions.sql.RowSqlTypes.BOOLEAN;
-import static org.apache.beam.sdk.extensions.sql.RowSqlTypes.INTEGER;
-import static org.apache.beam.sdk.extensions.sql.RowSqlTypes.VARCHAR;
+import static 
org.apache.beam.sdk.extensions.sql.impl.utils.CalciteUtils.BOOLEAN;
+import static 
org.apache.beam.sdk.extensions.sql.impl.utils.CalciteUtils.INTEGER;
+import static 
org.apache.beam.sdk.extensions.sql.impl.utils.CalciteUtils.VARCHAR;
 import static org.apache.beam.sdk.schemas.Schema.TypeName.ARRAY;
 import static org.apache.beam.sdk.schemas.Schema.TypeName.MAP;
 import static org.apache.beam.sdk.schemas.Schema.TypeName.ROW;
@@ -32,6 +32,7 @@
 import org.apache.beam.sdk.extensions.sql.meta.Table;
 import org.apache.beam.sdk.extensions.sql.meta.provider.text.TextTableProvider;
 import org.apache.beam.sdk.extensions.sql.meta.store.InMemoryMetaStore;
+import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.schemas.Schema.Field;
 import org.apache.calcite.tools.ValidationException;
 import org.junit.Test;
@@ -164,18 +165,18 @@ public void testExecute_createTableWithRowField() throws 
Exception {
                         "address",
                         ROW.type()
                             .withRowSchema(
-                                RowSqlTypes.builder()
-                                    .withVarcharField("street")
-                                    .withVarcharField("country")
+                                Schema.builder()
+                                    .addNullableField("street", 
Schema.FieldType.STRING)
+                                    .addNullableField("country", 
Schema.FieldType.STRING)
                                     .build()))
                     .withNullable(true),
                 Field.of(
                         "addressangular",
                         ROW.type()
                             .withRowSchema(
-                                RowSqlTypes.builder()
-                                    .withVarcharField("street")
-                                    .withVarcharField("country")
+                                Schema.builder()
+                                    .addNullableField("street", 
Schema.FieldType.STRING)
+                                    .addNullableField("country", 
Schema.FieldType.STRING)
                                     .build()))
                     .withNullable(true),
                 Field.of("isrobot", BOOLEAN).withNullable(true))
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationCovarianceTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationCovarianceTest.java
index d6e2a4d58b1..b3e83e17942 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationCovarianceTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationCovarianceTest.java
@@ -43,13 +43,13 @@
   @Before
   public void setUp() {
     Schema schema =
-        RowSqlTypes.builder()
-            .withDoubleField("f_double1")
-            .withDoubleField("f_double2")
-            .withDoubleField("f_double3")
-            .withIntegerField("f_int1")
-            .withIntegerField("f_int2")
-            .withIntegerField("f_int3")
+        Schema.builder()
+            .addDoubleField("f_double1")
+            .addDoubleField("f_double2")
+            .addDoubleField("f_double3")
+            .addInt32Field("f_int1")
+            .addInt32Field("f_int2")
+            .addInt32Field("f_int3")
             .build();
 
     List<Row> rowsInTableB =
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationTest.java
index f8fbe990166..e6322dd6d4f 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationTest.java
@@ -58,11 +58,11 @@
   @Before
   public void setUp() {
     Schema schemaInTableB =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withDoubleField("f_double")
-            .withIntegerField("f_int2")
-            .withDecimalField("f_decimal")
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addDoubleField("f_double")
+            .addInt32Field("f_int2")
+            .addDecimalField("f_decimal")
             .build();
 
     List<Row> rowsInTableB =
@@ -121,8 +121,7 @@ private void runAggregationWithoutWindow(PCollection<Row> 
input) throws Exceptio
 
     PCollection<Row> result = input.apply("testAggregationWithoutWindow", 
BeamSql.query(sql));
 
-    Schema resultType =
-        
RowSqlTypes.builder().withIntegerField("f_int2").withBigIntField("size").build();
+    Schema resultType = 
Schema.builder().addInt32Field("f_int2").addInt64Field("size").build();
 
     Row row = Row.withSchema(resultType).addValues(0, 4L).build();
 
@@ -166,35 +165,35 @@ private void runAggregationFunctions(PCollection<Row> 
input) throws Exception {
             .apply("testAggregationFunctions", BeamSql.query(sql));
 
     Schema resultType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int2")
-            .withBigIntField("size")
-            .withBigIntField("sum1")
-            .withBigIntField("avg1")
-            .withBigIntField("max1")
-            .withBigIntField("min1")
-            .withSmallIntField("sum2")
-            .withSmallIntField("avg2")
-            .withSmallIntField("max2")
-            .withSmallIntField("min2")
-            .withTinyIntField("sum3")
-            .withTinyIntField("avg3")
-            .withTinyIntField("max3")
-            .withTinyIntField("min3")
-            .withFloatField("sum4")
-            .withFloatField("avg4")
-            .withFloatField("max4")
-            .withFloatField("min4")
-            .withDoubleField("sum5")
-            .withDoubleField("avg5")
-            .withDoubleField("max5")
-            .withDoubleField("min5")
-            .withTimestampField("max6")
-            .withTimestampField("min6")
-            .withDoubleField("varpop1")
-            .withDoubleField("varsamp1")
-            .withIntegerField("varpop2")
-            .withIntegerField("varsamp2")
+        Schema.builder()
+            .addInt32Field("f_int2")
+            .addInt64Field("size")
+            .addInt64Field("sum1")
+            .addInt64Field("avg1")
+            .addInt64Field("max1")
+            .addInt64Field("min1")
+            .addInt16Field("sum2")
+            .addInt16Field("avg2")
+            .addInt16Field("max2")
+            .addInt16Field("min2")
+            .addByteField("sum3")
+            .addByteField("avg3")
+            .addByteField("max3")
+            .addByteField("min3")
+            .addFloatField("sum4")
+            .addFloatField("avg4")
+            .addFloatField("max4")
+            .addFloatField("min4")
+            .addDoubleField("sum5")
+            .addDoubleField("avg5")
+            .addDoubleField("max5")
+            .addDoubleField("min5")
+            .addDateTimeField("max6")
+            .addDateTimeField("min6")
+            .addDoubleField("varpop1")
+            .addDoubleField("varsamp1")
+            .addInt32Field("varpop2")
+            .addInt32Field("varsamp2")
             .build();
 
     Row row =
@@ -288,8 +287,7 @@ private void runDistinct(PCollection<Row> input) throws 
Exception {
 
     PCollection<Row> result = input.apply("testDistinct", BeamSql.query(sql));
 
-    Schema resultType =
-        
RowSqlTypes.builder().withIntegerField("f_int").withBigIntField("f_long").build();
+    Schema resultType = 
Schema.builder().addInt32Field("f_int").addInt64Field("f_long").build();
 
     List<Row> expectedRows =
         TestUtils.RowsBuilder.of(resultType)
@@ -328,10 +326,10 @@ private void runTumbleWindow(PCollection<Row> input) 
throws Exception {
             .apply("testTumbleWindow", BeamSql.query(sql));
 
     Schema resultType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int2")
-            .withBigIntField("size")
-            .withTimestampField("window_start")
+        Schema.builder()
+            .addInt32Field("f_int2")
+            .addInt64Field("size")
+            .addDateTimeField("window_start")
             .build();
 
     List<Row> expectedRows =
@@ -358,7 +356,7 @@ private void runTumbleWindow(PCollection<Row> input) throws 
Exception {
   @Category(UsesTestStream.class)
   public void testTriggeredTumble() throws Exception {
     Schema inputSchema =
-        
RowSqlTypes.builder().withIntegerField("f_int").withTimestampField("f_timestamp").build();
+        
Schema.builder().addInt32Field("f_int").addDateTimeField("f_timestamp").build();
 
     PCollection<Row> input =
         pipeline.apply(
@@ -384,7 +382,7 @@ public void testTriggeredTumble() throws Exception {
         "SELECT SUM(f_int) AS f_int_sum FROM PCOLLECTION"
             + " GROUP BY TUMBLE(f_timestamp, INTERVAL '1' HOUR)";
 
-    Schema outputSchema = 
RowSqlTypes.builder().withIntegerField("fn_int_sum").build();
+    Schema outputSchema = Schema.builder().addInt32Field("fn_int_sum").build();
 
     PCollection<Row> result =
         input
@@ -429,10 +427,10 @@ private void runHopWindow(PCollection<Row> input) throws 
Exception {
     PCollection<Row> result = input.apply("testHopWindow", BeamSql.query(sql));
 
     Schema resultType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int2")
-            .withBigIntField("size")
-            .withTimestampField("window_start")
+        Schema.builder()
+            .addInt32Field("f_int2")
+            .addInt64Field("size")
+            .addDateTimeField("window_start")
             .build();
 
     List<Row> expectedRows =
@@ -480,10 +478,10 @@ private void runSessionWindow(PCollection<Row> input) 
throws Exception {
             .apply("testSessionWindow", BeamSql.query(sql));
 
     Schema resultType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int2")
-            .withBigIntField("size")
-            .withTimestampField("window_start")
+        Schema.builder()
+            .addInt32Field("f_int2")
+            .addInt64Field("size")
+            .addDateTimeField("window_start")
             .build();
 
     List<Row> expectedRows =
@@ -555,10 +553,10 @@ public void testSupportsGlobalWindowWithCustomTrigger() 
throws Exception {
     DateTime startTime = new DateTime(2017, 1, 1, 0, 0, 0, 0);
 
     Schema type =
-        RowSqlTypes.builder()
-            .withIntegerField("f_intGroupingKey")
-            .withIntegerField("f_intValue")
-            .withTimestampField("f_timestamp")
+        Schema.builder()
+            .addInt32Field("f_intGroupingKey")
+            .addInt32Field("f_intValue")
+            .addDateTimeField("f_timestamp")
             .build();
 
     Object[] rows =
@@ -594,10 +592,10 @@ public void 
testSupportsNonGlobalWindowWithCustomTrigger() {
     DateTime startTime = new DateTime(2017, 1, 1, 0, 0, 0, 0);
 
     Schema type =
-        RowSqlTypes.builder()
-            .withIntegerField("f_intGroupingKey")
-            .withIntegerField("f_intValue")
-            .withTimestampField("f_timestamp")
+        Schema.builder()
+            .addInt32Field("f_intGroupingKey")
+            .addInt32Field("f_intValue")
+            .addDateTimeField("f_timestamp")
             .build();
 
     Object[] rows =
@@ -633,7 +631,7 @@ public void testSupportsNonGlobalWindowWithCustomTrigger() {
   }
 
   private List<Row> rowsWithSingleIntField(String fieldName, List<Integer> 
values) {
-    return 
TestUtils.rowsBuilderOf(RowSqlTypes.builder().withIntegerField(fieldName).build())
+    return 
TestUtils.rowsBuilderOf(Schema.builder().addInt32Field(fieldName).build())
         .addRows(values)
         .getRows();
   }
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationVarianceTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationVarianceTest.java
index 1d38688d162..2934090933d 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationVarianceTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationVarianceTest.java
@@ -43,10 +43,10 @@
   @Before
   public void setUp() {
     Schema schema =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withDoubleField("f_double")
-            .withIntegerField("f_int2")
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addDoubleField("f_double")
+            .addInt32Field("f_int2")
             .build();
 
     List<Row> rowsInTableB =
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslArrayTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslArrayTest.java
index 694d4d32c17..59d02d51e07 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslArrayTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslArrayTest.java
@@ -27,7 +27,6 @@
 import org.apache.beam.sdk.values.PCollectionTuple;
 import org.apache.beam.sdk.values.Row;
 import org.apache.beam.sdk.values.TupleTag;
-import org.apache.calcite.sql.type.SqlTypeName;
 import org.junit.Rule;
 import org.junit.Test;
 import org.junit.rules.ExpectedException;
@@ -36,9 +35,9 @@
 public class BeamSqlDslArrayTest {
 
   private static final Schema INPUT_SCHEMA =
-      RowSqlTypes.builder()
-          .withIntegerField("f_int")
-          .withArrayField("f_stringArr", SqlTypeName.VARCHAR)
+      Schema.builder()
+          .addInt32Field("f_int")
+          .addArrayField("f_stringArr", Schema.FieldType.STRING)
           .build();
 
   @Rule public final TestPipeline pipeline = TestPipeline.create();
@@ -49,9 +48,9 @@ public void testSelectArrayValue() {
     PCollection<Row> input = pCollectionOf2Elements();
 
     Schema resultType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withArrayField("f_arr", SqlTypeName.VARCHAR)
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addArrayField("f_arr", Schema.FieldType.STRING)
             .build();
 
     PCollection<Row> result =
@@ -71,9 +70,9 @@ public void testProjectArrayField() {
     PCollection<Row> input = pCollectionOf2Elements();
 
     Schema resultType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withArrayField("f_stringArr", SqlTypeName.VARCHAR)
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addArrayField("f_stringArr", Schema.FieldType.STRING)
             .build();
 
     PCollection<Row> result =
@@ -94,7 +93,7 @@ public void testProjectArrayField() {
   public void testAccessArrayElement() {
     PCollection<Row> input = pCollectionOf2Elements();
 
-    Schema resultType = 
RowSqlTypes.builder().withVarcharField("f_arrElem").build();
+    Schema resultType = Schema.builder().addStringField("f_arrElem").build();
 
     PCollection<Row> result =
         input.apply("sqlQuery", BeamSql.query("SELECT f_stringArr[0] FROM 
PCOLLECTION"));
@@ -115,7 +114,7 @@ public void testSingleElement() throws Exception {
         PBegin.in(pipeline)
             .apply("boundedInput1", 
Create.of(inputRow).withCoder(INPUT_SCHEMA.getRowCoder()));
 
-    Schema resultType = 
RowSqlTypes.builder().withVarcharField("f_arrElem").build();
+    Schema resultType = Schema.builder().addStringField("f_arrElem").build();
 
     PCollection<Row> result =
         input.apply("sqlQuery", BeamSql.query("SELECT ELEMENT(f_stringArr) 
FROM PCOLLECTION"));
@@ -129,7 +128,7 @@ public void testSingleElement() throws Exception {
   public void testCardinality() {
     PCollection<Row> input = pCollectionOf2Elements();
 
-    Schema resultType = 
RowSqlTypes.builder().withIntegerField("f_size").build();
+    Schema resultType = Schema.builder().addInt32Field("f_size").build();
 
     PCollection<Row> result =
         input.apply("sqlQuery", BeamSql.query("SELECT CARDINALITY(f_stringArr) 
FROM PCOLLECTION"));
@@ -151,7 +150,7 @@ public void testUnnestLiteral() {
     TupleTag<Row> mainTag = new TupleTag<Row>("main") {};
     PCollectionTuple inputTuple = PCollectionTuple.of(mainTag, input);
 
-    Schema resultType = 
RowSqlTypes.builder().withVarcharField("f_string").build();
+    Schema resultType = Schema.builder().addStringField("f_string").build();
 
     PCollection<Row> result =
         inputTuple.apply("sqlQuery", BeamSql.query("SELECT * FROM UNNEST 
(ARRAY ['a', 'b', 'c'])"));
@@ -174,7 +173,7 @@ public void testUnnestNamedLiteral() {
     TupleTag<Row> mainTag = new TupleTag<Row>("main") {};
     PCollectionTuple inputTuple = PCollectionTuple.of(mainTag, input);
 
-    Schema resultType = 
RowSqlTypes.builder().withVarcharField("f_string").build();
+    Schema resultType = Schema.builder().addStringField("f_string").build();
 
     PCollection<Row> result =
         inputTuple.apply(
@@ -209,8 +208,7 @@ public void testUnnestCrossJoin() {
     TupleTag<Row> mainTag = new TupleTag<Row>("main") {};
     PCollectionTuple inputTuple = PCollectionTuple.of(mainTag, input);
 
-    Schema resultType =
-        
RowSqlTypes.builder().withIntegerField("f_int").withVarcharField("f_string").build();
+    Schema resultType = 
Schema.builder().addInt32Field("f_int").addStringField("f_string").build();
 
     PCollection<Row> result =
         inputTuple.apply(
@@ -233,15 +231,17 @@ public void testUnnestCrossJoin() {
   @Test
   public void testSelectRowsFromArrayOfRows() {
     Schema elementSchema =
-        
RowSqlTypes.builder().withVarcharField("f_rowString").withIntegerField("f_rowInt").build();
+        
Schema.builder().addStringField("f_rowString").addInt32Field("f_rowInt").build();
 
     Schema resultSchema =
-        RowSqlTypes.builder().withArrayField("f_resultArray", 
elementSchema).build();
+        Schema.builder()
+            .addArrayField("f_resultArray", 
Schema.FieldType.row(elementSchema))
+            .build();
 
     Schema inputType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withArrayField("f_arrayOfRows", elementSchema)
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addArrayField("f_arrayOfRows", 
Schema.FieldType.row(elementSchema))
             .build();
 
     PCollection<Row> input =
@@ -290,14 +290,14 @@ public void testSelectRowsFromArrayOfRows() {
   @Test
   public void testSelectSingleRowFromArrayOfRows() {
     Schema elementSchema =
-        
RowSqlTypes.builder().withVarcharField("f_rowString").withIntegerField("f_rowInt").build();
+        
Schema.builder().addStringField("f_rowString").addInt32Field("f_rowInt").build();
 
     Schema resultSchema = elementSchema;
 
     Schema inputType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withArrayField("f_arrayOfRows", elementSchema)
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addArrayField("f_arrayOfRows", 
Schema.FieldType.row(elementSchema))
             .build();
 
     PCollection<Row> input =
@@ -336,14 +336,14 @@ public void testSelectSingleRowFromArrayOfRows() {
   @Test
   public void testSelectRowFieldFromArrayOfRows() {
     Schema elementSchema =
-        
RowSqlTypes.builder().withVarcharField("f_rowString").withIntegerField("f_rowInt").build();
+        
Schema.builder().addStringField("f_rowString").addInt32Field("f_rowInt").build();
 
-    Schema resultSchema = 
RowSqlTypes.builder().withVarcharField("f_stringField").build();
+    Schema resultSchema = 
Schema.builder().addStringField("f_stringField").build();
 
     Schema inputType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withArrayField("f_arrayOfRows", elementSchema)
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addArrayField("f_arrayOfRows", 
Schema.FieldType.row(elementSchema))
             .build();
 
     PCollection<Row> input =
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslBase.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslBase.java
index 02313c6126c..a8f5cbd6104 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslBase.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslBase.java
@@ -63,17 +63,17 @@
   @BeforeClass
   public static void prepareClass() throws ParseException {
     schemaInTableA =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withBigIntField("f_long")
-            .withSmallIntField("f_short")
-            .withTinyIntField("f_byte")
-            .withFloatField("f_float")
-            .withDoubleField("f_double")
-            .withVarcharField("f_string")
-            .withTimestampField("f_timestamp")
-            .withIntegerField("f_int2")
-            .withDecimalField("f_decimal")
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addInt64Field("f_long")
+            .addInt16Field("f_short")
+            .addByteField("f_byte")
+            .addFloatField("f_float")
+            .addDoubleField("f_double")
+            .addStringField("f_string")
+            .addDateTimeField("f_timestamp")
+            .addInt32Field("f_int2")
+            .addDecimalField("f_decimal")
             .build();
 
     rowsInTableA =
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslJoinTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslJoinTest.java
index 9e3ac1e664c..7c67c4fc09a 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslJoinTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslJoinTest.java
@@ -49,22 +49,22 @@
   @Rule public final TestPipeline pipeline = TestPipeline.create();
 
   private static final Schema SOURCE_ROW_TYPE =
-      RowSqlTypes.builder()
-          .withIntegerField("order_id")
-          .withIntegerField("site_id")
-          .withIntegerField("price")
+      Schema.builder()
+          .addNullableField("order_id", Schema.FieldType.INT32)
+          .addNullableField("site_id", Schema.FieldType.INT32)
+          .addNullableField("price", Schema.FieldType.INT32)
           .build();
 
   private static final RowCoder SOURCE_CODER = SOURCE_ROW_TYPE.getRowCoder();
 
   private static final Schema RESULT_ROW_TYPE =
-      RowSqlTypes.builder()
-          .withIntegerField("order_id")
-          .withIntegerField("site_id")
-          .withIntegerField("price")
-          .withIntegerField("order_id0")
-          .withIntegerField("site_id0")
-          .withIntegerField("price0")
+      Schema.builder()
+          .addNullableField("order_id", Schema.FieldType.INT32)
+          .addNullableField("site_id", Schema.FieldType.INT32)
+          .addNullableField("price", Schema.FieldType.INT32)
+          .addNullableField("order_id0", Schema.FieldType.INT32)
+          .addNullableField("site_id0", Schema.FieldType.INT32)
+          .addNullableField("price0", Schema.FieldType.INT32)
           .build();
 
   private static final RowCoder RESULT_CODER = RESULT_ROW_TYPE.getRowCoder();
@@ -297,11 +297,11 @@ public void 
testRejectsNonGlobalWindowsWithRepeatingTrigger() throws Exception {
     DateTime ts = new DateTime(2017, 1, 1, 1, 0, 0);
 
     return TestUtils.rowsBuilderOf(
-            RowSqlTypes.builder()
-                .withIntegerField("order_id")
-                .withIntegerField("price")
-                .withIntegerField("site_id")
-                .withTimestampField("timestamp")
+            Schema.builder()
+                .addInt32Field("order_id")
+                .addInt32Field("price")
+                .addInt32Field("site_id")
+                .addDateTimeField("timestamp")
                 .build())
         .addRows(
             1,
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslNestedRowsTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslNestedRowsTest.java
index 6ebcd2f0dad..285234a553f 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslNestedRowsTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslNestedRowsTest.java
@@ -25,7 +25,6 @@
 import org.apache.beam.sdk.values.PBegin;
 import org.apache.beam.sdk.values.PCollection;
 import org.apache.beam.sdk.values.Row;
-import org.apache.calcite.sql.type.SqlTypeName;
 import org.junit.Rule;
 import org.junit.Test;
 import org.junit.rules.ExpectedException;
@@ -39,22 +38,22 @@
   @Test
   public void testRowConstructorKeyword() {
     Schema nestedSchema =
-        RowSqlTypes.builder()
-            .withIntegerField("f_nestedInt")
-            .withVarcharField("f_nestedString")
-            .withIntegerField("f_nestedIntPlusOne")
+        Schema.builder()
+            .addInt32Field("f_nestedInt")
+            .addStringField("f_nestedString")
+            .addInt32Field("f_nestedIntPlusOne")
             .build();
 
     Schema resultSchema =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withIntegerField("f_int2")
-            .withVarcharField("f_varchar")
-            .withIntegerField("f_int3")
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addInt32Field("f_int2")
+            .addStringField("f_varchar")
+            .addInt32Field("f_int3")
             .build();
 
     Schema inputType =
-        RowSqlTypes.builder().withIntegerField("f_int").withRowField("f_row", 
nestedSchema).build();
+        Schema.builder().addInt32Field("f_int").addRowField("f_row", 
nestedSchema).build();
 
     PCollection<Row> input =
         PBegin.in(pipeline)
@@ -83,22 +82,22 @@ public void testRowConstructorKeyword() {
   public void testRowConstructorBraces() {
 
     Schema nestedSchema =
-        RowSqlTypes.builder()
-            .withIntegerField("f_nestedInt")
-            .withVarcharField("f_nestedString")
-            .withIntegerField("f_nestedIntPlusOne")
+        Schema.builder()
+            .addInt32Field("f_nestedInt")
+            .addStringField("f_nestedString")
+            .addInt32Field("f_nestedIntPlusOne")
             .build();
 
     Schema resultSchema =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withIntegerField("f_int2")
-            .withVarcharField("f_varchar")
-            .withIntegerField("f_int3")
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addInt32Field("f_int2")
+            .addStringField("f_varchar")
+            .addInt32Field("f_int3")
             .build();
 
     Schema inputType =
-        RowSqlTypes.builder().withIntegerField("f_int").withRowField("f_row", 
nestedSchema).build();
+        Schema.builder().addInt32Field("f_int").addRowField("f_row", 
nestedSchema).build();
 
     PCollection<Row> input =
         PBegin.in(pipeline)
@@ -127,19 +126,16 @@ public void testRowConstructorBraces() {
   public void testNestedRowFieldAccess() {
 
     Schema nestedSchema =
-        RowSqlTypes.builder()
-            .withIntegerField("f_nestedInt")
-            .withVarcharField("f_nestedString")
-            .withIntegerField("f_nestedIntPlusOne")
+        Schema.builder()
+            .addInt32Field("f_nestedInt")
+            .addStringField("f_nestedString")
+            .addInt32Field("f_nestedIntPlusOne")
             .build();
 
-    Schema resultSchema = 
RowSqlTypes.builder().withVarcharField("f_nestedString").build();
+    Schema resultSchema = 
Schema.builder().addStringField("f_nestedString").build();
 
     Schema inputType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withRowField("f_nestedRow", nestedSchema)
-            .build();
+        Schema.builder().addInt32Field("f_int").addRowField("f_nestedRow", 
nestedSchema).build();
 
     PCollection<Row> input =
         PBegin.in(pipeline)
@@ -174,21 +170,18 @@ public void testNestedRowFieldAccess() {
   public void testNestedRowArrayFieldAccess() {
 
     Schema resultSchema =
-        RowSqlTypes.builder().withArrayField("f_nestedArray", 
SqlTypeName.VARCHAR).build();
+        Schema.builder().addArrayField("f_nestedArray", 
Schema.FieldType.STRING).build();
 
     Schema nestedSchema =
-        RowSqlTypes.builder()
-            .withIntegerField("f_nestedInt")
-            .withVarcharField("f_nestedString")
-            .withIntegerField("f_nestedIntPlusOne")
-            .withArrayField("f_nestedArray", SqlTypeName.VARCHAR)
+        Schema.builder()
+            .addInt32Field("f_nestedInt")
+            .addStringField("f_nestedString")
+            .addInt32Field("f_nestedIntPlusOne")
+            .addArrayField("f_nestedArray", Schema.FieldType.STRING)
             .build();
 
     Schema inputType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withRowField("f_nestedRow", nestedSchema)
-            .build();
+        Schema.builder().addInt32Field("f_int").addRowField("f_nestedRow", 
nestedSchema).build();
 
     PCollection<Row> input =
         PBegin.in(pipeline)
@@ -228,22 +221,18 @@ public void testNestedRowArrayFieldAccess() {
   @Test
   public void testNestedRowArrayElementAccess() {
 
-    Schema resultSchema =
-        
RowSqlTypes.builder().withVarcharField("f_nestedArrayStringField").build();
+    Schema resultSchema = 
Schema.builder().addStringField("f_nestedArrayStringField").build();
 
     Schema nestedSchema =
-        RowSqlTypes.builder()
-            .withIntegerField("f_nestedInt")
-            .withVarcharField("f_nestedString")
-            .withIntegerField("f_nestedIntPlusOne")
-            .withArrayField("f_nestedArray", SqlTypeName.VARCHAR)
+        Schema.builder()
+            .addInt32Field("f_nestedInt")
+            .addStringField("f_nestedString")
+            .addInt32Field("f_nestedIntPlusOne")
+            .addArrayField("f_nestedArray", Schema.FieldType.STRING)
             .build();
 
     Schema inputType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withRowField("f_nestedRow", nestedSchema)
-            .build();
+        Schema.builder().addInt32Field("f_int").addRowField("f_nestedRow", 
nestedSchema).build();
 
     PCollection<Row> input =
         PBegin.in(pipeline)
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslProjectTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslProjectTest.java
index ecf8db4a0ca..15b583e6a6e 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslProjectTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslProjectTest.java
@@ -72,8 +72,7 @@ private void runPartialFields(PCollection<Row> input) throws 
Exception {
         PCollectionTuple.of(new TupleTag<>("TABLE_A"), input)
             .apply("testPartialFields", BeamSql.query(sql));
 
-    Schema resultType =
-        
RowSqlTypes.builder().withIntegerField("f_int").withBigIntField("f_long").build();
+    Schema resultType = 
Schema.builder().addInt32Field("f_int").addInt64Field("f_long").build();
 
     Row row = rowAtIndex(resultType, 0);
 
@@ -101,8 +100,7 @@ private void runPartialFieldsInMultipleRow(PCollection<Row> 
input) throws Except
         PCollectionTuple.of(new TupleTag<>("TABLE_A"), input)
             .apply("testPartialFieldsInMultipleRow", BeamSql.query(sql));
 
-    Schema resultType =
-        
RowSqlTypes.builder().withIntegerField("f_int").withBigIntField("f_long").build();
+    Schema resultType = 
Schema.builder().addInt32Field("f_int").addInt64Field("f_long").build();
 
     List<Row> expectedRows =
         IntStream.range(0, 4).mapToObj(i -> rowAtIndex(resultType, 
i)).collect(toList());
@@ -137,8 +135,7 @@ private void runPartialFieldsInRows(PCollection<Row> input) 
throws Exception {
         PCollectionTuple.of(new TupleTag<>("TABLE_A"), input)
             .apply("testPartialFieldsInRows", BeamSql.query(sql));
 
-    Schema resultType =
-        
RowSqlTypes.builder().withIntegerField("f_int").withBigIntField("f_long").build();
+    Schema resultType = 
Schema.builder().addInt32Field("f_int").addInt64Field("f_long").build();
 
     List<Row> expectedRows =
         IntStream.range(0, 4).mapToObj(i -> rowAtIndex(resultType, 
i)).collect(toList());
@@ -167,7 +164,7 @@ public void runLiteralField(PCollection<Row> input) throws 
Exception {
         PCollectionTuple.of(new TupleTag<>("TABLE_A"), input)
             .apply("testLiteralField", BeamSql.query(sql));
 
-    Schema resultType = 
RowSqlTypes.builder().withIntegerField("literal_field").build();
+    Schema resultType = 
Schema.builder().addInt32Field("literal_field").build();
 
     Row row = Row.withSchema(resultType).addValues(1).build();
 
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslUdfUdafTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslUdfUdafTest.java
index 7fe0004570a..dc99354ba62 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslUdfUdafTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslUdfUdafTest.java
@@ -34,8 +34,7 @@
   /** GROUP-BY with UDAF. */
   @Test
   public void testUdaf() throws Exception {
-    Schema resultType =
-        
RowSqlTypes.builder().withIntegerField("f_int2").withIntegerField("squaresum").build();
+    Schema resultType = 
Schema.builder().addInt32Field("f_int2").addInt32Field("squaresum").build();
 
     Row row = Row.withSchema(resultType).addValues(0, 30).build();
 
@@ -59,8 +58,7 @@ public void testUdaf() throws Exception {
   /** Test that an indirect subclass of a {@link CombineFn} works as a UDAF. 
BEAM-3777 */
   @Test
   public void testUdafMultiLevelDescendent() {
-    Schema resultType =
-        
RowSqlTypes.builder().withIntegerField("f_int2").withIntegerField("squaresum").build();
+    Schema resultType = 
Schema.builder().addInt32Field("f_int2").addInt32Field("squaresum").build();
 
     Row row = Row.withSchema(resultType).addValues(0, 354).build();
 
@@ -87,8 +85,7 @@ public void testRawCombineFnSubclass() {
     exceptions.expectMessage("CombineFn");
     pipeline.enableAbandonedNodeEnforcement(false);
 
-    Schema resultType =
-        
RowSqlTypes.builder().withIntegerField("f_int2").withIntegerField("squaresum").build();
+    Schema resultType = 
Schema.builder().addInt32Field("f_int2").addInt32Field("squaresum").build();
 
     Row row = Row.withSchema(resultType).addValues(0, 354).build();
 
@@ -102,8 +99,7 @@ public void testRawCombineFnSubclass() {
   /** test UDF. */
   @Test
   public void testUdf() throws Exception {
-    Schema resultType =
-        
RowSqlTypes.builder().withIntegerField("f_int").withIntegerField("cubicvalue").build();
+    Schema resultType = 
Schema.builder().addInt32Field("f_int").addInt32Field("cubicvalue").build();
     Row row = Row.withSchema(resultType).addValues(2, 8).build();
 
     String sql1 = "SELECT f_int, cubic1(f_int) as cubicvalue FROM PCOLLECTION 
WHERE f_int = 2";
@@ -124,7 +120,7 @@ public void testUdf() throws Exception {
             .apply("testUdf3", BeamSql.query(sql3).registerUdf("substr", 
UdfFnWithDefault.class));
 
     Schema subStrSchema =
-        
RowSqlTypes.builder().withIntegerField("f_int").withVarcharField("sub_string").build();
+        
Schema.builder().addInt32Field("f_int").addStringField("sub_string").build();
     Row subStrRow = Row.withSchema(subStrSchema).addValues(2, "s").build();
     PAssert.that(result3).containsInAnyOrder(subStrRow);
 
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlMapTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlMapTest.java
index 26940385581..810f80e9cb9 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlMapTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlMapTest.java
@@ -25,7 +25,6 @@
 import org.apache.beam.sdk.values.PBegin;
 import org.apache.beam.sdk.values.PCollection;
 import org.apache.beam.sdk.values.Row;
-import org.apache.calcite.sql.type.SqlTypeName;
 import org.junit.Rule;
 import org.junit.Test;
 import org.junit.rules.ExpectedException;
@@ -34,9 +33,9 @@
 public class BeamSqlMapTest {
 
   private static final Schema INPUT_ROW_TYPE =
-      RowSqlTypes.builder()
-          .withIntegerField("f_int")
-          .withMapField("f_intStringMap", SqlTypeName.VARCHAR, 
SqlTypeName.INTEGER)
+      Schema.builder()
+          .addInt32Field("f_int")
+          .addMapField("f_intStringMap", Schema.FieldType.STRING, 
Schema.FieldType.INT32)
           .build();
 
   @Rule public final TestPipeline pipeline = TestPipeline.create();
@@ -47,9 +46,10 @@ public void testSelectAll() {
     PCollection<Row> input = pCollectionOf2Elements();
 
     Schema resultType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withMapField("f_map", SqlTypeName.VARCHAR, SqlTypeName.INTEGER)
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addNullableField(
+                "f_map", Schema.FieldType.map(Schema.FieldType.STRING, 
Schema.FieldType.INT32))
             .build();
 
     PCollection<Row> result =
@@ -88,9 +88,9 @@ public void testSelectMapField() {
     PCollection<Row> input = pCollectionOf2Elements();
 
     Schema resultType =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withMapField("f_intStringMap", SqlTypeName.VARCHAR, 
SqlTypeName.INTEGER)
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addMapField("f_intStringMap", Schema.FieldType.STRING, 
Schema.FieldType.INT32)
             .build();
 
     PCollection<Row> result =
@@ -125,7 +125,8 @@ public void testSelectMapField() {
   public void testAccessMapElement() {
     PCollection<Row> input = pCollectionOf2Elements();
 
-    Schema resultType = 
RowSqlTypes.builder().withIntegerField("f_mapElem").build();
+    Schema resultType =
+        Schema.builder().addNullableField("f_mapElem", 
Schema.FieldType.INT32).build();
 
     PCollection<Row> result =
         input.apply("sqlQuery", BeamSql.query("SELECT f_intStringMap['key11'] 
FROM PCOLLECTION"));
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/interpreter/operator/BeamSqlDotExpressionTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/interpreter/operator/BeamSqlDotExpressionTest.java
index 1e95a28cb59..a6c783a2508 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/interpreter/operator/BeamSqlDotExpressionTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/interpreter/operator/BeamSqlDotExpressionTest.java
@@ -22,7 +22,6 @@
 import com.google.common.collect.ImmutableList;
 import com.google.common.collect.ImmutableMap;
 import java.util.List;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
 import org.apache.beam.sdk.values.Row;
@@ -41,8 +40,7 @@
 
   @Test
   public void testReturnsFieldValue() {
-    Schema schema =
-        
RowSqlTypes.builder().withVarcharField("f_string").withIntegerField("f_int").build();
+    Schema schema = 
Schema.builder().addStringField("f_string").addInt32Field("f_int").build();
 
     List<BeamSqlExpression> elements =
         ImmutableList.of(
@@ -58,8 +56,7 @@ public void testReturnsFieldValue() {
 
   @Test
   public void testThrowsForNonExistentField() {
-    Schema schema =
-        
RowSqlTypes.builder().withVarcharField("f_string").withIntegerField("f_int").build();
+    Schema schema = 
Schema.builder().addStringField("f_string").addInt32Field("f_int").build();
 
     List<BeamSqlExpression> elements =
         ImmutableList.of(
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/interpreter/operator/row/BeamSqlFieldAccessExpressionTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/interpreter/operator/row/BeamSqlFieldAccessExpressionTest.java
index d871dcedefc..fdc8a26df35 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/interpreter/operator/row/BeamSqlFieldAccessExpressionTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/interpreter/operator/row/BeamSqlFieldAccessExpressionTest.java
@@ -22,7 +22,6 @@
 import com.google.common.collect.ImmutableMap;
 import java.util.Arrays;
 import java.util.List;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import 
org.apache.beam.sdk.extensions.sql.impl.interpreter.operator.BeamSqlPrimitive;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
@@ -55,10 +54,10 @@ public void testAccessesFieldOfArray() {
   @Test
   public void testAccessesFieldOfRow() {
     Schema schema =
-        RowSqlTypes.builder()
-            .withVarcharField("f_string1")
-            .withVarcharField("f_string2")
-            .withVarcharField("f_string3")
+        Schema.builder()
+            .addStringField("f_string1")
+            .addStringField("f_string2")
+            .addStringField("f_string3")
             .build();
 
     BeamSqlPrimitive<Row> targetRow =
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/parser/BeamDDLNestedTypesTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/parser/BeamDDLNestedTypesTest.java
index fcba476de77..6fbd941f0f9 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/parser/BeamDDLNestedTypesTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/parser/BeamDDLNestedTypesTest.java
@@ -31,7 +31,6 @@
 import 
org.apache.beam.sdk.extensions.sql.utils.QuickCheckGenerators.AnyFieldType;
 import 
org.apache.beam.sdk.extensions.sql.utils.QuickCheckGenerators.PrimitiveTypes;
 import org.apache.beam.sdk.schemas.Schema;
-import org.apache.beam.sdk.schemas.Schema.Field;
 import org.apache.beam.sdk.schemas.Schema.FieldType;
 import org.apache.calcite.sql.SqlNode;
 import org.junit.runner.RunWith;
@@ -92,7 +91,7 @@ private Table executeCreateTableWith(String fieldType) {
   }
 
   private Schema newSimpleSchemaWith(FieldType fieldType) {
-    return Schema.builder().addField(Field.of("fieldname", 
fieldType).withNullable(true)).build();
+    return Schema.builder().addNullableField("fieldname", fieldType).build();
   }
 
   private String unparse(FieldType fieldType) {
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/parser/BeamDDLTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/parser/BeamDDLTest.java
index c66dab5bb1c..ca5409fecf3 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/parser/BeamDDLTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/parser/BeamDDLTest.java
@@ -25,8 +25,8 @@
 import com.alibaba.fastjson.JSONArray;
 import com.alibaba.fastjson.JSONObject;
 import java.util.stream.Stream;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.extensions.sql.impl.parser.impl.BeamSqlParserImpl;
+import org.apache.beam.sdk.extensions.sql.impl.utils.CalciteUtils;
 import org.apache.beam.sdk.extensions.sql.meta.Table;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.schemas.Schema.TypeName;
@@ -150,7 +150,7 @@ private static Table mockTable(
                     Schema.Field.of("id", TypeName.INT32.type())
                         .withNullable(true)
                         .withDescription("id"),
-                    Schema.Field.of("name", RowSqlTypes.VARCHAR)
+                    Schema.Field.of("name", CalciteUtils.VARCHAR)
                         .withNullable(true)
                         .withDescription("name"))
                 .collect(toSchema()))
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/schema/transform/BeamAggregationTransformTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/schema/transform/BeamAggregationTransformTest.java
index 11aee8a6e55..fd19285be44 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/schema/transform/BeamAggregationTransformTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/schema/transform/BeamAggregationTransformTest.java
@@ -24,7 +24,6 @@
 import org.apache.beam.sdk.coders.IterableCoder;
 import org.apache.beam.sdk.coders.KvCoder;
 import org.apache.beam.sdk.coders.RowCoder;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import 
org.apache.beam.sdk.extensions.sql.impl.transform.BeamAggregationTransforms;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.testing.PAssert;
@@ -356,39 +355,39 @@ private void prepareAggregationCalls() {
   private void prepareTypeAndCoder() {
     inRecordCoder = inputSchema.getRowCoder();
 
-    keyType = RowSqlTypes.builder().withIntegerField("f_int").build();
+    keyType = Schema.builder().addInt32Field("f_int").build();
 
     keyCoder = keyType.getRowCoder();
 
     aggPartType =
-        RowSqlTypes.builder()
-            .withBigIntField("count")
-            .withBigIntField("sum1")
-            .withBigIntField("avg1")
-            .withBigIntField("max1")
-            .withBigIntField("min1")
-            .withSmallIntField("sum2")
-            .withSmallIntField("avg2")
-            .withSmallIntField("max2")
-            .withSmallIntField("min2")
-            .withTinyIntField("sum3")
-            .withTinyIntField("avg3")
-            .withTinyIntField("max3")
-            .withTinyIntField("min3")
-            .withFloatField("sum4")
-            .withFloatField("avg4")
-            .withFloatField("max4")
-            .withFloatField("min4")
-            .withDoubleField("sum5")
-            .withDoubleField("avg5")
-            .withDoubleField("max5")
-            .withDoubleField("min5")
-            .withTimestampField("max7")
-            .withTimestampField("min7")
-            .withIntegerField("sum8")
-            .withIntegerField("avg8")
-            .withIntegerField("max8")
-            .withIntegerField("min8")
+        Schema.builder()
+            .addInt64Field("count")
+            .addInt64Field("sum1")
+            .addInt64Field("avg1")
+            .addInt64Field("max1")
+            .addInt64Field("min1")
+            .addInt16Field("sum2")
+            .addInt16Field("avg2")
+            .addInt16Field("max2")
+            .addInt16Field("min2")
+            .addByteField("sum3")
+            .addByteField("avg3")
+            .addByteField("max3")
+            .addByteField("min3")
+            .addFloatField("sum4")
+            .addFloatField("avg4")
+            .addFloatField("max4")
+            .addFloatField("min4")
+            .addDoubleField("sum5")
+            .addDoubleField("avg5")
+            .addDoubleField("max5")
+            .addDoubleField("min5")
+            .addDateTimeField("max7")
+            .addDateTimeField("min7")
+            .addInt32Field("sum8")
+            .addInt32Field("avg8")
+            .addInt32Field("max8")
+            .addInt32Field("min8")
             .build();
 
     aggCoder = aggPartType.getRowCoder();
@@ -447,35 +446,35 @@ private void prepareTypeAndCoder() {
 
   /** Row type of final output row. */
   private Schema prepareFinalSchema() {
-    return RowSqlTypes.builder()
-        .withIntegerField("f_int")
-        .withBigIntField("count")
-        .withBigIntField("sum1")
-        .withBigIntField("avg1")
-        .withBigIntField("max1")
-        .withBigIntField("min1")
-        .withSmallIntField("sum2")
-        .withSmallIntField("avg2")
-        .withSmallIntField("max2")
-        .withSmallIntField("min2")
-        .withTinyIntField("sum3")
-        .withTinyIntField("avg3")
-        .withTinyIntField("max3")
-        .withTinyIntField("min3")
-        .withFloatField("sum4")
-        .withFloatField("avg4")
-        .withFloatField("max4")
-        .withFloatField("min4")
-        .withDoubleField("sum5")
-        .withDoubleField("avg5")
-        .withDoubleField("max5")
-        .withDoubleField("min5")
-        .withTimestampField("max7")
-        .withTimestampField("min7")
-        .withIntegerField("sum8")
-        .withIntegerField("avg8")
-        .withIntegerField("max8")
-        .withIntegerField("min8")
+    return Schema.builder()
+        .addInt32Field("f_int")
+        .addInt64Field("count")
+        .addInt64Field("sum1")
+        .addInt64Field("avg1")
+        .addInt64Field("max1")
+        .addInt64Field("min1")
+        .addInt16Field("sum2")
+        .addInt16Field("avg2")
+        .addInt16Field("max2")
+        .addInt16Field("min2")
+        .addByteField("sum3")
+        .addByteField("avg3")
+        .addByteField("max3")
+        .addByteField("min3")
+        .addFloatField("sum4")
+        .addFloatField("avg4")
+        .addFloatField("max4")
+        .addFloatField("min4")
+        .addDoubleField("sum5")
+        .addDoubleField("avg5")
+        .addDoubleField("max5")
+        .addDoubleField("min5")
+        .addDateTimeField("max7")
+        .addDateTimeField("min7")
+        .addInt32Field("sum8")
+        .addInt32Field("avg8")
+        .addInt32Field("max8")
+        .addInt32Field("min8")
         .build();
   }
 
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/schema/transform/BeamTransformBaseTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/schema/transform/BeamTransformBaseTest.java
index 335863d5327..f6f2df20e10 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/schema/transform/BeamTransformBaseTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/schema/transform/BeamTransformBaseTest.java
@@ -18,7 +18,6 @@
 
 import java.text.ParseException;
 import java.util.List;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.extensions.sql.TestUtils;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.values.Row;
@@ -36,16 +35,16 @@
   @BeforeClass
   public static void prepareInput() throws NumberFormatException, 
ParseException {
     inputSchema =
-        RowSqlTypes.builder()
-            .withIntegerField("f_int")
-            .withBigIntField("f_long")
-            .withSmallIntField("f_short")
-            .withTinyIntField("f_byte")
-            .withFloatField("f_float")
-            .withDoubleField("f_double")
-            .withVarcharField("f_string")
-            .withTimestampField("f_timestamp")
-            .withIntegerField("f_int2")
+        Schema.builder()
+            .addInt32Field("f_int")
+            .addInt64Field("f_long")
+            .addInt16Field("f_short")
+            .addByteField("f_byte")
+            .addFloatField("f_float")
+            .addDoubleField("f_double")
+            .addStringField("f_string")
+            .addDateTimeField("f_timestamp")
+            .addInt32Field("f_int2")
             .build();
 
     inputRows =
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 7d3dd7e10aa..e56bcc5ee5c 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
@@ -28,7 +28,6 @@
 import java.util.List;
 import java.util.Map;
 import org.apache.beam.sdk.extensions.sql.BeamSql;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.extensions.sql.TestUtils;
 import org.apache.beam.sdk.extensions.sql.mock.MockedBoundedTable;
 import org.apache.beam.sdk.schemas.Schema;
@@ -60,19 +59,19 @@
           .build();
 
   private static final Schema ROW_TYPE =
-      RowSqlTypes.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")
+      Schema.builder()
+          .addDateTimeField("ts")
+          .addByteField("c_tinyint")
+          .addInt16Field("c_smallint")
+          .addInt32Field("c_integer")
+          .addInt64Field("c_bigint")
+          .addFloatField("c_float")
+          .addDoubleField("c_double")
+          .addDecimalField("c_decimal")
+          .addByteField("c_tinyint_max")
+          .addInt16Field("c_smallint_max")
+          .addInt32Field("c_integer_max")
+          .addInt64Field("c_bigint_max")
           .build();
 
   @Rule public final TestPipeline pipeline = TestPipeline.create();
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/integrationtest/BeamSqlComparisonOperatorsIntegrationTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/integrationtest/BeamSqlComparisonOperatorsIntegrationTest.java
index 1b4a22f6202..68948b7e1bc 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/integrationtest/BeamSqlComparisonOperatorsIntegrationTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/integrationtest/BeamSqlComparisonOperatorsIntegrationTest.java
@@ -19,7 +19,6 @@
 package org.apache.beam.sdk.extensions.sql.integrationtest;
 
 import java.math.BigDecimal;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.extensions.sql.mock.MockedBoundedTable;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.values.PCollection;
@@ -252,33 +251,33 @@ public void testIsNullAndIsNotNull() throws Exception {
   @Override
   protected PCollection<Row> getTestPCollection() {
     Schema type =
-        RowSqlTypes.builder()
-            .withTinyIntField("c_tinyint_0")
-            .withTinyIntField("c_tinyint_1")
-            .withTinyIntField("c_tinyint_2")
-            .withSmallIntField("c_smallint_0")
-            .withSmallIntField("c_smallint_1")
-            .withSmallIntField("c_smallint_2")
-            .withIntegerField("c_integer_0")
-            .withIntegerField("c_integer_1")
-            .withIntegerField("c_integer_2")
-            .withBigIntField("c_bigint_0")
-            .withBigIntField("c_bigint_1")
-            .withBigIntField("c_bigint_2")
-            .withFloatField("c_float_0")
-            .withFloatField("c_float_1")
-            .withFloatField("c_float_2")
-            .withDoubleField("c_double_0")
-            .withDoubleField("c_double_1")
-            .withDoubleField("c_double_2")
-            .withDecimalField("c_decimal_0")
-            .withDecimalField("c_decimal_1")
-            .withDecimalField("c_decimal_2")
-            .withVarcharField("c_varchar_0")
-            .withVarcharField("c_varchar_1")
-            .withVarcharField("c_varchar_2")
-            .withBooleanField("c_boolean_false")
-            .withBooleanField("c_boolean_true")
+        Schema.builder()
+            .addByteField("c_tinyint_0")
+            .addByteField("c_tinyint_1")
+            .addByteField("c_tinyint_2")
+            .addInt16Field("c_smallint_0")
+            .addInt16Field("c_smallint_1")
+            .addInt16Field("c_smallint_2")
+            .addInt32Field("c_integer_0")
+            .addInt32Field("c_integer_1")
+            .addInt32Field("c_integer_2")
+            .addInt64Field("c_bigint_0")
+            .addInt64Field("c_bigint_1")
+            .addInt64Field("c_bigint_2")
+            .addFloatField("c_float_0")
+            .addFloatField("c_float_1")
+            .addFloatField("c_float_2")
+            .addDoubleField("c_double_0")
+            .addDoubleField("c_double_1")
+            .addDoubleField("c_double_2")
+            .addDecimalField("c_decimal_0")
+            .addDecimalField("c_decimal_1")
+            .addDecimalField("c_decimal_2")
+            .addStringField("c_varchar_0")
+            .addStringField("c_varchar_1")
+            .addStringField("c_varchar_2")
+            .addBooleanField("c_boolean_false")
+            .addBooleanField("c_boolean_true")
             .build();
 
     try {
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/bigquery/BigQueryTableProviderTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/bigquery/BigQueryTableProviderTest.java
index 25ab11522b8..b5ae44a6936 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/bigquery/BigQueryTableProviderTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/bigquery/BigQueryTableProviderTest.java
@@ -24,7 +24,6 @@
 
 import java.util.stream.Stream;
 import org.apache.beam.sdk.extensions.sql.BeamSqlTable;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.extensions.sql.meta.Table;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.schemas.Schema.TypeName;
@@ -58,8 +57,8 @@ private static Table fakeTable(String name) {
         .location("project:dataset.table")
         .schema(
             Stream.of(
-                    Schema.Field.of("id", 
TypeName.INT32.type()).withNullable(true),
-                    Schema.Field.of("name", 
RowSqlTypes.VARCHAR).withNullable(true))
+                    Schema.Field.nullable("id", TypeName.INT32.type()),
+                    Schema.Field.nullable("name", Schema.FieldType.STRING))
                 .collect(toSchema()))
         .type("bigquery")
         .build();
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/kafka/KafkaTableProviderTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/kafka/KafkaTableProviderTest.java
index a91bacd9fd6..6d32744a8bc 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/kafka/KafkaTableProviderTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/kafka/KafkaTableProviderTest.java
@@ -27,7 +27,6 @@
 import com.google.common.collect.ImmutableList;
 import java.util.stream.Stream;
 import org.apache.beam.sdk.extensions.sql.BeamSqlTable;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.extensions.sql.meta.Table;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.schemas.Schema.TypeName;
@@ -69,8 +68,8 @@ private static Table mockTable(String name) {
         .location("kafka://localhost:2181/brokers?topic=test")
         .schema(
             Stream.of(
-                    Schema.Field.of("id", 
TypeName.INT32.type()).withNullable(true),
-                    Schema.Field.of("name", 
RowSqlTypes.VARCHAR).withNullable(true))
+                    Schema.Field.nullable("id", TypeName.INT32.type()),
+                    Schema.Field.nullable("name", Schema.FieldType.STRING))
                 .collect(toSchema()))
         .type("kafka")
         .properties(properties)
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubJsonIT.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubJsonIT.java
index 67184d453c8..f92a9014e71 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubJsonIT.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubJsonIT.java
@@ -26,7 +26,6 @@
 import java.util.Arrays;
 import java.util.List;
 import java.util.Set;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.extensions.sql.impl.BeamSqlEnv;
 import org.apache.beam.sdk.extensions.sql.meta.store.InMemoryMetaStore;
 import org.apache.beam.sdk.io.gcp.pubsub.PubsubIO;
@@ -53,7 +52,7 @@
 public class PubsubJsonIT implements Serializable {
 
   private static final Schema PAYLOAD_SCHEMA =
-      
RowSqlTypes.builder().withIntegerField("id").withVarcharField("name").build();
+      Schema.builder().addInt32Field("id").addStringField("name").build();
 
   @Rule public transient TestPubsub eventsTopic = TestPubsub.create();
   @Rule public transient TestPubsub dlqTopic = TestPubsub.create();
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubJsonTableProviderTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubJsonTableProviderTest.java
index d5a08545ab9..5bd6aaf58b1 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubJsonTableProviderTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubJsonTableProviderTest.java
@@ -18,12 +18,11 @@
 package org.apache.beam.sdk.extensions.sql.meta.provider.pubsub;
 
 import static junit.framework.TestCase.assertNotNull;
-import static org.apache.calcite.sql.type.SqlTypeName.VARCHAR;
+import static 
org.apache.beam.sdk.extensions.sql.impl.utils.CalciteUtils.VARCHAR;
 import static org.junit.Assert.assertEquals;
 
 import com.alibaba.fastjson.JSON;
 import org.apache.beam.sdk.extensions.sql.BeamSqlTable;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.extensions.sql.meta.Table;
 import org.apache.beam.sdk.schemas.Schema;
 import org.junit.Rule;
@@ -45,10 +44,10 @@ public void testTableTypePubsub() {
   public void testCreatesTable() {
     PubsubJsonTableProvider provider = new PubsubJsonTableProvider();
     Schema messageSchema =
-        RowSqlTypes.builder()
-            .withTimestampField("event_timestamp")
-            .withMapField("attributes", VARCHAR, VARCHAR)
-            .withRowField("payload", Schema.builder().build())
+        Schema.builder()
+            .addDateTimeField("event_timestamp")
+            .addMapField("attributes", VARCHAR, VARCHAR)
+            .addRowField("payload", Schema.builder().build())
             .build();
 
     Table tableDefinition = tableDefinition().schema(messageSchema).build();
@@ -63,9 +62,9 @@ public void testCreatesTable() {
   public void testThrowsIfTimestampFieldNotProvided() {
     PubsubJsonTableProvider provider = new PubsubJsonTableProvider();
     Schema messageSchema =
-        RowSqlTypes.builder()
-            .withMapField("attributes", VARCHAR, VARCHAR)
-            .withRowField("payload", Schema.builder().build())
+        Schema.builder()
+            .addMapField("attributes", VARCHAR, VARCHAR)
+            .addRowField("payload", Schema.builder().build())
             .build();
 
     Table tableDefinition = tableDefinition().schema(messageSchema).build();
@@ -79,9 +78,9 @@ public void testThrowsIfTimestampFieldNotProvided() {
   public void testThrowsIfAttributesFieldNotProvided() {
     PubsubJsonTableProvider provider = new PubsubJsonTableProvider();
     Schema messageSchema =
-        RowSqlTypes.builder()
-            .withTimestampField("event_timestamp")
-            .withRowField("payload", Schema.builder().build())
+        Schema.builder()
+            .addDateTimeField("event_timestamp")
+            .addRowField("payload", Schema.builder().build())
             .build();
 
     Table tableDefinition = tableDefinition().schema(messageSchema).build();
@@ -95,9 +94,9 @@ public void testThrowsIfAttributesFieldNotProvided() {
   public void testThrowsIfPayloadFieldNotProvided() {
     PubsubJsonTableProvider provider = new PubsubJsonTableProvider();
     Schema messageSchema =
-        RowSqlTypes.builder()
-            .withTimestampField("event_timestamp")
-            .withMapField("attributes", VARCHAR, VARCHAR)
+        Schema.builder()
+            .addDateTimeField("event_timestamp")
+            .addMapField("attributes", VARCHAR, VARCHAR)
             .build();
 
     Table tableDefinition = tableDefinition().schema(messageSchema).build();
@@ -111,11 +110,11 @@ public void testThrowsIfPayloadFieldNotProvided() {
   public void testThrowsIfExtraFieldsExist() {
     PubsubJsonTableProvider provider = new PubsubJsonTableProvider();
     Schema messageSchema =
-        RowSqlTypes.builder()
-            .withTimestampField("event_timestamp")
-            .withMapField("attributes", VARCHAR, VARCHAR)
-            .withVarcharField("someField")
-            .withRowField("payload", Schema.builder().build())
+        Schema.builder()
+            .addDateTimeField("event_timestamp")
+            .addMapField("attributes", VARCHAR, VARCHAR)
+            .addStringField("someField")
+            .addRowField("payload", Schema.builder().build())
             .build();
 
     Table tableDefinition = tableDefinition().schema(messageSchema).build();
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubMessageToRowTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubMessageToRowTest.java
index 683a9fa4d27..a843702691d 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubMessageToRowTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/pubsub/PubsubMessageToRowTest.java
@@ -20,9 +20,9 @@
 import static com.google.common.collect.Iterables.size;
 import static java.nio.charset.StandardCharsets.UTF_8;
 import static java.util.stream.Collectors.toSet;
+import static 
org.apache.beam.sdk.extensions.sql.impl.utils.CalciteUtils.VARCHAR;
 import static 
org.apache.beam.sdk.extensions.sql.meta.provider.pubsub.PubsubMessageToRow.DLQ_TAG;
 import static 
org.apache.beam.sdk.extensions.sql.meta.provider.pubsub.PubsubMessageToRow.MAIN_TAG;
-import static org.apache.calcite.sql.type.SqlTypeName.VARCHAR;
 import static org.junit.Assert.assertEquals;
 
 import com.google.common.collect.ImmutableMap;
@@ -32,7 +32,6 @@
 import java.util.Set;
 import java.util.function.Function;
 import java.util.stream.StreamSupport;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.io.gcp.pubsub.PubsubMessage;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.testing.PAssert;
@@ -56,14 +55,13 @@
 
   @Test
   public void testConvertsMessages() {
-    Schema payloadSchema =
-        
RowSqlTypes.builder().withIntegerField("id").withVarcharField("name").build();
+    Schema payloadSchema = 
Schema.builder().addInt32Field("id").addStringField("name").build();
 
     Schema messageSchema =
-        RowSqlTypes.builder()
-            .withTimestampField("event_timestamp")
-            .withMapField("attributes", VARCHAR, VARCHAR)
-            .withRowField("payload", payloadSchema)
+        Schema.builder()
+            .addDateTimeField("event_timestamp")
+            .addMapField("attributes", VARCHAR, VARCHAR)
+            .addRowField("payload", payloadSchema)
             .build();
 
     PCollection<Row> rows =
@@ -103,14 +101,13 @@ public void testConvertsMessages() {
 
   @Test
   public void testSendsInvalidToDLQ() {
-    Schema payloadSchema =
-        
RowSqlTypes.builder().withIntegerField("id").withVarcharField("name").build();
+    Schema payloadSchema = 
Schema.builder().addInt32Field("id").addStringField("name").build();
 
     Schema messageSchema =
-        RowSqlTypes.builder()
-            .withTimestampField("event_timestamp")
-            .withMapField("attributes", VARCHAR, VARCHAR)
-            .withRowField("payload", payloadSchema)
+        Schema.builder()
+            .addDateTimeField("event_timestamp")
+            .addMapField("attributes", VARCHAR, VARCHAR)
+            .addRowField("payload", payloadSchema)
             .build();
 
     PCollectionTuple outputs =
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/text/BeamTextCSVTableTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/text/BeamTextCSVTableTest.java
index 98099202cd0..916d4f935e9 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/text/BeamTextCSVTableTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/text/BeamTextCSVTableTest.java
@@ -30,7 +30,6 @@
 import java.nio.file.attribute.BasicFileAttributes;
 import java.util.Arrays;
 import java.util.List;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.testing.PAssert;
 import org.apache.beam.sdk.testing.TestPipeline;
@@ -55,12 +54,12 @@
    * <p>The types of the csv fields are: integer,bigint,float,double,string
    */
   private static Schema schema =
-      RowSqlTypes.builder()
-          .withIntegerField("id")
-          .withBigIntField("order_id")
-          .withFloatField("price")
-          .withDoubleField("amount")
-          .withVarcharField("user_name")
+      Schema.builder()
+          .addInt32Field("id")
+          .addInt64Field("order_id")
+          .addFloatField("price")
+          .addDoubleField("amount")
+          .addStringField("user_name")
           .build();
 
   private static Object[] data1 = new Object[] {1, 1L, 1.1F, 1.1, "james"};
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/text/TextTableProviderTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/text/TextTableProviderTest.java
index f246d2835db..d4dbca55808 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/text/TextTableProviderTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/text/TextTableProviderTest.java
@@ -25,7 +25,6 @@
 import com.alibaba.fastjson.JSONObject;
 import java.util.stream.Stream;
 import org.apache.beam.sdk.extensions.sql.BeamSqlTable;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.extensions.sql.meta.Table;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.schemas.Schema.TypeName;
@@ -77,8 +76,8 @@ private static Table mockTable(String name, String format) {
         .location("/home/admin/" + name)
         .schema(
             Stream.of(
-                    Schema.Field.of("id", 
TypeName.INT32.type()).withNullable(true),
-                    Schema.Field.of("name", 
RowSqlTypes.VARCHAR).withNullable(true))
+                    Schema.Field.nullable("id", TypeName.INT32.type()),
+                    Schema.Field.nullable("name", Schema.FieldType.STRING))
                 .collect(toSchema()))
         .type("text")
         .properties(properties)
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/store/InMemoryMetaStoreTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/store/InMemoryMetaStoreTest.java
index e53e648a5df..90f93fec2e7 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/store/InMemoryMetaStoreTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/store/InMemoryMetaStoreTest.java
@@ -27,7 +27,6 @@
 import java.util.Map;
 import java.util.stream.Stream;
 import org.apache.beam.sdk.extensions.sql.BeamSqlTable;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.extensions.sql.meta.Table;
 import org.apache.beam.sdk.extensions.sql.meta.provider.TableProvider;
 import org.apache.beam.sdk.extensions.sql.meta.provider.text.TextTableProvider;
@@ -86,7 +85,10 @@ public void testBuildBeamSqlTable() throws Exception {
     BeamSqlTable actualSqlTable = store.buildBeamSqlTable(table);
     assertNotNull(actualSqlTable);
     assertEquals(
-        
RowSqlTypes.builder().withIntegerField("id").withVarcharField("name").build(),
+        Schema.builder()
+            .addNullableField("id", Schema.FieldType.INT32)
+            .addNullableField("name", Schema.FieldType.STRING)
+            .build(),
         actualSqlTable.getSchema());
   }
 
@@ -120,8 +122,8 @@ private static Table mockTable(String name, String type) {
         .location("/home/admin/" + name)
         .schema(
             Stream.of(
-                    Schema.Field.of("id", 
TypeName.INT32.type()).withNullable(true),
-                    Schema.Field.of("name", 
RowSqlTypes.VARCHAR).withNullable(true))
+                    Schema.Field.nullable("id", TypeName.INT32.type()),
+                    Schema.Field.nullable("name", Schema.FieldType.STRING))
                 .collect(toSchema()))
         .type(type)
         .properties(new JSONObject())
diff --git 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryUtilsTest.java
 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryUtilsTest.java
index 99ccb9d98f0..888d7af63a7 100644
--- 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryUtilsTest.java
+++ 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryUtilsTest.java
@@ -43,11 +43,11 @@
 public class BigQueryUtilsTest {
   private static final Schema FLAT_TYPE = Schema
       .builder()
-      .addField(Schema.Field.nullable("id", Schema.FieldType.INT64))
-      .addField(Schema.Field.nullable("value", Schema.FieldType.DOUBLE))
-      .addField(Schema.Field.nullable("name", Schema.FieldType.STRING))
-      .addField(Schema.Field.nullable("timestamp", Schema.FieldType.DATETIME))
-      .addField(Schema.Field.nullable("valid", Schema.FieldType.BOOLEAN))
+      .addNullableField("id", Schema.FieldType.INT64)
+      .addNullableField("value", Schema.FieldType.DOUBLE)
+      .addNullableField("name", Schema.FieldType.STRING)
+      .addNullableField("timestamp", Schema.FieldType.DATETIME)
+      .addNullableField("valid", Schema.FieldType.BOOLEAN)
       .build();
 
   private static final Schema ARRAY_TYPE = Schema
@@ -57,7 +57,7 @@
 
   private static final Schema ROW_TYPE = Schema
       .builder()
-      .addField(Schema.Field.nullable("row", Schema.FieldType.row(FLAT_TYPE)))
+      .addNullableField("row", Schema.FieldType.row(FLAT_TYPE))
       .build();
 
   private static final Schema ARRAY_ROW_TYPE =
diff --git 
a/sdks/java/nexmark/src/main/java/org/apache/beam/sdk/nexmark/model/sql/adapter/ModelAdaptersMapping.java
 
b/sdks/java/nexmark/src/main/java/org/apache/beam/sdk/nexmark/model/sql/adapter/ModelAdaptersMapping.java
index e781a0a9b09..72403932144 100644
--- 
a/sdks/java/nexmark/src/main/java/org/apache/beam/sdk/nexmark/model/sql/adapter/ModelAdaptersMapping.java
+++ 
b/sdks/java/nexmark/src/main/java/org/apache/beam/sdk/nexmark/model/sql/adapter/ModelAdaptersMapping.java
@@ -23,13 +23,14 @@
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
+import org.apache.beam.sdk.extensions.sql.impl.utils.CalciteUtils;
 import org.apache.beam.sdk.nexmark.model.Auction;
 import org.apache.beam.sdk.nexmark.model.AuctionCount;
 import org.apache.beam.sdk.nexmark.model.AuctionPrice;
 import org.apache.beam.sdk.nexmark.model.Bid;
 import org.apache.beam.sdk.nexmark.model.NameCityStateId;
 import org.apache.beam.sdk.nexmark.model.Person;
+import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.values.Row;
 import org.joda.time.DateTime;
 
@@ -50,15 +51,15 @@
 
   private static ModelFieldsAdapter<Person> personAdapter() {
     return new ModelFieldsAdapter<Person>(
-        RowSqlTypes.builder()
-            .withBigIntField("id")
-            .withVarcharField("name")
-            .withVarcharField("emailAddress")
-            .withVarcharField("creditCard")
-            .withVarcharField("city")
-            .withVarcharField("state")
-            .withTimestampField("dateTime")
-            .withVarcharField("extra")
+        Schema.builder()
+            .addInt64Field("id")
+            .addStringField("name")
+            .addStringField("emailAddress")
+            .addStringField("creditCard")
+            .addStringField("city")
+            .addStringField("state")
+            .addField("dateTime", CalciteUtils.TIME)
+            .addStringField("extra")
             .build()) {
       @Override
       public List<Object> getFieldsValues(Person p) {
@@ -90,12 +91,12 @@ public Person getRowModel(Row row) {
 
   private static ModelFieldsAdapter<Bid> bidAdapter() {
     return new ModelFieldsAdapter<Bid>(
-        RowSqlTypes.builder()
-            .withBigIntField("auction")
-            .withBigIntField("bidder")
-            .withBigIntField("price")
-            .withTimestampField("dateTime")
-            .withVarcharField("extra")
+        Schema.builder()
+            .addInt64Field("auction")
+            .addInt64Field("bidder")
+            .addInt64Field("price")
+            .addField("dateTime", CalciteUtils.TIME)
+            .addStringField("extra")
             .build()) {
       @Override
       public List<Object> getFieldsValues(Bid b) {
@@ -121,17 +122,17 @@ public Bid getRowModel(Row row) {
 
   private static ModelFieldsAdapter<Auction> auctionAdapter() {
     return new ModelFieldsAdapter<Auction>(
-        RowSqlTypes.builder()
-            .withBigIntField("id")
-            .withVarcharField("itemName")
-            .withVarcharField("description")
-            .withBigIntField("initialBid")
-            .withBigIntField("reserve")
-            .withTimestampField("dateTime")
-            .withTimestampField("expires")
-            .withBigIntField("seller")
-            .withBigIntField("category")
-            .withVarcharField("extra")
+        Schema.builder()
+            .addInt64Field("id")
+            .addStringField("itemName")
+            .addStringField("description")
+            .addInt64Field("initialBid")
+            .addInt64Field("reserve")
+            .addField("dateTime", CalciteUtils.TIMESTAMP)
+            .addField("expires", CalciteUtils.TIMESTAMP)
+            .addInt64Field("seller")
+            .addInt64Field("category")
+            .addStringField("extra")
             .build()) {
       @Override
       public List<Object> getFieldsValues(Auction a) {
@@ -148,6 +149,7 @@ public Bid getRowModel(Row row) {
                 a.category,
                 a.extra));
       }
+
       @Override
       public Auction getRowModel(Row row) {
         return new Auction(
@@ -167,9 +169,9 @@ public Auction getRowModel(Row row) {
 
   private static ModelFieldsAdapter<AuctionCount> auctionCountAdapter() {
     return new ModelFieldsAdapter<AuctionCount>(
-        RowSqlTypes.builder()
-            .withBigIntField("auction")
-            .withBigIntField("num")
+        Schema.builder()
+            .addInt64Field("auction")
+            .addInt64Field("num")
             .build()) {
       @Override
       public List<Object> getFieldsValues(AuctionCount a) {
@@ -189,9 +191,9 @@ public AuctionCount getRowModel(Row row) {
 
   private static ModelFieldsAdapter<AuctionPrice> auctionPriceAdapter() {
     return new ModelFieldsAdapter<AuctionPrice>(
-        RowSqlTypes.builder()
-            .withBigIntField("auction")
-            .withBigIntField("price")
+        Schema.builder()
+            .addInt64Field("auction")
+            .addInt64Field("price")
             .build()) {
       @Override
       public List<Object> getFieldsValues(AuctionPrice a) {
@@ -211,11 +213,11 @@ public AuctionPrice getRowModel(Row row) {
 
   private static ModelFieldsAdapter<NameCityStateId> nameCityStateIdAdapter() {
     return new ModelFieldsAdapter<NameCityStateId>(
-        RowSqlTypes.builder()
-            .withVarcharField("name")
-            .withVarcharField("city")
-            .withVarcharField("state")
-            .withBigIntField("id")
+        Schema.builder()
+            .addStringField("name")
+            .addStringField("city")
+            .addStringField("state")
+            .addInt64Field("id")
             .build()) {
       @Override
       public List<Object> getFieldsValues(NameCityStateId a) {
diff --git 
a/sdks/java/nexmark/src/main/java/org/apache/beam/sdk/nexmark/queries/sql/SqlQuery3.java
 
b/sdks/java/nexmark/src/main/java/org/apache/beam/sdk/nexmark/queries/sql/SqlQuery3.java
index df9d1768bc8..f19909ec73c 100644
--- 
a/sdks/java/nexmark/src/main/java/org/apache/beam/sdk/nexmark/queries/sql/SqlQuery3.java
+++ 
b/sdks/java/nexmark/src/main/java/org/apache/beam/sdk/nexmark/queries/sql/SqlQuery3.java
@@ -21,7 +21,6 @@
 
 import org.apache.beam.sdk.coders.RowCoder;
 import org.apache.beam.sdk.extensions.sql.BeamSql;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.nexmark.NexmarkConfiguration;
 import org.apache.beam.sdk.nexmark.model.Auction;
 import org.apache.beam.sdk.nexmark.model.Event;
@@ -29,6 +28,7 @@
 import org.apache.beam.sdk.nexmark.model.Person;
 import org.apache.beam.sdk.nexmark.model.sql.ToRow;
 import org.apache.beam.sdk.nexmark.queries.Query3;
+import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.transforms.Filter;
 import org.apache.beam.sdk.transforms.PTransform;
 import org.apache.beam.sdk.transforms.ParDo;
@@ -146,12 +146,12 @@ private RowCoder getRecordCoder(Class modelClass) {
 
   private static RowCoder createRecordCoder() {
     return
-        RowSqlTypes
+        Schema
             .builder()
-            .withVarcharField("name")
-            .withVarcharField("city")
-            .withVarcharField("state")
-            .withBigIntField("id")
+            .addStringField("name")
+            .addStringField("city")
+            .addStringField("state")
+            .addInt64Field("id")
             .build()
             .getRowCoder();
   }
diff --git 
a/sdks/java/nexmark/src/test/java/org/apache/beam/sdk/nexmark/model/sql/RowSizeTest.java
 
b/sdks/java/nexmark/src/test/java/org/apache/beam/sdk/nexmark/model/sql/RowSizeTest.java
index 51ff3ccf7db..2dbd8be377c 100644
--- 
a/sdks/java/nexmark/src/test/java/org/apache/beam/sdk/nexmark/model/sql/RowSizeTest.java
+++ 
b/sdks/java/nexmark/src/test/java/org/apache/beam/sdk/nexmark/model/sql/RowSizeTest.java
@@ -24,7 +24,7 @@
 
 import com.google.common.collect.Iterables;
 import java.math.BigDecimal;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
+import org.apache.beam.sdk.extensions.sql.impl.utils.CalciteUtils;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.testing.PAssert;
 import org.apache.beam.sdk.testing.TestPipeline;
@@ -42,20 +42,20 @@
  */
 public class RowSizeTest {
 
-  private static final Schema ROW_TYPE = RowSqlTypes.builder()
-      .withTinyIntField("f_tinyint")
-      .withSmallIntField("f_smallint")
-      .withIntegerField("f_int")
-      .withBigIntField("f_bigint")
-      .withFloatField("f_float")
-      .withDoubleField("f_double")
-      .withDecimalField("f_decimal")
-      .withBooleanField("f_boolean")
-      .withTimeField("f_time")
-      .withDateField("f_date")
-      .withTimestampField("f_timestamp")
-      .withCharField("f_char")
-      .withVarcharField("f_varchar")
+  private static final Schema ROW_TYPE = Schema.builder()
+      .addByteField("f_tinyint")
+      .addInt16Field("f_smallint")
+      .addInt32Field("f_int")
+      .addInt64Field("f_bigint")
+      .addFloatField("f_float")
+      .addDoubleField("f_double")
+      .addDecimalField("f_decimal")
+      .addBooleanField("f_boolean")
+      .addField("f_time", CalciteUtils.TIME)
+      .addField("f_date", CalciteUtils.DATE)
+      .addDateTimeField("f_timestamp")
+      .addField("f_char", CalciteUtils.CHAR)
+      .addField("f_varchar", CalciteUtils.VARCHAR)
       .build();
 
   private static final long ROW_SIZE = 96L;
diff --git 
a/sdks/java/nexmark/src/test/java/org/apache/beam/sdk/nexmark/model/sql/adapter/ModelAdaptersMappingTest.java
 
b/sdks/java/nexmark/src/test/java/org/apache/beam/sdk/nexmark/model/sql/adapter/ModelAdaptersMappingTest.java
index 4d9b353c715..cd3955d63a4 100644
--- 
a/sdks/java/nexmark/src/test/java/org/apache/beam/sdk/nexmark/model/sql/adapter/ModelAdaptersMappingTest.java
+++ 
b/sdks/java/nexmark/src/test/java/org/apache/beam/sdk/nexmark/model/sql/adapter/ModelAdaptersMappingTest.java
@@ -23,7 +23,6 @@
 import static org.junit.Assert.assertTrue;
 
 import java.util.List;
-import org.apache.beam.sdk.extensions.sql.RowSqlTypes;
 import org.apache.beam.sdk.nexmark.model.Auction;
 import org.apache.beam.sdk.nexmark.model.Bid;
 import org.apache.beam.sdk.nexmark.model.Person;
@@ -39,42 +38,42 @@
   private static final Person PERSON =
       new Person(3L, "name", "email", "cc", "city", "state", 329823L, "extra");
 
-  private static final Schema PERSON_ROW_TYPE = RowSqlTypes.builder()
-      .withBigIntField("id")
-      .withVarcharField("name")
-      .withVarcharField("emailAddress")
-      .withVarcharField("creditCard")
-      .withVarcharField("city")
-      .withVarcharField("state")
-      .withTimestampField("dateTime")
-      .withVarcharField("extra")
+  private static final Schema PERSON_ROW_TYPE = Schema.builder()
+      .addInt64Field("id")
+      .addStringField("name")
+      .addStringField("emailAddress")
+      .addStringField("creditCard")
+      .addStringField("city")
+      .addStringField("state")
+      .addDateTimeField("dateTime")
+      .addStringField("extra")
       .build();
 
   private static final Bid BID =
       new Bid(5L, 3L, 123123L, 43234234L, "extra2");
 
-  private static final Schema BID_ROW_TYPE = RowSqlTypes.builder()
-      .withBigIntField("auction")
-      .withBigIntField("bidder")
-      .withBigIntField("price")
-      .withTimestampField("dateTime")
-      .withVarcharField("extra")
+  private static final Schema BID_ROW_TYPE = Schema.builder()
+      .addInt64Field("auction")
+      .addInt64Field("bidder")
+      .addInt64Field("price")
+      .addDateTimeField("dateTime")
+      .addStringField("extra")
       .build();
 
   private static final Auction AUCTION =
       new Auction(5L, "item", "desc", 342L, 321L, 3423342L, 2349234L, 3L, 1L, 
"extra3");
 
-  private static final Schema AUCTION_ROW_TYPE = RowSqlTypes.builder()
-      .withBigIntField("id")
-      .withVarcharField("itemName")
-      .withVarcharField("description")
-      .withBigIntField("initialBid")
-      .withBigIntField("reserve")
-      .withTimestampField("dateTime")
-      .withTimestampField("expires")
-      .withBigIntField("seller")
-      .withBigIntField("category")
-      .withVarcharField("extra")
+  private static final Schema AUCTION_ROW_TYPE = Schema.builder()
+      .addInt64Field("id")
+      .addStringField("itemName")
+      .addStringField("description")
+      .addInt64Field("initialBid")
+      .addInt64Field("reserve")
+      .addDateTimeField("dateTime")
+      .addDateTimeField("expires")
+      .addInt64Field("seller")
+      .addInt64Field("category")
+      .addStringField("extra")
       .build();
 
   @Test


 

----------------------------------------------------------------
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]


Issue Time Tracking
-------------------

    Worklog Id:     (was: 107383)
    Time Spent: 1.5h  (was: 1h 20m)

> remove RowSqlTypeBuilder
> ------------------------
>
>                 Key: BEAM-4082
>                 URL: https://issues.apache.org/jira/browse/BEAM-4082
>             Project: Beam
>          Issue Type: Sub-task
>          Components: sdk-java-core
>            Reporter: Kenneth Knowles
>            Assignee: Kenneth Knowles
>            Priority: Major
>          Time Spent: 1.5h
>  Remaining Estimate: 0h
>
> Quoting the PR, this should be removed.



--
This message was sent by Atlassian JIRA
(v7.6.3#76005)

Reply via email to