Implement bigquery calcite schema (#11)

This commit is contained in:
Kuan-Po Tseng
2022-05-26 15:48:00 +08:00
committed by GitHub
parent e5d6cb5a75
commit 36809385f1
9 changed files with 335 additions and 0 deletions
+5
View File
@@ -25,6 +25,11 @@
<artifactId>guava</artifactId>
</dependency>
<dependency>
<groupId>io.cml</groupId>
<artifactId>cml-spi</artifactId>
</dependency>
<dependency>
<groupId>io.cml</groupId>
<artifactId>trino-parser</artifactId>
@@ -0,0 +1,39 @@
/*
* Licensed 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 io.cml.calcite;
import org.apache.calcite.schema.Table;
import org.apache.calcite.schema.impl.AbstractSchema;
import java.util.Map;
import static java.util.Objects.requireNonNull;
class CmlSchema
extends AbstractSchema
{
private final Map<String, Table> tableMap;
public CmlSchema(Map<String, Table> tableMap)
{
this.tableMap = requireNonNull(tableMap, "tableMap is null");
}
@Override
protected Map<String, Table> getTableMap()
{
return tableMap;
}
}
@@ -0,0 +1,148 @@
/*
* Licensed 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 io.cml.calcite;
import com.google.common.collect.ImmutableList;
import io.cml.spi.metadata.ColumnMetadata;
import io.cml.spi.metadata.TableMetadata;
import io.cml.spi.type.PGType;
import io.trino.sql.parser.ParsingOptions;
import io.trino.sql.parser.SqlParser;
import io.trino.sql.tree.Statement;
import org.apache.calcite.adapter.java.JavaTypeFactory;
import org.apache.calcite.config.CalciteConnectionConfigImpl;
import org.apache.calcite.jdbc.CalciteSchema;
import org.apache.calcite.jdbc.JavaTypeFactoryImpl;
import org.apache.calcite.plan.ConventionTraitDef;
import org.apache.calcite.plan.RelOptCluster;
import org.apache.calcite.plan.RelOptPlanner;
import org.apache.calcite.plan.volcano.VolcanoPlanner;
import org.apache.calcite.prepare.CalciteCatalogReader;
import org.apache.calcite.rel.type.RelDataType;
import org.apache.calcite.rel.type.RelDataTypeFactory;
import org.apache.calcite.rex.RexBuilder;
import org.apache.calcite.schema.SchemaPlus;
import org.apache.calcite.sql.SqlDialect;
import org.apache.calcite.sql.dialect.BigQuerySqlDialect;
import org.apache.calcite.tools.Frameworks;
import java.util.List;
import static com.google.common.collect.ImmutableMap.toImmutableMap;
import static io.cml.spi.type.BigIntType.BIGINT;
import static io.cml.spi.type.BooleanType.BOOLEAN;
import static io.cml.spi.type.DoubleType.DOUBLE;
import static io.cml.spi.type.IntegerType.INTEGER;
import static io.cml.spi.type.VarcharType.VARCHAR;
public final class CmlSchemaUtil
{
private CmlSchemaUtil() {}
public enum Dialect
{
BIGQUERY(BigQuerySqlDialect.DEFAULT);
private final SqlDialect sqlDialect;
Dialect(SqlDialect sqlDialect)
{
this.sqlDialect = sqlDialect;
}
}
public static String convertQuery(Dialect dialect, SchemaPlusInfo schemaPlusInfo, String sql)
{
SqlParser sqlParser = new SqlParser();
Statement stmt = sqlParser.createStatement(sql, new ParsingOptions());
RelOptCluster cluster = newCluster();
SchemaPlus schemaPlus = get(schemaPlusInfo);
CalciteCatalogReader reader = new CalciteCatalogReader(
CalciteSchema.from(schemaPlus),
ImmutableList.of(),
cluster.getTypeFactory(),
CalciteConnectionConfigImpl.DEFAULT);
// TODO: uncomment this when CalciteRelConverter finished
// RelNode relNode = CalciteRelConverter.convert(cluster, reader, stmt);
// RelToSqlConverter relToSqlConverter = new RelToSqlConverter(dialect.sqlDialect);
// SqlNode sqlNode = relToSqlConverter.visitRoot(relNode).asStatement();
// SqlPrettyWriter sqlPrettyWriter = new SqlPrettyWriter();
// return sqlPrettyWriter.format(sqlNode);
return "";
}
private static RelOptCluster newCluster()
{
RelDataTypeFactory typeFactory = new JavaTypeFactoryImpl();
RelOptPlanner planner = new VolcanoPlanner();
planner.addRelTraitDef(ConventionTraitDef.INSTANCE);
return RelOptCluster.create(planner, new RexBuilder(typeFactory));
}
private static SchemaPlus get(SchemaPlusInfo schemaPlusInfo)
{
SchemaPlus rootSchema = Frameworks.createRootSchema(true);
schemaPlusInfo.getSchemaTableMap()
.forEach((schema, tables) -> rootSchema.add(schema, toCmlSchema(tables)));
return rootSchema;
}
private static CmlTable toCmlTable(TableMetadata tableMetadata)
{
JavaTypeFactoryImpl typeFactory = new JavaTypeFactoryImpl();
RelDataTypeFactory.Builder builder = new RelDataTypeFactory.Builder(typeFactory);
for (ColumnMetadata columnMetadata : tableMetadata.getColumns()) {
builder.add(columnMetadata.getName(), toRelDataType(typeFactory, columnMetadata.getType()));
}
return new CmlTable(builder.build());
}
private static CmlSchema toCmlSchema(List<TableMetadata> tables)
{
return new CmlSchema(tables.stream().collect(
toImmutableMap(
table -> table.getTable().getTableName(),
CmlSchemaUtil::toCmlTable,
// TODO: handle case sensitive table name
(a, b) -> a)));
}
// TODO: handle nested types
private static RelDataType toRelDataType(JavaTypeFactory typeFactory, PGType<?> pgType)
{
if (pgType.equals(BOOLEAN)) {
return typeFactory.createJavaType(Boolean.class);
}
if (pgType.equals(INTEGER)) {
return typeFactory.createJavaType(Integer.class);
}
if (pgType.equals(BIGINT)) {
return typeFactory.createJavaType(Long.class);
}
if (pgType.equals(VARCHAR)) {
return typeFactory.createJavaType(String.class);
}
if (pgType.equals(DOUBLE)) {
return typeFactory.createJavaType(Double.class);
}
throw new UnsupportedOperationException(pgType.type() + " not supported yet");
}
}
@@ -0,0 +1,36 @@
/*
* Licensed 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 io.cml.calcite;
import org.apache.calcite.rel.type.RelDataType;
import org.apache.calcite.rel.type.RelDataTypeFactory;
import org.apache.calcite.schema.impl.AbstractTable;
class CmlTable
extends AbstractTable
{
private final RelDataType rowType;
CmlTable(RelDataType relDataType)
{
this.rowType = relDataType;
}
@Override
public RelDataType getRowType(RelDataTypeFactory typeFactory)
{
return rowType;
}
}
@@ -0,0 +1,35 @@
/*
* Licensed 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 io.cml.calcite;
import io.cml.spi.metadata.TableMetadata;
import java.util.List;
import java.util.Map;
public class SchemaPlusInfo
{
private final Map<String, List<TableMetadata>> schemaTableMap;
public SchemaPlusInfo(Map<String, List<TableMetadata>> schemaTableMap)
{
this.schemaTableMap = schemaTableMap;
}
public Map<String, List<TableMetadata>> getSchemaTableMap()
{
return schemaTableMap;
}
}
+5
View File
@@ -73,6 +73,11 @@
<artifactId>guice</artifactId>
</dependency>
<dependency>
<groupId>io.cml</groupId>
<artifactId>cml-calcite</artifactId>
</dependency>
<dependency>
<groupId>io.cml</groupId>
<artifactId>cml-spi</artifactId>
@@ -0,0 +1,51 @@
/*
* Licensed 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 io.cml.connector.bigquery;
import io.cml.calcite.CmlSchemaUtil;
import io.cml.calcite.SchemaPlusInfo;
import io.cml.spi.connector.Connector;
import io.cml.spi.metadata.TableMetadata;
import javax.inject.Inject;
import java.util.List;
import java.util.Map;
import static com.google.common.collect.ImmutableMap.toImmutableMap;
import static java.util.Objects.requireNonNull;
import static java.util.function.Function.identity;
public class BigQuerySqlConverter
{
private final SchemaPlusInfo schemaPlusInfo;
@Inject
public BigQuerySqlConverter(Connector bigQueryConnector)
{
final Connector connector = requireNonNull(bigQueryConnector, "BigQueryConnector is null");
List<String> schemas = connector.listSchemas();
Map<String, List<TableMetadata>> schemaTableMap = schemas.stream()
.collect(toImmutableMap(identity(), schema -> connector.listTables(schema)));
this.schemaPlusInfo = new SchemaPlusInfo(schemaTableMap);
}
public String convertSql(String sql)
{
return CmlSchemaUtil.convertQuery(CmlSchemaUtil.Dialect.BIGQUERY, schemaPlusInfo, sql);
}
}
@@ -28,6 +28,7 @@ import io.airlift.configuration.AbstractConfigurationAwareModule;
import io.cml.connector.bigquery.BigQueryConfig;
import io.cml.connector.bigquery.BigQueryConnector;
import io.cml.connector.bigquery.BigQueryCredentialsSupplier;
import io.cml.connector.bigquery.BigQuerySqlConverter;
import io.cml.pgcatalog.builder.BigQueryPgCatalogTableBuilder;
import io.cml.pgcatalog.builder.BigQueryPgFunctionBuilder;
import io.cml.pgcatalog.builder.PgCatalogTableBuilder;
@@ -50,6 +51,7 @@ public class BigQueryConnectorModule
binder.bind(PgCatalogTableBuilder.class).to(BigQueryPgCatalogTableBuilder.class).in(Scopes.SINGLETON);
binder.bind(PgFunctionBuilder.class).to(BigQueryPgFunctionBuilder.class).in(Scopes.SINGLETON);
binder.bind(PgMetadata.class).to(BigQueryPgMetadata.class).in(Scopes.SINGLETON);
binder.bind(BigQuerySqlConverter.class).in(Scopes.SINGLETON);
configBinder(binder).bindConfig(BigQueryConfig.class);
}
+14
View File
@@ -188,6 +188,12 @@
<version>1.15</version>
</dependency>
<dependency>
<groupId>io.cml</groupId>
<artifactId>cml-calcite</artifactId>
<version>${project.version}</version>
</dependency>
<dependency>
<groupId>io.cml</groupId>
<artifactId>cml-main</artifactId>
@@ -245,10 +251,18 @@
<groupId>commons-logging</groupId>
<artifactId>commons-logging</artifactId>
</exclusion>
<exclusion>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
</exclusion>
<exclusion>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
</exclusion>
<exclusion>
<groupId>org.yaml</groupId>
<artifactId>snakeyaml</artifactId>
</exclusion>
</exclusions>
</dependency>