From 262eea27a8c65119a91e664877983e7360f20cec Mon Sep 17 00:00:00 2001 From: Danny McCormick Date: Fri, 12 Jun 2026 19:36:43 +0000 Subject: [PATCH 1/3] SQL Database DDL & Usability Improvements: Support CREATE/DROP DATABASE and USE short syntax --- .../src/main/codegen/includes/parserImpls.ftl | 4 +-- .../sql/impl/parser/SqlDdlNodes.java | 4 +-- .../sql/meta/catalog/InMemoryCatalog.java | 7 ++-- .../sql/BeamSqlCliDatabaseTest.java | 36 +++++++++++++++++++ 4 files changed, 45 insertions(+), 6 deletions(-) diff --git a/sdks/java/extensions/sql/src/main/codegen/includes/parserImpls.ftl b/sdks/java/extensions/sql/src/main/codegen/includes/parserImpls.ftl index 94c0161c492c..cb8eec438728 100644 --- a/sdks/java/extensions/sql/src/main/codegen/includes/parserImpls.ftl +++ b/sdks/java/extensions/sql/src/main/codegen/includes/parserImpls.ftl @@ -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) : { @@ -391,7 +391,7 @@ SqlCall SqlUseDatabase(Span s, String scope) : { s.add(this); } - + [ ] databaseName = CompoundIdentifier() { return new SqlUseDatabase( diff --git a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/parser/SqlDdlNodes.java b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/parser/SqlDdlNodes.java index 6d6be5d5a127..43b6d86d188a 100644 --- a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/parser/SqlDdlNodes.java +++ b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/parser/SqlDdlNodes.java @@ -55,7 +55,7 @@ public static SqlNode column( } /** Returns the schema in which to create an object. */ - static Pair schema( + public static Pair schema( CalcitePrepare.Context context, boolean mutable, SqlIdentifier id) { CalciteSchema rootSchema = mutable ? context.getMutableRootSchema() : context.getRootSchema(); @Nullable CalciteSchema schema = null; @@ -72,7 +72,7 @@ static Pair schema( return Pair.of(checkStateNotNull(schema, "Got null sub-schema for path '%s'", path), name(id)); } - private static @Nullable CalciteSchema childSchema(CalciteSchema rootSchema, List path) { + public static @Nullable CalciteSchema childSchema(CalciteSchema rootSchema, List path) { @Nullable CalciteSchema schema = rootSchema; for (String p : path) { if (schema == null) { diff --git a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/catalog/InMemoryCatalog.java b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/catalog/InMemoryCatalog.java index cdee6c930224..5efc3010b9c1 100644 --- a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/catalog/InMemoryCatalog.java +++ b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/catalog/InMemoryCatalog.java @@ -18,7 +18,6 @@ package org.apache.beam.sdk.extensions.sql.meta.catalog; import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument; -import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState; import java.util.Collection; import java.util.Collections; @@ -111,12 +110,16 @@ public Collection databases() { @Override public boolean dropDatabase(String database, boolean cascade) { - checkState(!cascade, "%s does not support CASCADE.", getClass().getSimpleName()); + MetaStore metaStore = metaStores.get(database); + if (!cascade && metaStore != null && !metaStore.getTables().isEmpty()) { + throw new IllegalStateException("Database '" + database + "' is not empty."); + } boolean removed = databases.remove(database); if (database.equals(currentDatabase)) { currentDatabase = null; } + metaStores.remove(database); return removed; } diff --git a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java index 588caa78a2b7..2905f8efcaad 100644 --- a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java +++ b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java @@ -83,6 +83,15 @@ public void testUseDatabase() { assertEquals("my_database2", catalogManager.currentCatalog().currentDatabase()); } + @Test + public void testUseDatabaseWithoutDatabaseKeyword() { + assertEquals(DEFAULT, catalogManager.currentCatalog().currentDatabase()); + cli.execute("CREATE DATABASE my_database"); + assertEquals(DEFAULT, catalogManager.currentCatalog().currentDatabase()); + cli.execute("USE my_database"); + assertEquals("my_database", catalogManager.currentCatalog().currentDatabase()); + } + @Test public void testUseDatabase_doesNotExist() { assertEquals(DEFAULT, catalogManager.currentCatalog().currentDatabase()); @@ -126,6 +135,33 @@ public void testDropDatabase_nonexistent() { cli.execute("DROP DATABASE my_database"); } + @Test + public void testDropDatabase_notEmpty_restrict() { + cli.execute("CREATE DATABASE db_1"); + cli.execute("USE db_1"); + + TestTableProvider testTableProvider = new TestTableProvider(); + catalogManager.registerTableProvider(testTableProvider); + cli.execute("CREATE EXTERNAL TABLE person(id int, name varchar, age int) TYPE 'test'"); + + thrown.expect(RuntimeException.class); + thrown.expectMessage("Database 'db_1' is not empty."); + cli.execute("DROP DATABASE db_1"); + } + + @Test + public void testDropDatabase_notEmpty_cascade() { + cli.execute("CREATE DATABASE db_1"); + cli.execute("USE db_1"); + + TestTableProvider testTableProvider = new TestTableProvider(); + catalogManager.registerTableProvider(testTableProvider); + cli.execute("CREATE EXTERNAL TABLE person(id int, name varchar, age int) TYPE 'test'"); + + cli.execute("DROP DATABASE db_1 CASCADE"); + assertFalse(catalogManager.currentCatalog().databaseExists("db_1")); + } + @Test public void testCreateInsertDropTableUsingDefaultDatabase() { Catalog catalog = catalogManager.currentCatalog(); From 0a4871537ce6906af6a8c7790b0a8ae6db7ce6fd Mon Sep 17 00:00:00 2001 From: Danny McCormick Date: Fri, 12 Jun 2026 20:17:23 +0000 Subject: [PATCH 2/3] Address code review feedback: add guard in InMemoryCatalog.dropDatabase and add tests for DROP DATABASE IF EXISTS --- .../sql/meta/catalog/InMemoryCatalog.java | 3 +++ .../extensions/sql/BeamSqlCliDatabaseTest.java | 16 ++++++++++++++++ 2 files changed, 19 insertions(+) diff --git a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/catalog/InMemoryCatalog.java b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/catalog/InMemoryCatalog.java index 5efc3010b9c1..647de314ab10 100644 --- a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/catalog/InMemoryCatalog.java +++ b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/catalog/InMemoryCatalog.java @@ -110,6 +110,9 @@ public Collection databases() { @Override public boolean dropDatabase(String database, boolean cascade) { + if (!databases.contains(database)) { + return false; + } MetaStore metaStore = metaStores.get(database); if (!cascade && metaStore != null && !metaStore.getTables().isEmpty()) { throw new IllegalStateException("Database '" + database + "' is not empty."); diff --git a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java index 2905f8efcaad..4378ec59a306 100644 --- a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java +++ b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java @@ -135,6 +135,22 @@ public void testDropDatabase_nonexistent() { cli.execute("DROP DATABASE my_database"); } + @Test + public void testDropDatabase_ifExists_nonexistent() { + assertFalse(catalogManager.currentCatalog().databaseExists("my_database")); + // Should not throw exception + cli.execute("DROP DATABASE IF EXISTS my_database"); + assertFalse(catalogManager.currentCatalog().databaseExists("my_database")); + } + + @Test + public void testDropDatabase_ifExists_exists() { + cli.execute("CREATE DATABASE my_database"); + assertTrue(catalogManager.currentCatalog().databaseExists("my_database")); + cli.execute("DROP DATABASE IF EXISTS my_database"); + assertFalse(catalogManager.currentCatalog().databaseExists("my_database")); + } + @Test public void testDropDatabase_notEmpty_restrict() { cli.execute("CREATE DATABASE db_1"); From 4314fd89913e76a3dbbdd2292015fbfc1849207e Mon Sep 17 00:00:00 2001 From: Danny McCormick Date: Fri, 12 Jun 2026 20:25:12 +0000 Subject: [PATCH 3/3] Address new code review feedback: optimize InMemoryCatalog, add null check in SqlDdlNodes, and make tests more precise --- .../beam/sdk/extensions/sql/impl/parser/SqlDdlNodes.java | 3 +++ .../sdk/extensions/sql/meta/catalog/InMemoryCatalog.java | 8 ++++---- .../beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java | 2 +- 3 files changed, 8 insertions(+), 5 deletions(-) diff --git a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/parser/SqlDdlNodes.java b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/parser/SqlDdlNodes.java index 43b6d86d188a..f8d7e6f73851 100644 --- a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/parser/SqlDdlNodes.java +++ b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/parser/SqlDdlNodes.java @@ -73,6 +73,9 @@ public static Pair schema( } public static @Nullable CalciteSchema childSchema(CalciteSchema rootSchema, List path) { + if (path == null) { + return null; + } @Nullable CalciteSchema schema = rootSchema; for (String p : path) { if (schema == null) { diff --git a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/catalog/InMemoryCatalog.java b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/catalog/InMemoryCatalog.java index 647de314ab10..68e80c2340fa 100644 --- a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/catalog/InMemoryCatalog.java +++ b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/catalog/InMemoryCatalog.java @@ -110,20 +110,20 @@ public Collection databases() { @Override public boolean dropDatabase(String database, boolean cascade) { - if (!databases.contains(database)) { - return false; - } MetaStore metaStore = metaStores.get(database); if (!cascade && metaStore != null && !metaStore.getTables().isEmpty()) { throw new IllegalStateException("Database '" + database + "' is not empty."); } boolean removed = databases.remove(database); + if (!removed) { + return false; + } if (database.equals(currentDatabase)) { currentDatabase = null; } metaStores.remove(database); - return removed; + return true; } @Override diff --git a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java index 4378ec59a306..54682911fe11 100644 --- a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java +++ b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlCliDatabaseTest.java @@ -160,7 +160,7 @@ public void testDropDatabase_notEmpty_restrict() { catalogManager.registerTableProvider(testTableProvider); cli.execute("CREATE EXTERNAL TABLE person(id int, name varchar, age int) TYPE 'test'"); - thrown.expect(RuntimeException.class); + thrown.expect(CalciteContextException.class); thrown.expectMessage("Database 'db_1' is not empty."); cli.execute("DROP DATABASE db_1"); }