Skip to content

Commit 28b8aa4

Browse files
committed
SQL API Extensions, Usability Improvements, and Subquery Decorrelation
1 parent d521df8 commit 28b8aa4

9 files changed

Lines changed: 320 additions & 28 deletions

File tree

sdks/java/extensions/sql/src/main/codegen/includes/parserImpls.ftl

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -381,7 +381,7 @@ SqlCreate SqlCreateDatabase(Span s, boolean replace) :
381381
}
382382

383383
/**
384-
* USE DATABASE ( catalog_name '.' )? database_name
384+
* USE [ DATABASE ] ( catalog_name '.' )? database_name
385385
*/
386386
SqlCall SqlUseDatabase(Span s, String scope) :
387387
{
@@ -391,7 +391,7 @@ SqlCall SqlUseDatabase(Span s, String scope) :
391391
<USE> {
392392
s.add(this);
393393
}
394-
<DATABASE>
394+
[ <DATABASE> ]
395395
databaseName = CompoundIdentifier()
396396
{
397397
return new SqlUseDatabase(

sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/BeamCalciteTable.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,7 @@ public class BeamCalciteTable extends AbstractQueryableTable
5353
private final Map<String, String> pipelineOptionsMap;
5454
private @Nullable PipelineOptions pipelineOptions;
5555

56-
BeamCalciteTable(
56+
public BeamCalciteTable(
5757
BeamSqlTable beamTable,
5858
Map<String, String> pipelineOptionsMap,
5959
@Nullable PipelineOptions pipelineOptions) {

sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/BeamSqlEnv.java

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,17 +47,23 @@
4747
import org.apache.beam.sdk.transforms.SerializableFunction;
4848
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.jdbc.CalcitePrepare;
4949
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.plan.RelOptUtil;
50+
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.rel.RelNode;
5051
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.schema.Function;
5152
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.sql.SqlKind;
53+
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.tools.RelBuilder;
5254
import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.tools.RuleSet;
5355
import org.checkerframework.checker.nullness.qual.Nullable;
56+
import org.slf4j.Logger;
57+
import org.slf4j.LoggerFactory;
5458

5559
/**
5660
* Contains the metadata of tables/UDF functions, and exposes APIs to
5761
* query/validate/optimize/translate SQL statements.
5862
*/
5963
@Internal
6064
public class BeamSqlEnv {
65+
private static final Logger LOG = LoggerFactory.getLogger(BeamSqlEnv.class);
66+
6167
JdbcConnection connection;
6268
QueryPlanner planner;
6369

@@ -116,6 +122,31 @@ public BeamRelNode parseQuery(String query, QueryParameters queryParameters)
116122
return planner.convertToBeamRel(query, queryParameters);
117123
}
118124

125+
public QueryPlanner getPlanner() {
126+
return planner;
127+
}
128+
129+
public RelBuilder getRelBuilder() {
130+
return planner.getRelBuilder();
131+
}
132+
133+
public BeamRelNode convertToBeamRel(RelNode relNode) {
134+
return planner.convertToBeamRel(relNode, QueryParameters.ofNone());
135+
}
136+
137+
public RelNode parseLogicalPlan(String query) throws ParseException {
138+
return planner.parseToRel(query, QueryParameters.ofNone());
139+
}
140+
141+
public void registerSchemaFunction(String name, Function function) {
142+
connection.getCurrentSchemaPlus().add(name, function);
143+
}
144+
145+
public org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.sql.SqlOperatorTable
146+
getOperatorTable() {
147+
return planner.getOperatorTable();
148+
}
149+
119150
public boolean isDdl(String sqlStatement) throws ParseException {
120151
return planner.parse(sqlStatement).getKind().belongsTo(SqlKind.DDL);
121152
}
@@ -196,6 +227,7 @@ public BeamSqlEnvBuilder setCurrentSchema(String name) {
196227

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

263295
configureSchemas(jdbcConnection);
264296

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

267300
// The planner may choose to add its own builtin functions to the schema, so load user-defined

0 commit comments

Comments
 (0)