fmorillo7694 commented on code in PR #206: URL: https://github.com/apache/flink-connector-aws/pull/206#discussion_r4125206267
########## flink-catalog-aws/flink-catalog-aws-glue/src/test/java/org/apache/flink/table/catalog/glue/util/GlueTableUtilsTest.java: ########## @@ -0,0 +1,362 @@ +/* + * 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.flink.table.catalog.glue.util; + +import org.apache.flink.table.api.DataTypes; +import org.apache.flink.table.api.Schema; +import org.apache.flink.table.catalog.ObjectPath; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import software.amazon.awssdk.services.glue.model.Column; +import software.amazon.awssdk.services.glue.model.StorageDescriptor; +import software.amazon.awssdk.services.glue.model.Table; + +import java.util.Arrays; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +/** + * Unit tests for the GlueTableUtils class. Tests the utility methods for working with AWS Glue + * tables. + */ +class GlueTableUtilsTest { + + private GlueTypeConverter glueTypeConverter; + private GlueTableUtils glueTableUtils; + + // Test data + private static final String TEST_CONNECTOR_TYPE = "kinesis"; + private static final String TEST_TABLE_LOCATION = "arn://..."; + private static final String TEST_TABLE_NAME = "test_table"; + private static final String TEST_COLUMN_NAME = "test_column"; + + @BeforeEach + void setUp() { + // Initialize GlueTypeConverter directly as it is already implemented + glueTypeConverter = new GlueTypeConverter(); + glueTableUtils = new GlueTableUtils(glueTypeConverter); + } + + @Test + void testBuildStorageDescriptor() { + // Prepare test data + List<Column> glueColumns = + Arrays.asList(Column.builder().name(TEST_COLUMN_NAME).type("string").build()); + + // Build the StorageDescriptor + StorageDescriptor storageDescriptor = + glueTableUtils.buildStorageDescriptor( + new HashMap<>(), glueColumns, TEST_TABLE_LOCATION); + + // Assert that the StorageDescriptor is not null and contains the correct location + Assertions.assertNotNull(storageDescriptor, "StorageDescriptor should not be null"); + Assertions.assertEquals( + TEST_TABLE_LOCATION, storageDescriptor.location(), "Table location should match"); + Assertions.assertEquals( + 1, storageDescriptor.columns().size(), "StorageDescriptor should have one column"); + Assertions.assertEquals( + TEST_COLUMN_NAME, + storageDescriptor.columns().get(0).name(), + "Column name should match"); + } + + @Test + void testExtractTableLocationWithLocationKey() { + // Prepare table properties with a connector type and location + Map<String, String> tableProperties = new HashMap<>(); + tableProperties.put("connector", TEST_CONNECTOR_TYPE); + tableProperties.put( + "stream.arn", TEST_TABLE_LOCATION); // Mimicking a location key for kinesis + + ObjectPath tablePath = new ObjectPath("test_database", TEST_TABLE_NAME); + + // Extract table location + String location = glueTableUtils.extractTableLocation(tableProperties, tablePath); + + // Assert that the correct location is used + Assertions.assertEquals( + TEST_TABLE_LOCATION, location, "Table location should match the location key"); + } + + @Test + void testExtractTableLocationWithDefaultLocation() { + // Prepare table properties without a location key + Map<String, String> tableProperties = new HashMap<>(); + tableProperties.put("connector", TEST_CONNECTOR_TYPE); // No actual location key here + + ObjectPath tablePath = new ObjectPath("test_database", TEST_TABLE_NAME); + + // Extract table location + String location = glueTableUtils.extractTableLocation(tableProperties, tablePath); + + // Assert that the default location is used + String expectedLocation = + tablePath.getDatabaseName() + "/tables/" + tablePath.getObjectName(); + Assertions.assertEquals(expectedLocation, location, "Default location should be used"); + } + + @Test + void testMapFlinkColumnToGlueColumn() { + // Prepare a Flink column to convert + org.apache.flink.table.catalog.Column flinkColumn = + org.apache.flink.table.catalog.Column.physical( + TEST_COLUMN_NAME, + DataTypes.STRING() // Fix: DataTypes.STRING() instead of DataType.STRING() + ); + + // Convert Flink column to Glue column + Column glueColumn = glueTableUtils.mapFlinkColumnToGlueColumn(flinkColumn); + + // Assert that the Glue column is correctly mapped + Assertions.assertNotNull(glueColumn, "Converted Glue column should not be null"); + Assertions.assertEquals( + TEST_COLUMN_NAME, + glueColumn.name(), + "Column name should be preserved as declared (no lowercasing)"); + Assertions.assertEquals( + "string", glueColumn.type(), "Column type should match the expected Glue type"); + } + + @Test + void testGetSchemaFromGlueTable() { + // Prepare a Glue table with columns + List<Column> glueColumns = + Arrays.asList( + Column.builder().name(TEST_COLUMN_NAME).type("string").build(), + Column.builder().name("another_column").type("int").build()); + StorageDescriptor storageDescriptor = + StorageDescriptor.builder().columns(glueColumns).build(); + Table glueTable = Table.builder().storageDescriptor(storageDescriptor).build(); + + // Get the schema from the Glue table + Schema schema = glueTableUtils.getSchemaFromGlueTable(glueTable); + + // Assert that the schema is correctly constructed + Assertions.assertNotNull(schema, "Schema should not be null"); + Assertions.assertEquals(2, schema.getColumns().size(), "Schema should have two columns"); + } + + @Test + void testColumnNameCaseSensitivity() { + // 1. Define Flink columns with mixed case names + org.apache.flink.table.catalog.Column upperCaseColumn = + org.apache.flink.table.catalog.Column.physical( + "UpperCaseColumn", DataTypes.STRING()); + + org.apache.flink.table.catalog.Column mixedCaseColumn = + org.apache.flink.table.catalog.Column.physical("mixedCaseColumn", DataTypes.INT()); + + org.apache.flink.table.catalog.Column lowerCaseColumn = + org.apache.flink.table.catalog.Column.physical( + "lowercase_column", DataTypes.BOOLEAN()); + + // 2. Convert Flink columns to Glue columns + Column glueUpperCase = glueTableUtils.mapFlinkColumnToGlueColumn(upperCaseColumn); + Column glueMixedCase = glueTableUtils.mapFlinkColumnToGlueColumn(mixedCaseColumn); + Column glueLowerCase = glueTableUtils.mapFlinkColumnToGlueColumn(lowerCaseColumn); + + // 3. Verify Glue column names are stored lowercase (Glue lowercases on write, so we + // store lowercase deterministically) with the original case in the parameters. + Assertions.assertEquals( + "uppercasecolumn", + glueUpperCase.name(), + "Glue column name should be stored lowercase"); + Assertions.assertEquals( + "mixedcasecolumn", + glueMixedCase.name(), + "Glue column name should be stored lowercase"); + Assertions.assertEquals( + "lowercase_column", + glueLowerCase.name(), + "Glue column name should be stored lowercase"); + + // 4. Verify the originalName parameter carries the declared case (only when needed) + Assertions.assertEquals( + "UpperCaseColumn", + glueUpperCase.parameters().get("originalName"), + "originalName parameter should preserve the declared case"); + Assertions.assertEquals( + "mixedCaseColumn", + glueMixedCase.parameters().get("originalName"), + "originalName parameter should preserve the declared case"); + Assertions.assertFalse( + glueLowerCase.parameters() != null + && glueLowerCase.parameters().containsKey("originalName"), + "already-lowercase columns need no originalName parameter"); + + // 5. Create a Glue table with these columns + List<Column> glueColumns = Arrays.asList(glueUpperCase, glueMixedCase, glueLowerCase); + StorageDescriptor storageDescriptor = + StorageDescriptor.builder().columns(glueColumns).build(); + Table glueTable = Table.builder().storageDescriptor(storageDescriptor).build(); + + // 6. Convert back to Flink schema + Schema schema = glueTableUtils.getSchemaFromGlueTable(glueTable); + + // 7. Verify that original case is preserved in schema + List<String> columnNames = + schema.getColumns().stream().map(col -> col.getName()).collect(Collectors.toList()); + + Assertions.assertEquals(3, columnNames.size(), "Schema should have three columns"); + Assertions.assertTrue( + columnNames.contains("UpperCaseColumn"), + "Schema should contain the uppercase column with original case"); + Assertions.assertTrue( + columnNames.contains("mixedCaseColumn"), + "Schema should contain the mixed case column with original case"); + Assertions.assertTrue( + columnNames.contains("lowercase_column"), + "Schema should contain the lowercase column with original case"); + } + + @Test + void testEndToEndColumnNameCasePreservation() { + // This test simulates a more complete lifecycle with table creation and JSON parsing + + // 1. Create Flink columns with mixed case (representing original source) + List<org.apache.flink.table.catalog.Column> flinkColumns = + Arrays.asList( + org.apache.flink.table.catalog.Column.physical("ID", DataTypes.INT()), + org.apache.flink.table.catalog.Column.physical( + "UserName", DataTypes.STRING()), + org.apache.flink.table.catalog.Column.physical( + "timestamp", DataTypes.TIMESTAMP()), + org.apache.flink.table.catalog.Column.physical( + "DATA_VALUE", DataTypes.STRING())); + + // 2. Convert to Glue columns (simulating what happens in table creation) + List<Column> glueColumns = + flinkColumns.stream() + .map(glueTableUtils::mapFlinkColumnToGlueColumn) + .collect(Collectors.toList()); + + // 3. Verify Glue columns are stored lowercase (real Glue lowercases on write) with the + // declared case preserved via the originalName parameter. + for (int i = 0; i < flinkColumns.size(); i++) { + String originalName = flinkColumns.get(i).getName(); + Column glueColumn = glueColumns.get(i); + + Assertions.assertEquals( + originalName.toLowerCase(), + glueColumn.name(), + "Glue column name should be stored lowercase"); + if (!originalName.equals(originalName.toLowerCase())) { + Assertions.assertEquals( + originalName, + glueColumn.parameters().get("originalName"), + "originalName parameter should preserve the declared case"); + } + } + + // 4. Create a Glue table with these columns (simulating storage in Glue) + StorageDescriptor storageDescriptor = + StorageDescriptor.builder().columns(glueColumns).build(); + Table glueTable = Table.builder().storageDescriptor(storageDescriptor).build(); + + // 5. Convert back to Flink schema (simulating table retrieval for queries) + Schema schema = glueTableUtils.getSchemaFromGlueTable(glueTable); + + // 6. Verify original case is preserved in the resulting schema + List<String> resultColumnNames = + schema.getColumns().stream().map(col -> col.getName()).collect(Collectors.toList()); + + for (org.apache.flink.table.catalog.Column originalColumn : flinkColumns) { + String originalName = originalColumn.getName(); + Assertions.assertTrue( + resultColumnNames.contains(originalName), + "Result schema should contain original column name with case preserved: " + + originalName); + } + + // 7. Verify that a JSON string matching the original schema can be parsed correctly + // This is a simulation of the real-world scenario where properly cased column names + // are needed for JSON parsing + String jsonExample = + "{\"ID\":1,\"UserName\":\"test\",\"timestamp\":\"2023-01-01 12:00:00\",\"DATA_VALUE\":\"sample\"}"; + + // We don't actually parse the JSON here since that would require external dependencies, + // but this illustrates the scenario where correct case is important + + Assertions.assertEquals( + "ID", resultColumnNames.get(0), "First column should maintain original case"); + Assertions.assertEquals( + "UserName", + resultColumnNames.get(1), + "Second column should maintain original case"); + Assertions.assertEquals( + "timestamp", + resultColumnNames.get(2), + "Third column should maintain original case"); + Assertions.assertEquals( + "DATA_VALUE", + resultColumnNames.get(3), + "Fourth column should maintain original case"); + } + + @Test + void testGetSchemaFromGlueTableHonorsLegacyOriginalNameParameter() { + // Tables written by older catalog versions carry lowercased names plus an + // "originalName" column parameter; reads must still surface the declared name. + Column legacyColumn = + Column.builder() + .name("username") + .type("string") + .parameters(java.util.Collections.singletonMap("originalName", "UserName")) + .build(); + StorageDescriptor storageDescriptor = + StorageDescriptor.builder().columns(Arrays.asList(legacyColumn)).build(); + Table glueTable = Table.builder().storageDescriptor(storageDescriptor).build(); + + Schema schema = glueTableUtils.getSchemaFromGlueTable(glueTable); + + Assertions.assertEquals(1, schema.getColumns().size(), "Schema should have one column"); + Assertions.assertEquals( + "UserName", + schema.getColumns().get(0).getName(), + "Legacy originalName parameter should still be honored on read"); + } + + @Test + void testGetSchemaFromGlueTableIncludesPartitionColumns() { Review Comment: Rewritten with the name-keyed design; the unit test now builds the Glue-order column map (SD columns then partition key) and asserts the declared order and types come back. ########## flink-catalog-aws/flink-catalog-aws-glue/src/test/java/org/apache/flink/table/catalog/glue/util/RealGlueCleanupExtension.java: ########## @@ -0,0 +1,87 @@ +/* + * 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.flink.table.catalog.glue.util; + +import org.junit.jupiter.api.extension.AfterEachCallback; +import org.junit.jupiter.api.extension.BeforeEachCallback; +import org.junit.jupiter.api.extension.ExtensionContext; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.Database; +import software.amazon.awssdk.services.glue.model.EntityNotFoundException; +import software.amazon.awssdk.services.glue.model.GetDatabasesResponse; +import software.amazon.awssdk.services.glue.model.Table; + +import java.util.HashSet; +import java.util.Set; + +/** + * Keeps a real AWS Glue account clean when the catalog test suites run against the real service + * (see {@link GlueTestClientFactory}): before each test it snapshots the account's database names, + * and after the test it deletes any database (and its tables) that the test created. The + * fake-backed default mode is a no-op. + * + * <p>The delta approach means pre-existing databases in the account are never touched, and tests + * keep their fixed database names without colliding across tests. + */ +public class RealGlueCleanupExtension implements BeforeEachCallback, AfterEachCallback { + + private GlueClient cleanupClient; Review Comment: Yes — it now implements `AfterAllCallback` and closes the client. While in there I found the real cause of the "Database does not exist" flakes we had been attributing to Glue eventual consistency: the extension deleted every database that appeared in the account during a test, but surefire runs the suites in **four forks against the same account**, so one fork's cleanup deleted another fork's live database mid-test. `uniqueName` now carries a per-JVM token and registers the name; the extension deletes only names this JVM handed out. The real-Glue run went from 10 flakes to 0. ########## flink-catalog-aws/flink-catalog-aws-glue/src/test/java/org/apache/flink/table/catalog/glue/operator/GlueDatabaseOperationsTest.java: ########## @@ -0,0 +1,336 @@ +/* + * 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.flink.table.catalog.glue.operator; + +import org.apache.flink.table.catalog.CatalogDatabase; +import org.apache.flink.table.catalog.CatalogDatabaseImpl; +import org.apache.flink.table.catalog.exceptions.CatalogException; +import org.apache.flink.table.catalog.exceptions.DatabaseAlreadyExistException; +import org.apache.flink.table.catalog.exceptions.DatabaseNotExistException; +import org.apache.flink.table.catalog.glue.util.GlueTestClientFactory; +import org.apache.flink.table.catalog.glue.util.RealGlueCleanupExtension; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.InvalidInputException; +import software.amazon.awssdk.services.glue.model.OperationTimeoutException; +import software.amazon.awssdk.services.glue.model.ResourceNumberLimitExceededException; + +import java.util.Collections; +import java.util.List; + +import static org.assertj.core.api.Assumptions.assumeThat; + +/** + * Unit tests for the GlueDatabaseOperations class. These tests verify the functionality for + * database operations such as create, drop, get, and list in the AWS Glue service. + */ +@ExtendWith(RealGlueCleanupExtension.class) +class GlueDatabaseOperationsTest { + + private GlueClient glueClient; + private GlueDatabaseOperator glueDatabaseOperations; + private String db1; + private String db2; + private String testDbUpper; + + @BeforeEach + void setUp() { + glueClient = GlueTestClientFactory.createClient(); + glueDatabaseOperations = new GlueDatabaseOperator(glueClient, "testCatalog"); + db1 = GlueTestClientFactory.uniqueName("db1"); + db2 = GlueTestClientFactory.uniqueName("db2"); + testDbUpper = GlueTestClientFactory.uniqueName("TestDB"); + } + + @Test + void testCreateDatabase() throws DatabaseAlreadyExistException, DatabaseNotExistException { + CatalogDatabase catalogDatabase = new CatalogDatabaseImpl(Collections.emptyMap(), "test"); + glueDatabaseOperations.createDatabase(db1, catalogDatabase); + Assertions.assertTrue(glueDatabaseOperations.glueDatabaseExists(db1)); + Assertions.assertEquals( + "test", glueDatabaseOperations.getDatabase(db1).getDescription().orElse(null)); + } + + @Test + void testCreateDatabaseWithUppercaseLetters() + throws DatabaseAlreadyExistException, DatabaseNotExistException { + CatalogDatabase catalogDatabase = new CatalogDatabaseImpl(Collections.emptyMap(), "test"); + // Uppercase letters should now be accepted with case preservation + Assertions.assertDoesNotThrow( + () -> glueDatabaseOperations.createDatabase(testDbUpper, catalogDatabase)); + + // Verify database was created and exists + Assertions.assertTrue(glueDatabaseOperations.glueDatabaseExists(testDbUpper)); + + // Verify the database can be retrieved + CatalogDatabase retrieved = glueDatabaseOperations.getDatabase(testDbUpper); + Assertions.assertNotNull(retrieved); + Assertions.assertEquals("test", retrieved.getDescription().orElse(null)); + } + + @Test + void testCreateDatabaseWithHyphens() { + CatalogDatabase catalogDatabase = new CatalogDatabaseImpl(Collections.emptyMap(), "test"); + CatalogException exception = + Assertions.assertThrows( + CatalogException.class, + () -> glueDatabaseOperations.createDatabase("db-1", catalogDatabase)); + Assertions.assertTrue( + exception.getMessage().contains("letters, numbers, and underscores"), + "Exception message should mention allowed characters"); + } + + @Test + void testCreateDatabaseWithSpecialCharacters() { + CatalogDatabase catalogDatabase = new CatalogDatabaseImpl(Collections.emptyMap(), "test"); + CatalogException exception = + Assertions.assertThrows( + CatalogException.class, + () -> glueDatabaseOperations.createDatabase("db.1", catalogDatabase)); + Assertions.assertTrue( + exception.getMessage().contains("letters, numbers, and underscores"), + "Exception message should mention allowed characters"); + } + + @Test + void testCreateDatabaseAlreadyExists() throws DatabaseAlreadyExistException { + CatalogDatabase catalogDatabase = + new CatalogDatabaseImpl(Collections.emptyMap(), "Description"); + glueDatabaseOperations.createDatabase(db1, catalogDatabase); + Assertions.assertThrows( + DatabaseAlreadyExistException.class, + () -> glueDatabaseOperations.createDatabase(db1, catalogDatabase)); + } + + @Test + void testCreateDatabaseInvalidInput() throws DatabaseAlreadyExistException { + CatalogDatabase catalogDatabase = + new CatalogDatabaseImpl(Collections.emptyMap(), "Description"); + fakeClient() + .setNextException( + InvalidInputException.builder().message("Invalid database name").build()); + Assertions.assertThrows( + CatalogException.class, + () -> glueDatabaseOperations.createDatabase(db1, catalogDatabase)); + } + + @Test + void testCreateDatabaseResourceLimitExceeded() throws DatabaseAlreadyExistException { + CatalogDatabase catalogDatabase = + new CatalogDatabaseImpl(Collections.emptyMap(), "Description"); + fakeClient() + .setNextException( + ResourceNumberLimitExceededException.builder() + .message("Resource limit exceeded") + .build()); + Assertions.assertThrows( + CatalogException.class, + () -> glueDatabaseOperations.createDatabase(db1, catalogDatabase)); + } + + @Test + void testCreateDatabaseTimeout() throws DatabaseAlreadyExistException { + CatalogDatabase catalogDatabase = + new CatalogDatabaseImpl(Collections.emptyMap(), "Description"); + fakeClient() + .setNextException( + OperationTimeoutException.builder().message("Operation timed out").build()); + Assertions.assertThrows( + CatalogException.class, + () -> glueDatabaseOperations.createDatabase(db1, catalogDatabase)); + } + + @Test + void testDropDatabase() throws DatabaseAlreadyExistException { + CatalogDatabase catalogDatabase = + new CatalogDatabaseImpl(Collections.emptyMap(), "Description"); + glueDatabaseOperations.createDatabase(db1, catalogDatabase); + Assertions.assertDoesNotThrow(() -> glueDatabaseOperations.dropGlueDatabase(db1)); + Assertions.assertFalse(glueDatabaseOperations.glueDatabaseExists(db1)); + } + + @Test + void testDropDatabaseNotFound() { + Assertions.assertThrows( + DatabaseNotExistException.class, + () -> glueDatabaseOperations.dropGlueDatabase(db1)); + } + + @Test + void testDropDatabaseInvalidInput() { + fakeClient() + .setNextException( + InvalidInputException.builder().message("Invalid database name").build()); + Assertions.assertThrows( + CatalogException.class, () -> glueDatabaseOperations.dropGlueDatabase(db1)); + } + + @Test + void testDropDatabaseTimeout() { + fakeClient() + .setNextException( + OperationTimeoutException.builder().message("Operation timed out").build()); + Assertions.assertThrows( + CatalogException.class, () -> glueDatabaseOperations.dropGlueDatabase(db1)); + } + + @Test + void testListDatabases() throws DatabaseAlreadyExistException { + CatalogDatabase catalogDatabase1 = new CatalogDatabaseImpl(Collections.emptyMap(), "test1"); + CatalogDatabase catalogDatabase2 = new CatalogDatabaseImpl(Collections.emptyMap(), "test2"); + glueDatabaseOperations.createDatabase(db1, catalogDatabase1); + glueDatabaseOperations.createDatabase(db2, catalogDatabase2); + + List<String> databaseNames = glueDatabaseOperations.listDatabases(); + Assertions.assertTrue(databaseNames.contains(db1)); + Assertions.assertTrue(databaseNames.contains(db2)); + } + + @Test + void testListDatabasesTimeout() { + fakeClient() + .setNextException( + OperationTimeoutException.builder().message("Operation timed out").build()); + Assertions.assertThrows( + CatalogException.class, () -> glueDatabaseOperations.listDatabases()); + } + + @Test + void testListDatabasesResourceLimitExceeded() { + fakeClient() + .setNextException( + ResourceNumberLimitExceededException.builder() + .message("Resource limit exceeded") + .build()); + Assertions.assertThrows( + CatalogException.class, () -> glueDatabaseOperations.listDatabases()); + } + + @Test + void testGetDatabase() throws DatabaseNotExistException, DatabaseAlreadyExistException { + CatalogDatabase catalogDatabase = + new CatalogDatabaseImpl(Collections.emptyMap(), "comment"); + glueDatabaseOperations.createDatabase(db1, catalogDatabase); + CatalogDatabase retrievedDatabase = glueDatabaseOperations.getDatabase(db1); + Assertions.assertNotNull(retrievedDatabase); + Assertions.assertEquals("comment", retrievedDatabase.getComment()); + } + + @Test + void testGetDatabaseNotFound() { + Assertions.assertThrows( + DatabaseNotExistException.class, () -> glueDatabaseOperations.getDatabase(db1)); + } + + @Test + void testGetDatabaseInvalidInput() { + fakeClient() + .setNextException( + InvalidInputException.builder().message("Invalid database name").build()); + Assertions.assertThrows( Review Comment: Done — both tests assert the exact message (`Invalid database name 'x': …` / `Timed out looking up database 'x'`) and the cause type. ########## flink-catalog-aws/flink-catalog-aws-glue/src/test/java/org/apache/flink/table/catalog/glue/util/GlueTypeConverterTest.java: ########## @@ -0,0 +1,226 @@ +/* + * 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.flink.table.catalog.glue.util; + +import org.apache.flink.table.api.DataTypes; +import org.apache.flink.table.catalog.glue.exception.UnsupportedDataTypeMappingException; +import org.apache.flink.table.types.DataType; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +class GlueTypeConverterTest { + + private final GlueTypeConverter converter = new GlueTypeConverter(); + + @Test Review Comment: Done — see the parameterised tables in `GlueTypeConverterTest` (thread on `GlueTypeConverter:61`). ########## flink-catalog-aws/flink-catalog-aws-glue/src/test/java/org/apache/flink/table/catalog/glue/GlueCatalogSchemaFidelityTest.java: ########## @@ -0,0 +1,235 @@ +/* + * 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.flink.table.catalog.glue; + +import org.apache.flink.table.api.DataTypes; +import org.apache.flink.table.api.Schema; +import org.apache.flink.table.catalog.CatalogBaseTable; +import org.apache.flink.table.catalog.CatalogDatabaseImpl; +import org.apache.flink.table.catalog.CatalogTable; +import org.apache.flink.table.catalog.CatalogView; +import org.apache.flink.table.catalog.Column; +import org.apache.flink.table.catalog.ObjectPath; +import org.apache.flink.table.catalog.ResolvedCatalogTable; +import org.apache.flink.table.catalog.ResolvedCatalogView; +import org.apache.flink.table.catalog.ResolvedSchema; +import org.apache.flink.table.catalog.UniqueConstraint; +import org.apache.flink.table.catalog.WatermarkSpec; +import org.apache.flink.table.catalog.glue.util.GlueTestClientFactory; +import org.apache.flink.table.catalog.glue.util.RealGlueCleanupExtension; +import org.apache.flink.table.expressions.ExpressionVisitor; +import org.apache.flink.table.expressions.ResolvedExpression; +import org.apache.flink.table.types.DataType; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.GetTableRequest; +import software.amazon.awssdk.services.glue.model.Table; + +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Tests that schema features which AWS Glue columns cannot represent - computed columns, metadata + * columns, watermarks, and primary keys - survive a full create/read round-trip, and that views do + * not leak their columns into Glue partition keys. + */ +@ExtendWith(RealGlueCleanupExtension.class) +class GlueCatalogSchemaFidelityTest { + + private GlueClient glueClient; + private GlueCatalog glueCatalog; + private String databaseName; + private String glueDatabaseName; + + @BeforeEach + void setUp() throws Exception { + glueClient = GlueTestClientFactory.createClient(); + glueCatalog = new GlueCatalog("test_catalog", "default", "us-east-1", glueClient); + databaseName = GlueTestClientFactory.uniqueName("fidelitydb"); + glueDatabaseName = databaseName.toLowerCase(); + glueCatalog.createDatabase( + databaseName, new CatalogDatabaseImpl(new HashMap<>(), "fidelity tests"), false); + } + + @AfterEach + void tearDown() { + if (glueCatalog != null) { + glueCatalog.close(); + } + } + + @Test + void testWatermarkPrimaryKeyAndNonPhysicalColumnsRoundTrip() throws Exception { Review Comment: Added: partitioned round-trip with every column type asserted (types resolved through a real planner `DataTypeFactory`, since restored types are `UnresolvedDataType`s), column comments, view types, a foreign `EXTERNAL_TABLE` with Hive `varchar(255)`/`char(10)`/`decimal(10,2)`/`array`/`map` columns, a foreign `VIRTUAL_VIEW`, `alterTable` refusing to overwrite a view or change partition keys, reserved options, and an unresolved table getting a clear message. ########## flink-catalog-aws/flink-catalog-aws-glue/src/test/java/org/apache/flink/table/catalog/glue/GlueCatalogSqlMotoITCase.java: ########## @@ -0,0 +1,188 @@ +/* + * 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.flink.table.catalog.glue; + +import org.apache.flink.table.api.EnvironmentSettings; +import org.apache.flink.table.api.TableEnvironment; +import org.apache.flink.types.Row; +import org.apache.flink.util.CollectionUtil; + +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.testcontainers.containers.GenericContainer; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import software.amazon.awssdk.auth.credentials.AwsBasicCredentials; +import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider; +import software.amazon.awssdk.regions.Region; +import software.amazon.awssdk.services.glue.GlueClient; + +import java.net.URI; +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * SQL-path integration test for the Glue catalog against a moto Glue emulator. + * + * <p>Unlike {@link GlueCatalogMotoITCase}, which drives {@link GlueCatalog} methods directly, this + * test exercises the full user path: {@code CREATE CATALOG ... WITH ('type'='glue')} discovers + * {@code GlueCatalogFactory} via SPI, the factory builds its own {@code GlueClient}, and all + * catalog operations flow through Flink SQL DDL and the planner down to the Glue wire protocol. + * + * <p>The factory-built client is pointed at moto through the catalog's {@code aws.endpoint} option + * (handled by the shared AWS client-creation path) and system-property credentials, so this test + * also covers the factory's AWS option pass-through. + */ +@Testcontainers +class GlueCatalogSqlMotoITCase { + + private static final int MOTO_PORT = 5000; + + @Container + private static final GenericContainer<?> MOTO = + new GenericContainer<>("motoserver/moto:5.0.28").withExposedPorts(MOTO_PORT); + + private static GlueClient seedClient; + private static TableEnvironment tEnv; + + @BeforeAll + static void setUp() { + String endpoint = + String.format("http://%s:%d", MOTO.getHost(), MOTO.getMappedPort(MOTO_PORT)); + + // Credentials come from system properties (first in the default credentials chain); + // the endpoint is routed to moto via the catalog's own 'aws.endpoint' option below. + System.setProperty("aws.accessKeyId", "testing"); + System.setProperty("aws.secretAccessKey", "testing"); + + // Seed the default database (USE CATALOG validates it exists). + seedClient = + GlueClient.builder() + .endpointOverride(URI.create(endpoint)) + .region(Region.US_EAST_1) + .credentialsProvider( + StaticCredentialsProvider.create( + AwsBasicCredentials.create("testing", "testing"))) + .build(); + seedClient.createDatabase(builder -> builder.databaseInput(db -> db.name("default"))); + + tEnv = TableEnvironment.create(EnvironmentSettings.inStreamingMode()); + tEnv.executeSql( + "CREATE CATALOG glue_moto WITH (" + + "'type' = 'glue', " + + "'region' = 'us-east-1', " + + "'aws.endpoint' = '" + + endpoint + + "', " + + "'default-database' = 'default')"); + tEnv.executeSql("USE CATALOG glue_moto"); + } + + @AfterAll + static void tearDown() { + System.clearProperty("aws.accessKeyId"); + System.clearProperty("aws.secretAccessKey"); + if (seedClient != null) { + seedClient.close(); + } + } + + private static List<Row> sql(String statement) { + return CollectionUtil.iteratorToList(tEnv.executeSql(statement).collect()); + } + + @Test + void testShowDatabasesThroughSql() { + List<Row> databases = sql("SHOW DATABASES"); + + assertThat(databases).extracting(row -> row.getField(0)).contains("default"); + } + + @Test + void testDatabaseDdlThroughSql() { + tEnv.executeSql("CREATE DATABASE sql_ddl_db COMMENT 'created via SQL DDL'"); + + assertThat(sql("SHOW DATABASES")).extracting(row -> row.getField(0)).contains("sql_ddl_db"); + } + + @Test + void testTableDdlRoundTripThroughSql() { + tEnv.executeSql("CREATE DATABASE sql_table_db"); + tEnv.executeSql( + "CREATE TABLE sql_table_db.orders (" + + " user_id STRING," + + " order_total DOUBLE" + + ") WITH (" + + " 'connector' = 'kinesis'," + + " 'stream.arn' = 'arn:aws:kinesis:us-east-1:000000000000:stream/orders'" + + ")"); + + assertThat(sql("SHOW TABLES IN sql_table_db")) + .extracting(row -> row.getField(0)) + .contains("orders"); + + // Schema must survive the wire round-trip through moto and come back through DESCRIBE. + List<Row> columns = sql("DESCRIBE sql_table_db.orders"); + assertThat(columns).hasSize(2); + assertThat(columns.get(0).getField(0)).isEqualTo("user_id"); + assertThat(columns.get(0).getField(1)).isEqualTo("STRING"); + assertThat(columns.get(1).getField(0)).isEqualTo("order_total"); + assertThat(columns.get(1).getField(1)).isEqualTo("DOUBLE"); + + tEnv.executeSql("DROP TABLE sql_table_db.orders"); + assertThat(sql("SHOW TABLES IN sql_table_db")) + .extracting(row -> row.getField(0)) + .doesNotContain("orders"); + } + + @Test + void testSchemaFidelityRoundTripThroughSql() { + // Note: watermark and computed-column DDL is exercised in the fake-backed and + // real-AWS tiers instead. Validating any SQL expression makes the planner probe the + // catalog's function APIs, and moto does not implement the Glue UDF API (HTTP 500). + tEnv.executeSql("CREATE DATABASE sql_fidelity_db"); + tEnv.executeSql( + "CREATE TABLE sql_fidelity_db.events (" + + " userId STRING," + + " eventTime TIMESTAMP(3)," + + " price DOUBLE," + + " kafkaOffset BIGINT METADATA FROM 'offset' VIRTUAL," + + " PRIMARY KEY (userId) NOT ENFORCED" + + ") WITH (" + + " 'connector' = 'kinesis'," + + " 'stream.arn' = 'arn:aws:kinesis:us-east-1:000000000000:stream/events'" + + ")"); + + // Read the table back through the catalog: primary key and metadata columns must + // survive the round-trip through Glue table parameters. + String createTable = + sql("SHOW CREATE TABLE sql_fidelity_db.events").get(0).getField(0).toString(); + + assertThat(createTable) + .contains("`kafkaOffset` BIGINT METADATA FROM 'offset' VIRTUAL") + .contains("PRIMARY KEY (`userId`) NOT ENFORCED"); + + // Column order must be preserved, including the interleaved metadata column. + assertThat(sql("DESCRIBE sql_fidelity_db.events")) + .extracting(row -> row.getField(0)) + .containsExactly("userId", "eventTime", "price", "kafkaOffset"); Review Comment: Done — `DESCRIBE` is asserted per column for name, type and nullability, and a partitioned wire round-trip was added (`testPartitionedTableRoundTripPreservesDeclaredOrderAndTypes`). -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
