Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -381,7 +381,7 @@ SqlCreate SqlCreateDatabase(Span s, boolean replace) :
}

/**
* USE DATABASE ( catalog_name '.' )? database_name
* USE [ DATABASE ] ( catalog_name '.' )? database_name
*/
SqlCall SqlUseDatabase(Span s, String scope) :
{
Expand All @@ -391,7 +391,7 @@ SqlCall SqlUseDatabase(Span s, String scope) :
<USE> {
s.add(this);
}
<DATABASE>
[ <DATABASE> ]
databaseName = CompoundIdentifier()
{
return new SqlUseDatabase(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ public class BeamCalciteTable extends AbstractQueryableTable
private final Map<String, String> pipelineOptionsMap;
private @Nullable PipelineOptions pipelineOptions;

BeamCalciteTable(
public BeamCalciteTable(
BeamSqlTable beamTable,
Map<String, String> pipelineOptionsMap,
@Nullable PipelineOptions pipelineOptions) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,17 +47,23 @@
import org.apache.beam.sdk.transforms.SerializableFunction;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.jdbc.CalcitePrepare;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.plan.RelOptUtil;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.rel.RelNode;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.schema.Function;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.sql.SqlKind;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.tools.RelBuilder;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.tools.RuleSet;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

/**
* Contains the metadata of tables/UDF functions, and exposes APIs to
* query/validate/optimize/translate SQL statements.
*/
@Internal
public class BeamSqlEnv {
private static final Logger LOG = LoggerFactory.getLogger(BeamSqlEnv.class);

JdbcConnection connection;
QueryPlanner planner;

Expand Down Expand Up @@ -116,6 +122,31 @@ public BeamRelNode parseQuery(String query, QueryParameters queryParameters)
return planner.convertToBeamRel(query, queryParameters);
}

public QueryPlanner getPlanner() {
return planner;
}

public RelBuilder getRelBuilder() {
return planner.getRelBuilder();
}

public BeamRelNode convertToBeamRel(RelNode relNode) {
return planner.convertToBeamRel(relNode, QueryParameters.ofNone());
}

public RelNode parseLogicalPlan(String query) throws ParseException {
return planner.parseToRel(query, QueryParameters.ofNone());
}

public void registerSchemaFunction(String name, Function function) {
connection.getCurrentSchemaPlus().add(name, function);
}

public org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.sql.SqlOperatorTable
getOperatorTable() {
return planner.getOperatorTable();
}

public boolean isDdl(String sqlStatement) throws ParseException {
return planner.parse(sqlStatement).getKind().belongsTo(SqlKind.DDL);
}
Expand Down Expand Up @@ -196,6 +227,7 @@ public BeamSqlEnvBuilder setCurrentSchema(String name) {

/** Set the ruleSet used for query optimizer. */
public BeamSqlEnvBuilder setRuleSets(Collection<RuleSet> ruleSets) {
LOG.info("Setting BeamSqlEnv rulesets to: {}", ruleSets);
this.ruleSets = ruleSets;
return this;
}
Expand Down Expand Up @@ -262,6 +294,7 @@ public BeamSqlEnv build() {

configureSchemas(jdbcConnection);

LOG.info("Instantiating planner with ruleSets: {}", ruleSets);
QueryPlanner planner = instantiatePlanner(jdbcConnection, ruleSets);

// The planner may choose to add its own builtin functions to the schema, so load user-defined
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,12 +31,15 @@
import org.apache.beam.sdk.extensions.sql.impl.rel.BeamRelNode;
import org.apache.beam.sdk.extensions.sql.impl.rel.BeamSqlRelUtils;
import org.apache.beam.sdk.extensions.sql.impl.udf.BeamBuiltinFunctionProvider;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.vendor.calcite.v1_40_0.com.google.common.collect.Table;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.config.CalciteConnectionConfig;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.jdbc.CalciteSchema;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.plan.Contexts;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.plan.ConventionTraitDef;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.plan.RelOptCluster;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.plan.RelOptCost;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.plan.RelOptPlanner;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.plan.RelOptPlanner.CannotPlanException;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.plan.RelTraitDef;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.plan.RelTraitSet;
Expand All @@ -52,6 +55,7 @@
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.rel.metadata.RelMetadataProvider;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.rel.metadata.RelMetadataQuery;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.rel.type.RelDataType;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.rel.type.RelDataTypeFactory;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.rex.RexBuilder;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.rex.RexDynamicParam;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.rex.RexNode;
Expand All @@ -64,10 +68,14 @@
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.sql.parser.SqlParser;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.sql.parser.SqlParserImplFactory;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.sql.util.SqlOperatorTables;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.sql.validate.SqlConformance;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.sql.validate.SqlConformanceEnum;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.sql2rel.SqlToRelConverter;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.tools.FrameworkConfig;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.tools.Frameworks;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.tools.Planner;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.tools.Program;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.tools.RelBuilder;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.tools.RelConversionException;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.tools.RuleSet;
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.tools.ValidationException;
Expand All @@ -90,11 +98,52 @@ public class CalciteQueryPlanner implements QueryPlanner {

private final Planner planner;
private final JdbcConnection connection;
private final FrameworkConfig config;

// Cannot be final because of wacky initialization logic
private RelOptCluster relOptCluster;
private CalciteCatalogReader catalogReader;
private RelDataTypeFactory typeFactory;
private RelOptPlanner calcitePlanner;

/** Called by {@link BeamSqlEnv}.instantiatePlanner() reflectively. */
public CalciteQueryPlanner(JdbcConnection connection, Collection<RuleSet> ruleSets) {
this.connection = connection;
this.planner = Frameworks.getPlanner(defaultConfig(connection, ruleSets));
this.config = defaultConfig(connection, ruleSets);
this.planner = Frameworks.getPlanner(config);

Frameworks.withPlanner(
(cluster, relOptSchema, rootSchema) -> {
// CAPTURE THE COMPONENTS HERE
this.relOptCluster = cluster;
this.catalogReader = (CalciteCatalogReader) relOptSchema;
this.typeFactory = cluster.getTypeFactory();
this.calcitePlanner = cluster.getPlanner();

// ... any other setup from the original lambda ...
// e.g., planner.setExecutor(executor);

return null;
},
config);

if (this.relOptCluster == null || this.catalogReader == null) {
throw new IllegalStateException("Failed to initialize Calcite components");
}
}

/**
* Returns a RelBuilder instance configured with the same Calcite components used by this
* QueryPlanner.
*/
@Override
public RelBuilder getRelBuilder() {
return RelBuilder.create(config);
}

@Override
public SqlOperatorTable getOperatorTable() {
return config.getOperatorTable();
}

public static final Factory FACTORY =
Expand All @@ -103,6 +152,7 @@ public CalciteQueryPlanner(JdbcConnection connection, Collection<RuleSet> ruleSe
public QueryPlanner createPlanner(
JdbcConnection jdbcConnection, Collection<RuleSet> ruleSets) {
loadBuiltinFunctions(jdbcConnection);
LOG.info("Factory creating planner with ruleSets: {}", ruleSets);
return new CalciteQueryPlanner(jdbcConnection, ruleSets);
}

Expand All @@ -120,12 +170,20 @@ private void loadBuiltinFunctions(JdbcConnection jdbcConnection) {

public FrameworkConfig defaultConfig(JdbcConnection connection, Collection<RuleSet> ruleSets) {
final CalciteConnectionConfig config = connection.config();
// Resolve the parser conformance. Calcite's Avatica JDBC connect path silently drops the
// {@code conformance} connection property (it is not in the driver's registered property set),
// so {@code config.conformance()} is always DEFAULT here even when callers set it via
// {@code BeamSqlPipelineOptions.calciteConnectionProperties}. We therefore read that map
// directly from the connection's pipeline options and let it override. This keeps the behavior
// opt-in: with no {@code conformance} property the connection's own (DEFAULT) value is used, so
// existing Beam SQL behavior is unchanged.
final SqlConformance conformance = resolveConformance(connection, config);
final SqlParser.ConfigBuilder parserConfig =
SqlParser.configBuilder()
.setQuotedCasing(config.quotedCasing())
.setUnquotedCasing(config.unquotedCasing())
.setQuoting(config.quoting())
.setConformance(config.conformance())
.setConformance(conformance)
.setCaseSensitive(config.caseSensitive());
final SqlParserImplFactory parserFactory =
config.parserFactory(SqlParserImplFactory.class, null);
Expand All @@ -150,6 +208,7 @@ public FrameworkConfig defaultConfig(JdbcConnection connection, Collection<RuleS
// Revert the flag flip of CALCITE-3870 which led to missing rules
SqlToRelConverter.Config sqlToRelConfig = SqlToRelConverter.config().withExpand(true);

LOG.info("Creating config with rulesets: {}", ruleSets);
return Frameworks.newConfigBuilder()
.parserConfig(parserConfig.build())
.defaultSchema(defaultSchema)
Expand All @@ -163,6 +222,43 @@ public FrameworkConfig defaultConfig(JdbcConnection connection, Collection<RuleS
.build();
}

/**
* Resolves the {@link SqlConformance} for the parser. Prefers an explicit {@code conformance}
* entry in {@link BeamSqlPipelineOptions#getCalciteConnectionProperties()} (looked up
* case-insensitively, value matched against {@link SqlConformanceEnum}); otherwise falls back to
* the connection's own conformance. This is the bridge for {@code conformance=BABEL}, which the
* Avatica JDBC connect path drops, enabling Spark-SQL spellings (e.g. the {@code !=} operator)
* the default conformance rejects. Returns the connection default on any unrecognized value.
*/
private static SqlConformance resolveConformance(
JdbcConnection connection, CalciteConnectionConfig config) {
PipelineOptions options = connection.getPipelineOptions();
if (options == null) {
return config.conformance();
}
BeamSqlPipelineOptions sqlOptions = options.as(BeamSqlPipelineOptions.class);
Map<String, String> props = sqlOptions.getCalciteConnectionProperties();
if (props == null) {
return config.conformance();
}
String value = null;
for (Map.Entry<String, String> e : props.entrySet()) {
if ("conformance".equalsIgnoreCase(e.getKey())) {
value = e.getValue();
break;
}
}
if (value == null) {
return config.conformance();
}
try {
return SqlConformanceEnum.valueOf(value.trim().toUpperCase());
} catch (IllegalArgumentException ex) {
LOG.warn("Unrecognized calcite conformance '{}', using {}", value, config.conformance());
return config.conformance();
}
}

/** Parse input SQL query, and return a {@link SqlNode} as grammar tree. */
@Override
public SqlNode parse(String sqlStatement) throws ParseException {
Expand All @@ -179,15 +275,14 @@ public SqlNode parse(String sqlStatement) throws ParseException {

/**
* It parses and validate the input query, then convert into a {@link BeamRelNode} tree. Note that
* query parameters are not yet supported.
* query parameters are now supported for positional parameters.
*/
@Override
public BeamRelNode convertToBeamRel(String sqlStatement, QueryParameters queryParameters)
throws ParseException, SqlConversionException {
Preconditions.checkArgument(
queryParameters.getKind() == Kind.NONE || queryParameters.getKind() == Kind.POSITIONAL,
"Beam SQL Calcite dialect only supports positional query parameters.");
BeamRelNode beamRelNode;
"Beam SQL Calcite dialect only supports positional query parameters or no parameters.");
try {
SqlNode parsed = planner.parse(sqlStatement);
TableResolutionUtils.setupCustomTableResolution(connection, parsed);
Expand All @@ -203,12 +298,59 @@ public BeamRelNode convertToBeamRel(String sqlStatement, QueryParameters queryPa
relNode,
new ParameterBinder(root.rel.getCluster().getRexBuilder(), queryParameters));
}
LOG.info("SQLPlan>\n{}", BeamSqlRelUtils.explainLazily(root.rel));
return convertToBeamRel(relNode, queryParameters);
} catch (RelConversionException | CannotPlanException e) {
throw new SqlConversionException(
String.format("Unable to convert query %s", sqlStatement), e);
} catch (SqlParseException | ValidationException e) {
throw new ParseException(String.format("Unable to parse query %s", sqlStatement), e);
} finally {
planner.close();
}
}

private static RelNode bindParameters(RelNode rel, RexShuttle binder) {
RelNode newRel = rel.accept(binder);
java.util.List<RelNode> newInputs = new java.util.ArrayList<>();
for (RelNode input : newRel.getInputs()) {
newInputs.add(bindParameters(input, binder));
}
return newRel.copy(newRel.getTraitSet(), newInputs);
}

@Override
public RelNode parseToRel(String sqlStatement, QueryParameters queryParameters)
throws ParseException, SqlConversionException {
Preconditions.checkArgument(
queryParameters.getKind() == Kind.NONE,
"Beam SQL Calcite dialect does not yet support query parameters.");
try {
SqlNode parsed = planner.parse(sqlStatement);
TableResolutionUtils.setupCustomTableResolution(connection, parsed);
SqlNode validated = planner.validate(parsed);
// root of original logical plan
RelRoot root = planner.rel(validated);
return root.rel;
} catch (RelConversionException e) {
throw new SqlConversionException(
String.format("Unable to convert query %s", sqlStatement), e);
} catch (SqlParseException | ValidationException e) {
throw new ParseException(String.format("Unable to parse query %s", sqlStatement), e);
} finally {
planner.close();
}
}

@Override
public BeamRelNode convertToBeamRel(RelNode relNode, QueryParameters queryParameters) {
RelNode beamRelNode;
try {
LOG.info("SQLPlan>\n{}", BeamSqlRelUtils.explainLazily(relNode));
RelTraitSet desiredTraits =
relNode
.getTraitSet()
.replace(BeamLogicalConvention.INSTANCE)
.replace(root.collation)
// .replace(root.collation)
.simplify();
// beam physical plan
relNode
Expand All @@ -224,26 +366,23 @@ public BeamRelNode convertToBeamRel(String sqlStatement, QueryParameters queryPa
RelMetadataQuery.THREAD_PROVIDERS.set(
JaninoRelMetadataProvider.of(relNode.getCluster().getMetadataProvider()));
relNode.getCluster().invalidateMetadataQuery();
beamRelNode = (BeamRelNode) planner.transform(0, desiredTraits, relNode);
Program program = config.getPrograms().get(0);
LOG.info("Desired traits: {}", desiredTraits);
beamRelNode =
program.run(
relNode.getCluster().getPlanner(),
relNode,
desiredTraits,
ImmutableList.of(),
ImmutableList.of());
LOG.info("BEAMPlan>\n{}", BeamSqlRelUtils.explainLazily(beamRelNode));
} catch (RelConversionException | CannotPlanException e) {
} catch (CannotPlanException e) {
throw new SqlConversionException(
String.format("Unable to convert query %s", sqlStatement), e);
} catch (SqlParseException | ValidationException e) {
throw new ParseException(String.format("Unable to parse query %s", sqlStatement), e);
String.format("Unable to convert relNode to Beam: %s", relNode), e);
} finally {
planner.close();
}
return beamRelNode;
}

private static RelNode bindParameters(RelNode rel, RexShuttle binder) {
RelNode newRel = rel.accept(binder);
java.util.List<RelNode> newInputs = new java.util.ArrayList<>();
for (RelNode input : newRel.getInputs()) {
newInputs.add(bindParameters(input, binder));
}
return newRel.copy(newRel.getTraitSet(), newInputs);
return (BeamRelNode) beamRelNode;
}

// It needs to be public so that the generated code in Calcite can access it.
Expand Down
Loading
Loading