fmorillo7694 commented on code in PR #206: URL: https://github.com/apache/flink-connector-aws/pull/206#discussion_r4125202099
########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/util/GlueTableUtils.java: ########## @@ -0,0 +1,236 @@ +/* + * 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.Schema; +import org.apache.flink.table.catalog.ObjectPath; +import org.apache.flink.table.types.DataType; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +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.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Map; + +/** + * Utility class for working with Glue tables, including transforming Glue-specific metadata into + * Flink-compatible objects. + */ +public class GlueTableUtils { + + /** Logger for logging Glue table operations. */ + private static final Logger LOG = LoggerFactory.getLogger(GlueTableUtils.class); + + /** Glue type converter for type conversions between Flink and Glue types. */ + private final GlueTypeConverter glueTypeConverter; + + /** + * Constructor to initialize GlueTableUtils with a GlueTypeConverter. + * + * @param glueTypeConverter The GlueTypeConverter instance for type mapping. + */ + public GlueTableUtils(GlueTypeConverter glueTypeConverter) { + this.glueTypeConverter = glueTypeConverter; + } + + /** + * Builds a Glue StorageDescriptor from the given table properties, columns, and location. + * + * @param tableProperties Table properties for the Glue table. + * @param glueColumns Columns to be included in the StorageDescriptor. + * @param tableLocation Location of the Glue table. + * @return A newly built StorageDescriptor object. + */ + public StorageDescriptor buildStorageDescriptor( + Map<String, String> tableProperties, List<Column> glueColumns, String tableLocation) { + + return StorageDescriptor.builder().columns(glueColumns).location(tableLocation).build(); + } + + /** + * Extracts the table location based on the table properties and the table path. First, it + * checks for a location key from the connector registry. If no such key is found, it uses a + * default path based on the table path. + * + * @param tableProperties Table properties containing the connector and location. + * @param tablePath The Flink ObjectPath representing the table. + * @return The location of the Glue table. + */ + public String extractTableLocation(Map<String, String> tableProperties, ObjectPath tablePath) { + String connectorType = tableProperties.get("connector"); + if (connectorType != null) { + String locationKey = ConnectorRegistry.getLocationKey(connectorType); + if (locationKey != null && tableProperties.containsKey(locationKey)) { + String location = tableProperties.get(locationKey); + return location; + } + } + + String defaultLocation = + tablePath.getDatabaseName() + "/tables/" + tablePath.getObjectName(); + return defaultLocation; + } + + /** + * Converts a Flink column to a Glue column. The column's data type is converted using the + * GlueTypeConverter. + * + * @param flinkColumn The Flink column to be converted. + * @return The corresponding Glue column. + */ + public Column mapFlinkColumnToGlueColumn(org.apache.flink.table.catalog.Column flinkColumn) { + String glueType = glueTypeConverter.toGlueDataType(flinkColumn.getDataType()); + + // AWS Glue lowercases column names on CreateTable/UpdateTable (verified empirically: + // a column created as "userId" is stored and returned as "userid"). To preserve the + // declared case, store the lowercased name explicitly and stash the original name in + // the "originalName" column parameter, which the read path restores. + String originalName = flinkColumn.getName(); + String glueName = originalName.toLowerCase(); + + Column.Builder builder = Column.builder().name(glueName).type(glueType); + if (!glueName.equals(originalName)) { + builder.parameters( + Collections.singletonMap( + GlueCatalogConstants.ORIGINAL_COLUMN_NAME, originalName)); + } + return builder.build(); + } + + /** + * Converts a Glue table into a Flink schema. Each Glue column is mapped to a Flink column using + * the GlueTypeConverter. Partition columns (stored at the Glue table level, not in the storage + * descriptor) are appended after the data columns so that declared partition keys are part of + * the Flink schema, as required by {@code CatalogTable}. Computed and metadata columns, + * watermarks, and the primary key - which Glue columns cannot represent - are restored from the + * {@code flink.schema.*} table parameters written by {@link GlueFlinkSchemaProperties}. + * + * @param glueTable The Glue table from which the schema will be derived. + * @return A Flink schema constructed from the Glue table's columns. + */ + public Schema getSchemaFromGlueTable(Table glueTable) { + Schema.Builder schemaBuilder = Schema.newBuilder(); + + List<GlueFlinkSchemaProperties.PhysicalColumnSpec> physicalColumns = new ArrayList<>(); + java.util.Set<String> notNullColumns = + GlueFlinkSchemaProperties.getNotNullColumns(glueTable.parameters()); + List<Column> columns = + glueTable.storageDescriptor() != null + ? glueTable.storageDescriptor().columns() + : Collections.emptyList(); + for (Column column : columns) { + String columnName = getColumnName(column); + DataType convertedType = glueTypeConverter.toFlinkDataType(column.type()); + // Glue type strings carry no nullability; restore the declared NOT NULL + // constraints (required for primary key columns to resolve). + DataType flinkDataType = + notNullColumns.contains(columnName) ? convertedType.notNull() : convertedType; + physicalColumns.add( + new GlueFlinkSchemaProperties.PhysicalColumnSpec(columnName, flinkDataType)); + } + + // Partition columns live in Table.partitionKeys(), not in the storage descriptor. + // Their declared case is restored from the table-level parameter (Glue rejects + // column-level parameters on partition columns). + if (glueTable.partitionKeys() != null && !glueTable.partitionKeys().isEmpty()) { + List<String> partitionKeyNames = getPartitionKeyNames(glueTable); + List<Column> partitionColumns = glueTable.partitionKeys(); + for (int i = 0; i < partitionColumns.size(); i++) { Review Comment: You were right, and it was worse than a wrong order: it corrupted types (see your follow-up on `GlueFlinkSchemaProperties:112`). Fixed at the root by keying every per-column parameter by **name** and persisting `flink.schema.column-order`, so read-back is independent of how Glue orders storage-descriptor columns vs. partition keys. Regression tests: `testPartitionedTableDeclaredOrderAndTypesSurviveRoundTrip` (partition key declared first, every type asserted), the same over the wire in `GlueCatalogSqlMotoITCase`, and `testPartitionedTableEndToEnd` against real Glue. ########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/util/GlueTableUtils.java: ########## @@ -0,0 +1,236 @@ +/* + * 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.Schema; +import org.apache.flink.table.catalog.ObjectPath; +import org.apache.flink.table.types.DataType; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +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.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Map; + +/** + * Utility class for working with Glue tables, including transforming Glue-specific metadata into + * Flink-compatible objects. + */ +public class GlueTableUtils { + + /** Logger for logging Glue table operations. */ + private static final Logger LOG = LoggerFactory.getLogger(GlueTableUtils.class); + + /** Glue type converter for type conversions between Flink and Glue types. */ + private final GlueTypeConverter glueTypeConverter; + + /** + * Constructor to initialize GlueTableUtils with a GlueTypeConverter. + * + * @param glueTypeConverter The GlueTypeConverter instance for type mapping. + */ + public GlueTableUtils(GlueTypeConverter glueTypeConverter) { + this.glueTypeConverter = glueTypeConverter; + } + + /** + * Builds a Glue StorageDescriptor from the given table properties, columns, and location. + * + * @param tableProperties Table properties for the Glue table. + * @param glueColumns Columns to be included in the StorageDescriptor. + * @param tableLocation Location of the Glue table. + * @return A newly built StorageDescriptor object. + */ + public StorageDescriptor buildStorageDescriptor( + Map<String, String> tableProperties, List<Column> glueColumns, String tableLocation) { + + return StorageDescriptor.builder().columns(glueColumns).location(tableLocation).build(); + } + + /** + * Extracts the table location based on the table properties and the table path. First, it + * checks for a location key from the connector registry. If no such key is found, it uses a + * default path based on the table path. + * + * @param tableProperties Table properties containing the connector and location. + * @param tablePath The Flink ObjectPath representing the table. + * @return The location of the Glue table. + */ + public String extractTableLocation(Map<String, String> tableProperties, ObjectPath tablePath) { + String connectorType = tableProperties.get("connector"); + if (connectorType != null) { + String locationKey = ConnectorRegistry.getLocationKey(connectorType); + if (locationKey != null && tableProperties.containsKey(locationKey)) { + String location = tableProperties.get(locationKey); + return location; + } + } + + String defaultLocation = + tablePath.getDatabaseName() + "/tables/" + tablePath.getObjectName(); + return defaultLocation; + } + + /** + * Converts a Flink column to a Glue column. The column's data type is converted using the + * GlueTypeConverter. + * + * @param flinkColumn The Flink column to be converted. + * @return The corresponding Glue column. + */ + public Column mapFlinkColumnToGlueColumn(org.apache.flink.table.catalog.Column flinkColumn) { + String glueType = glueTypeConverter.toGlueDataType(flinkColumn.getDataType()); + + // AWS Glue lowercases column names on CreateTable/UpdateTable (verified empirically: + // a column created as "userId" is stored and returned as "userid"). To preserve the + // declared case, store the lowercased name explicitly and stash the original name in + // the "originalName" column parameter, which the read path restores. + String originalName = flinkColumn.getName(); + String glueName = originalName.toLowerCase(); + + Column.Builder builder = Column.builder().name(glueName).type(glueType); + if (!glueName.equals(originalName)) { + builder.parameters( + Collections.singletonMap( + GlueCatalogConstants.ORIGINAL_COLUMN_NAME, originalName)); + } + return builder.build(); + } + + /** + * Converts a Glue table into a Flink schema. Each Glue column is mapped to a Flink column using + * the GlueTypeConverter. Partition columns (stored at the Glue table level, not in the storage + * descriptor) are appended after the data columns so that declared partition keys are part of + * the Flink schema, as required by {@code CatalogTable}. Computed and metadata columns, + * watermarks, and the primary key - which Glue columns cannot represent - are restored from the + * {@code flink.schema.*} table parameters written by {@link GlueFlinkSchemaProperties}. + * + * @param glueTable The Glue table from which the schema will be derived. + * @return A Flink schema constructed from the Glue table's columns. + */ + public Schema getSchemaFromGlueTable(Table glueTable) { + Schema.Builder schemaBuilder = Schema.newBuilder(); + + List<GlueFlinkSchemaProperties.PhysicalColumnSpec> physicalColumns = new ArrayList<>(); + java.util.Set<String> notNullColumns = + GlueFlinkSchemaProperties.getNotNullColumns(glueTable.parameters()); + List<Column> columns = + glueTable.storageDescriptor() != null + ? glueTable.storageDescriptor().columns() + : Collections.emptyList(); + for (Column column : columns) { + String columnName = getColumnName(column); + DataType convertedType = glueTypeConverter.toFlinkDataType(column.type()); + // Glue type strings carry no nullability; restore the declared NOT NULL + // constraints (required for primary key columns to resolve). + DataType flinkDataType = + notNullColumns.contains(columnName) ? convertedType.notNull() : convertedType; + physicalColumns.add( + new GlueFlinkSchemaProperties.PhysicalColumnSpec(columnName, flinkDataType)); + } + + // Partition columns live in Table.partitionKeys(), not in the storage descriptor. + // Their declared case is restored from the table-level parameter (Glue rejects + // column-level parameters on partition columns). + if (glueTable.partitionKeys() != null && !glueTable.partitionKeys().isEmpty()) { + List<String> partitionKeyNames = getPartitionKeyNames(glueTable); + List<Column> partitionColumns = glueTable.partitionKeys(); + for (int i = 0; i < partitionColumns.size(); i++) { + Column partitionColumn = partitionColumns.get(i); + DataType convertedType = glueTypeConverter.toFlinkDataType(partitionColumn.type()); + String partitionKeyName = partitionKeyNames.get(i); + DataType flinkDataType = + notNullColumns.contains(partitionKeyName) + ? convertedType.notNull() + : convertedType; + physicalColumns.add( + new GlueFlinkSchemaProperties.PhysicalColumnSpec( + partitionKeyName, flinkDataType)); + } + } + + // Merge in computed/metadata columns and re-apply watermarks and the primary key. + GlueFlinkSchemaProperties.applySchemaWithNonPhysicalColumns( + glueTable.parameters(), physicalColumns, schemaBuilder); + + return schemaBuilder.build(); + } + + /** + * Returns the Flink-facing partition key names of a Glue table, in declared order. The original + * (case-preserved) names come from the table-level {@link + * GlueCatalogConstants#ORIGINAL_PARTITION_KEYS} parameter when present (written by this catalog + * because Glue rejects column-level parameters on partition columns), falling back to the + * per-column resolution for tables written by other writers or older versions. + * + * @param glueTable The Glue table. + * @return Ordered partition key names to expose to Flink; empty when not partitioned. + */ + public static List<String> getPartitionKeyNames(Table glueTable) { + if (glueTable.partitionKeys() == null || glueTable.partitionKeys().isEmpty()) { + return Collections.emptyList(); + } + List<Column> partitionColumns = glueTable.partitionKeys(); + if (glueTable.parameters() != null + && glueTable + .parameters() + .containsKey(GlueCatalogConstants.ORIGINAL_PARTITION_KEYS)) { + String[] originalNames = + glueTable + .parameters() + .get(GlueCatalogConstants.ORIGINAL_PARTITION_KEYS) + .split(",", -1); + if (originalNames.length == partitionColumns.size()) { + return java.util.Arrays.asList(originalNames); + } + LOG.warn( + "Ignoring malformed {} parameter on table {}: {} entries for {} partition keys", + GlueCatalogConstants.ORIGINAL_PARTITION_KEYS, + glueTable.name(), + originalNames.length, + partitionColumns.size()); + } + List<String> names = new java.util.ArrayList<>(partitionColumns.size()); + for (Column partitionColumn : partitionColumns) { + names.add(getColumnName(partitionColumn)); + } + return names; + } + + /** + * Returns the Flink-facing name of a Glue column: the original (case-preserved) name from the + * "originalName" column parameter when present, otherwise the Glue-stored name. Glue lowercases + * column names on write, so this parameter is how the declared case survives the round-trip. + * + * @param column The Glue column. + * @return The column name to expose to Flink. + */ + public static String getColumnName(Column column) { + if (column.parameters() != null + && column.parameters().containsKey(GlueCatalogConstants.ORIGINAL_COLUMN_NAME)) { Review Comment: Consistency fixed: the column parameter is now `flink.original-column-name`. The old `originalName` key is still read so tables written by the earlier revision of this branch keep working. ########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/util/GlueTypeConverter.java: ########## @@ -0,0 +1,315 @@ +/* + * 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.apache.flink.table.types.logical.ArrayType; +import org.apache.flink.table.types.logical.DecimalType; +import org.apache.flink.table.types.logical.LogicalType; +import org.apache.flink.table.types.logical.LogicalTypeRoot; +import org.apache.flink.table.types.logical.MapType; +import org.apache.flink.table.types.logical.RowType; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.ArrayList; +import java.util.List; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +/** + * Utility class for converting Flink types to Glue types and vice versa. Supports the conversion of + * common primitive, array, map, and struct types. + */ +public class GlueTypeConverter { + + /** Logger for tracking Glue type conversions. */ + private static final Logger LOG = LoggerFactory.getLogger(GlueTypeConverter.class); + + /** Regular expressions for handling specific Glue types. */ + private static final Pattern DECIMAL_PATTERN = Pattern.compile("decimal\\((\\d+),(\\d+)\\)"); + + private static final Pattern ARRAY_PATTERN = Pattern.compile("array<(.+)>"); + private static final Pattern MAP_PATTERN = Pattern.compile("map<(.+),(.+)>"); + private static final Pattern STRUCT_PATTERN = Pattern.compile("struct<(.+)>"); + + /** + * Converts a Flink DataType to its corresponding Glue type as a string. + * + * @param flinkType The Flink DataType to be converted. + * @return The Glue type as a string. + */ + public String toGlueDataType(DataType flinkType) { Review Comment: Done — `GlueTypeConverterTest` is now parameterised over the full tables: 26 Glue→Flink cases (every Glue primitive in the spellings other engines write, incl. `integer`, bare `decimal`, `char(n)`/`varchar(n)`, whitespace/case variants, nested complex types) and 22 Flink→Glue cases (every supported type root incl. `TIMESTAMP_LTZ`, `TIME`, `BINARY`, NOT NULL), plus the lossy read-backs and the rejected Flink types. ########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/GlueCatalog.java: ########## @@ -0,0 +1,1850 @@ +/* + * 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.annotation.PublicEvolving; +import org.apache.flink.annotation.VisibleForTesting; +import org.apache.flink.connector.aws.config.AWSConfigConstants; +import org.apache.flink.connector.aws.util.AWSClientUtil; +import org.apache.flink.connector.aws.util.AWSGeneralUtil; +import org.apache.flink.table.api.Schema; +import org.apache.flink.table.catalog.AbstractCatalog; +import org.apache.flink.table.catalog.CatalogBaseTable; +import org.apache.flink.table.catalog.CatalogDatabase; +import org.apache.flink.table.catalog.CatalogFunction; +import org.apache.flink.table.catalog.CatalogPartition; +import org.apache.flink.table.catalog.CatalogPartitionImpl; +import org.apache.flink.table.catalog.CatalogPartitionSpec; +import org.apache.flink.table.catalog.CatalogTable; +import org.apache.flink.table.catalog.CatalogView; +import org.apache.flink.table.catalog.ObjectPath; +import org.apache.flink.table.catalog.ResolvedCatalogBaseTable; +import org.apache.flink.table.catalog.exceptions.CatalogException; +import org.apache.flink.table.catalog.exceptions.DatabaseAlreadyExistException; +import org.apache.flink.table.catalog.exceptions.DatabaseNotEmptyException; +import org.apache.flink.table.catalog.exceptions.DatabaseNotExistException; +import org.apache.flink.table.catalog.exceptions.FunctionAlreadyExistException; +import org.apache.flink.table.catalog.exceptions.FunctionNotExistException; +import org.apache.flink.table.catalog.exceptions.PartitionAlreadyExistsException; +import org.apache.flink.table.catalog.exceptions.PartitionNotExistException; +import org.apache.flink.table.catalog.exceptions.PartitionSpecInvalidException; +import org.apache.flink.table.catalog.exceptions.TableAlreadyExistException; +import org.apache.flink.table.catalog.exceptions.TableNotExistException; +import org.apache.flink.table.catalog.exceptions.TableNotPartitionedException; +import org.apache.flink.table.catalog.exceptions.TablePartitionedException; +import org.apache.flink.table.catalog.glue.operator.GlueDatabaseOperator; +import org.apache.flink.table.catalog.glue.operator.GlueFunctionOperator; +import org.apache.flink.table.catalog.glue.operator.GluePartitionOperator; +import org.apache.flink.table.catalog.glue.operator.GlueTableOperator; +import org.apache.flink.table.catalog.glue.util.GlueCatalogConstants; +import org.apache.flink.table.catalog.glue.util.GlueFlinkSchemaProperties; +import org.apache.flink.table.catalog.glue.util.GlueTableUtils; +import org.apache.flink.table.catalog.glue.util.GlueTypeConverter; +import org.apache.flink.table.catalog.stats.CatalogColumnStatistics; +import org.apache.flink.table.catalog.stats.CatalogTableStatistics; +import org.apache.flink.table.expressions.Expression; +import org.apache.flink.table.functions.FunctionIdentifier; +import org.apache.flink.util.Preconditions; +import org.apache.flink.util.StringUtils; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import software.amazon.awssdk.http.apache.ApacheHttpClient; +import software.amazon.awssdk.regions.Region; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.Partition; +import software.amazon.awssdk.services.glue.model.PartitionInput; +import software.amazon.awssdk.services.glue.model.StorageDescriptor; +import software.amazon.awssdk.services.glue.model.Table; +import software.amazon.awssdk.services.glue.model.TableInput; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Properties; +import java.util.stream.Collectors; + +/** + * GlueCatalog is an implementation of the Flink AbstractCatalog that interacts with AWS Glue. This + * class allows Flink to perform various catalog operations such as creating, deleting, and + * retrieving databases and tables from Glue. It encapsulates AWS Glue's API and provides a + * Flink-compatible interface. + * + * <p>This catalog uses GlueClient to interact with AWS Glue services, and operations related to + * databases and tables are delegated to respective helper classes like GlueDatabaseOperations and + * GlueTableOperations. + */ +@PublicEvolving +public class GlueCatalog extends AbstractCatalog { + + private static final Logger LOG = LoggerFactory.getLogger(GlueCatalog.class); + + private GlueClient glueClient; + private GlueTypeConverter glueTypeConverter; + private GlueDatabaseOperator glueDatabaseOperations; + private GlueTableOperator glueTableOperations; + private GlueFunctionOperator glueFunctionsOperations; + private GluePartitionOperator gluePartitionOperations; + private GlueTableUtils glueTableUtils; + + /** + * Constructs a GlueCatalog with a provided Glue client. + * + * @param name the name of the catalog + * @param defaultDatabase the default database for the catalog + * @param region the AWS region to be used for Glue operations + * @param glueClient Glue Client so we can decide which one to use for testing + */ + @VisibleForTesting + GlueCatalog(String name, String defaultDatabase, String region, GlueClient glueClient) { + super(name, defaultDatabase); + + // Validate region parameter + Preconditions.checkNotNull(region, "region cannot be null"); + Preconditions.checkArgument(!region.trim().isEmpty(), "region cannot be empty"); + + // Initialize GlueClient in the constructor + if (glueClient != null) { + setup(glueClient); + } else { + // If no GlueClient is provided, initialize it using the default region + GlueClient client = GlueClient.builder().region(Region.of(region)).build(); + setup(client); + } + } + + /** + * Constructs a GlueCatalog with default client configuration. + * + * @param name the name of the catalog + * @param defaultDatabase the default database for the catalog + * @param region the AWS region to be used for Glue operations + */ + public GlueCatalog(String name, String defaultDatabase, String region) { + this(name, defaultDatabase, region, new Properties()); + } + + /** + * Constructs a GlueCatalog whose Glue client is built from the given AWS client properties, + * using the same client-creation path as the other AWS connectors ({@link AWSClientUtil}). This + * makes the standard {@code aws.*} settings available to the catalog - for example {@code + * aws.credentials.provider} to select a credential mode, {@code aws.endpoint} to point at a + * Glue-compatible endpoint, and the {@code aws.http-client.*} options. + * + * @param name the name of the catalog + * @param defaultDatabase the default database for the catalog + * @param region the AWS region to be used for Glue operations + * @param glueClientProperties AWS client properties, keyed by {@link AWSConfigConstants} + */ + public GlueCatalog( + String name, String defaultDatabase, String region, Properties glueClientProperties) { + super(name, defaultDatabase); + + // Validate region parameter + Preconditions.checkNotNull(region, "region cannot be null"); + Preconditions.checkArgument(!region.trim().isEmpty(), "region cannot be empty"); + Preconditions.checkNotNull(glueClientProperties, "glueClientProperties cannot be null"); + + Properties clientProperties = new Properties(); + clientProperties.putAll(glueClientProperties); + // The explicit region argument wins over any aws.region property. + clientProperties.setProperty(AWSConfigConstants.AWS_REGION, region); + AWSGeneralUtil.validateAwsConfiguration(clientProperties); + + GlueClient client = + AWSClientUtil.createAwsSyncClient( + clientProperties, + AWSGeneralUtil.createSyncHttpClient( + clientProperties, ApacheHttpClient.builder()), + GlueClient.builder(), + GlueCatalogConstants.BASE_GLUE_USER_AGENT_PREFIX_FORMAT, + GlueCatalogConstants.GLUE_CLIENT_USER_AGENT_PREFIX); + setup(client); + } + + /** + * Private helper method to set up the GlueCatalog with a GlueClient instance. This method + * initializes all the necessary components and operators. + * + * @param glueClient the GlueClient to use for AWS Glue operations + */ + private void setup(GlueClient glueClient) { + this.glueClient = glueClient; + this.glueTypeConverter = new GlueTypeConverter(); + this.glueTableUtils = new GlueTableUtils(glueTypeConverter); + this.glueDatabaseOperations = new GlueDatabaseOperator(glueClient, getName()); + this.glueTableOperations = new GlueTableOperator(glueClient, getName()); + this.glueFunctionsOperations = new GlueFunctionOperator(glueClient, getName()); + this.gluePartitionOperations = new GluePartitionOperator(glueClient, getName()); + } + + /** + * Validates that a database exists, throwing DatabaseNotExistException if it doesn't. + * + * @param databaseName the name of the database to validate + * @throws DatabaseNotExistException if the database does not exist + * @throws CatalogException if an error occurs while checking database existence + */ + private void validateDatabaseExists(String databaseName) + throws DatabaseNotExistException, CatalogException { + Preconditions.checkArgument( + !StringUtils.isNullOrWhitespaceOnly(databaseName), + "databaseName cannot be null or empty"); + + if (!databaseExists(databaseName)) { + throw new DatabaseNotExistException(getName(), databaseName); + } + } + + /** + * Opens the GlueCatalog and initializes necessary resources. + * + * @throws CatalogException if an error occurs during the opening process + */ + @Override + public void open() throws CatalogException { + LOG.info("Opening GlueCatalog with client: {}", glueClient); + } + + /** + * Closes the GlueCatalog and releases resources. + * + * @throws CatalogException if an error occurs during the closing process + */ + @Override + public void close() throws CatalogException { + if (glueClient != null) { + LOG.info("Closing GlueCatalog client"); + // The AWS SDK close() is best-effort and does not surface exceptions, + // so no retry logic is required here. + glueClient.close(); + } + } + + /** + * Lists all the databases available in the Glue catalog. + * + * @return a list of database names + * @throws CatalogException if an error occurs while listing the databases + */ + @Override + public List<String> listDatabases() throws CatalogException { + return glueDatabaseOperations.listDatabases(); + } + + /** + * Retrieves a specific database by its name. + * + * @param databaseName the name of the database to retrieve + * @return the database if found + * @throws DatabaseNotExistException if the database does not exist + * @throws CatalogException if an error occurs while retrieving the database + */ + @Override + public CatalogDatabase getDatabase(String databaseName) + throws DatabaseNotExistException, CatalogException { + Preconditions.checkArgument( + !StringUtils.isNullOrWhitespaceOnly(databaseName), + "databaseName cannot be null or empty"); + + // Use case-insensitive database name resolution + String glueDatabaseName = findGlueDatabaseName(databaseName); + if (glueDatabaseName == null) { + throw new DatabaseNotExistException(getName(), databaseName); + } + + return glueDatabaseOperations.getDatabase(glueDatabaseName); + } + + /** + * Checks if a database exists in Glue. + * + * @param databaseName the name of the database + * @return true if the database exists, false otherwise + * @throws CatalogException if an error occurs while checking the database + */ + @Override + public boolean databaseExists(String databaseName) throws CatalogException { + Preconditions.checkArgument( + !StringUtils.isNullOrWhitespaceOnly(databaseName), + "databaseName cannot be null or empty"); + + // Use case-insensitive database name resolution + return findGlueDatabaseName(databaseName) != null; + } + + /** + * Creates a new database in Glue. + * + * @param databaseName the name of the database to create + * @param catalogDatabase the catalog database object containing database metadata + * @param ifNotExists flag indicating whether to ignore the error if the database already exists + * @throws DatabaseAlreadyExistException if the database already exists and ifNotExists is false + * @throws CatalogException if an error occurs while creating the database + */ + @Override + public void createDatabase( + String databaseName, CatalogDatabase catalogDatabase, boolean ifNotExists) + throws DatabaseAlreadyExistException, CatalogException { + Preconditions.checkArgument( + !StringUtils.isNullOrWhitespaceOnly(databaseName), + "databaseName cannot be null or empty"); + Preconditions.checkNotNull(catalogDatabase, "CatalogDatabase cannot be null"); + + // Check for exact case match first + boolean exactExists = databaseExists(databaseName); + if (exactExists && !ifNotExists) { + throw new DatabaseAlreadyExistException(getName(), databaseName); + } + if (exactExists) { + return; // Database exists with exact case, and IF NOT EXISTS is true + } + + // Check for case-insensitive collision (Glue limitation) + String conflictingDatabase = findCaseInsensitiveConflict(databaseName); + if (conflictingDatabase != null) { + String message = + String.format( + "Cannot create database '%s' because it conflicts with existing database '%s'. " + + "AWS Glue stores database names in lowercase, so '%s' and '%s' would both be stored as '%s'.", + databaseName, + conflictingDatabase, + databaseName, + conflictingDatabase, + databaseName.toLowerCase()); + throw new DatabaseAlreadyExistException( + getName(), databaseName, new CatalogException(message)); + } Review Comment: Done. `createDatabase` calls Glue directly and maps `AlreadyExistsException` → `DatabaseAlreadyExistException`; the pre-check and `findCaseInsensitiveConflict` are gone (it was unreachable anyway: `databaseExists("FOO")` was already true when `Foo` existed, so the helpful message never fired). The lowercase-storage explanation now rides as the *cause* of the `DatabaseAlreadyExistException`, where it is reachable. ########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/GlueCatalog.java: ########## @@ -0,0 +1,1850 @@ +/* + * 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.annotation.PublicEvolving; +import org.apache.flink.annotation.VisibleForTesting; +import org.apache.flink.connector.aws.config.AWSConfigConstants; +import org.apache.flink.connector.aws.util.AWSClientUtil; +import org.apache.flink.connector.aws.util.AWSGeneralUtil; +import org.apache.flink.table.api.Schema; +import org.apache.flink.table.catalog.AbstractCatalog; +import org.apache.flink.table.catalog.CatalogBaseTable; +import org.apache.flink.table.catalog.CatalogDatabase; +import org.apache.flink.table.catalog.CatalogFunction; +import org.apache.flink.table.catalog.CatalogPartition; +import org.apache.flink.table.catalog.CatalogPartitionImpl; +import org.apache.flink.table.catalog.CatalogPartitionSpec; +import org.apache.flink.table.catalog.CatalogTable; +import org.apache.flink.table.catalog.CatalogView; +import org.apache.flink.table.catalog.ObjectPath; +import org.apache.flink.table.catalog.ResolvedCatalogBaseTable; +import org.apache.flink.table.catalog.exceptions.CatalogException; +import org.apache.flink.table.catalog.exceptions.DatabaseAlreadyExistException; +import org.apache.flink.table.catalog.exceptions.DatabaseNotEmptyException; +import org.apache.flink.table.catalog.exceptions.DatabaseNotExistException; +import org.apache.flink.table.catalog.exceptions.FunctionAlreadyExistException; +import org.apache.flink.table.catalog.exceptions.FunctionNotExistException; +import org.apache.flink.table.catalog.exceptions.PartitionAlreadyExistsException; +import org.apache.flink.table.catalog.exceptions.PartitionNotExistException; +import org.apache.flink.table.catalog.exceptions.PartitionSpecInvalidException; +import org.apache.flink.table.catalog.exceptions.TableAlreadyExistException; +import org.apache.flink.table.catalog.exceptions.TableNotExistException; +import org.apache.flink.table.catalog.exceptions.TableNotPartitionedException; +import org.apache.flink.table.catalog.exceptions.TablePartitionedException; +import org.apache.flink.table.catalog.glue.operator.GlueDatabaseOperator; +import org.apache.flink.table.catalog.glue.operator.GlueFunctionOperator; +import org.apache.flink.table.catalog.glue.operator.GluePartitionOperator; +import org.apache.flink.table.catalog.glue.operator.GlueTableOperator; +import org.apache.flink.table.catalog.glue.util.GlueCatalogConstants; +import org.apache.flink.table.catalog.glue.util.GlueFlinkSchemaProperties; +import org.apache.flink.table.catalog.glue.util.GlueTableUtils; +import org.apache.flink.table.catalog.glue.util.GlueTypeConverter; +import org.apache.flink.table.catalog.stats.CatalogColumnStatistics; +import org.apache.flink.table.catalog.stats.CatalogTableStatistics; +import org.apache.flink.table.expressions.Expression; +import org.apache.flink.table.functions.FunctionIdentifier; +import org.apache.flink.util.Preconditions; +import org.apache.flink.util.StringUtils; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import software.amazon.awssdk.http.apache.ApacheHttpClient; +import software.amazon.awssdk.regions.Region; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.Partition; +import software.amazon.awssdk.services.glue.model.PartitionInput; +import software.amazon.awssdk.services.glue.model.StorageDescriptor; +import software.amazon.awssdk.services.glue.model.Table; +import software.amazon.awssdk.services.glue.model.TableInput; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Properties; +import java.util.stream.Collectors; + +/** + * GlueCatalog is an implementation of the Flink AbstractCatalog that interacts with AWS Glue. This + * class allows Flink to perform various catalog operations such as creating, deleting, and + * retrieving databases and tables from Glue. It encapsulates AWS Glue's API and provides a + * Flink-compatible interface. + * + * <p>This catalog uses GlueClient to interact with AWS Glue services, and operations related to + * databases and tables are delegated to respective helper classes like GlueDatabaseOperations and + * GlueTableOperations. + */ +@PublicEvolving +public class GlueCatalog extends AbstractCatalog { + + private static final Logger LOG = LoggerFactory.getLogger(GlueCatalog.class); + + private GlueClient glueClient; + private GlueTypeConverter glueTypeConverter; + private GlueDatabaseOperator glueDatabaseOperations; + private GlueTableOperator glueTableOperations; + private GlueFunctionOperator glueFunctionsOperations; + private GluePartitionOperator gluePartitionOperations; + private GlueTableUtils glueTableUtils; + + /** + * Constructs a GlueCatalog with a provided Glue client. + * + * @param name the name of the catalog + * @param defaultDatabase the default database for the catalog + * @param region the AWS region to be used for Glue operations + * @param glueClient Glue Client so we can decide which one to use for testing + */ + @VisibleForTesting + GlueCatalog(String name, String defaultDatabase, String region, GlueClient glueClient) { + super(name, defaultDatabase); + + // Validate region parameter + Preconditions.checkNotNull(region, "region cannot be null"); + Preconditions.checkArgument(!region.trim().isEmpty(), "region cannot be empty"); + + // Initialize GlueClient in the constructor + if (glueClient != null) { + setup(glueClient); + } else { + // If no GlueClient is provided, initialize it using the default region + GlueClient client = GlueClient.builder().region(Region.of(region)).build(); + setup(client); + } + } + + /** + * Constructs a GlueCatalog with default client configuration. + * + * @param name the name of the catalog + * @param defaultDatabase the default database for the catalog + * @param region the AWS region to be used for Glue operations + */ + public GlueCatalog(String name, String defaultDatabase, String region) { + this(name, defaultDatabase, region, new Properties()); + } + + /** + * Constructs a GlueCatalog whose Glue client is built from the given AWS client properties, + * using the same client-creation path as the other AWS connectors ({@link AWSClientUtil}). This + * makes the standard {@code aws.*} settings available to the catalog - for example {@code + * aws.credentials.provider} to select a credential mode, {@code aws.endpoint} to point at a + * Glue-compatible endpoint, and the {@code aws.http-client.*} options. + * + * @param name the name of the catalog + * @param defaultDatabase the default database for the catalog + * @param region the AWS region to be used for Glue operations + * @param glueClientProperties AWS client properties, keyed by {@link AWSConfigConstants} + */ + public GlueCatalog( + String name, String defaultDatabase, String region, Properties glueClientProperties) { + super(name, defaultDatabase); + + // Validate region parameter + Preconditions.checkNotNull(region, "region cannot be null"); + Preconditions.checkArgument(!region.trim().isEmpty(), "region cannot be empty"); + Preconditions.checkNotNull(glueClientProperties, "glueClientProperties cannot be null"); + + Properties clientProperties = new Properties(); + clientProperties.putAll(glueClientProperties); + // The explicit region argument wins over any aws.region property. + clientProperties.setProperty(AWSConfigConstants.AWS_REGION, region); + AWSGeneralUtil.validateAwsConfiguration(clientProperties); + + GlueClient client = + AWSClientUtil.createAwsSyncClient( + clientProperties, + AWSGeneralUtil.createSyncHttpClient( + clientProperties, ApacheHttpClient.builder()), + GlueClient.builder(), + GlueCatalogConstants.BASE_GLUE_USER_AGENT_PREFIX_FORMAT, + GlueCatalogConstants.GLUE_CLIENT_USER_AGENT_PREFIX); + setup(client); + } + + /** + * Private helper method to set up the GlueCatalog with a GlueClient instance. This method + * initializes all the necessary components and operators. + * + * @param glueClient the GlueClient to use for AWS Glue operations + */ + private void setup(GlueClient glueClient) { + this.glueClient = glueClient; + this.glueTypeConverter = new GlueTypeConverter(); + this.glueTableUtils = new GlueTableUtils(glueTypeConverter); + this.glueDatabaseOperations = new GlueDatabaseOperator(glueClient, getName()); + this.glueTableOperations = new GlueTableOperator(glueClient, getName()); + this.glueFunctionsOperations = new GlueFunctionOperator(glueClient, getName()); + this.gluePartitionOperations = new GluePartitionOperator(glueClient, getName()); + } + + /** + * Validates that a database exists, throwing DatabaseNotExistException if it doesn't. + * + * @param databaseName the name of the database to validate + * @throws DatabaseNotExistException if the database does not exist + * @throws CatalogException if an error occurs while checking database existence + */ + private void validateDatabaseExists(String databaseName) + throws DatabaseNotExistException, CatalogException { + Preconditions.checkArgument( + !StringUtils.isNullOrWhitespaceOnly(databaseName), + "databaseName cannot be null or empty"); + + if (!databaseExists(databaseName)) { + throw new DatabaseNotExistException(getName(), databaseName); + } + } + + /** + * Opens the GlueCatalog and initializes necessary resources. + * + * @throws CatalogException if an error occurs during the opening process + */ + @Override + public void open() throws CatalogException { + LOG.info("Opening GlueCatalog with client: {}", glueClient); + } + + /** + * Closes the GlueCatalog and releases resources. + * + * @throws CatalogException if an error occurs during the closing process + */ + @Override + public void close() throws CatalogException { + if (glueClient != null) { + LOG.info("Closing GlueCatalog client"); + // The AWS SDK close() is best-effort and does not surface exceptions, + // so no retry logic is required here. + glueClient.close(); + } + } + + /** + * Lists all the databases available in the Glue catalog. + * + * @return a list of database names + * @throws CatalogException if an error occurs while listing the databases + */ + @Override + public List<String> listDatabases() throws CatalogException { + return glueDatabaseOperations.listDatabases(); + } + + /** + * Retrieves a specific database by its name. + * + * @param databaseName the name of the database to retrieve + * @return the database if found + * @throws DatabaseNotExistException if the database does not exist + * @throws CatalogException if an error occurs while retrieving the database + */ + @Override + public CatalogDatabase getDatabase(String databaseName) + throws DatabaseNotExistException, CatalogException { + Preconditions.checkArgument( + !StringUtils.isNullOrWhitespaceOnly(databaseName), + "databaseName cannot be null or empty"); + + // Use case-insensitive database name resolution + String glueDatabaseName = findGlueDatabaseName(databaseName); + if (glueDatabaseName == null) { + throw new DatabaseNotExistException(getName(), databaseName); + } + + return glueDatabaseOperations.getDatabase(glueDatabaseName); + } + + /** + * Checks if a database exists in Glue. + * + * @param databaseName the name of the database + * @return true if the database exists, false otherwise + * @throws CatalogException if an error occurs while checking the database + */ + @Override + public boolean databaseExists(String databaseName) throws CatalogException { + Preconditions.checkArgument( + !StringUtils.isNullOrWhitespaceOnly(databaseName), + "databaseName cannot be null or empty"); + + // Use case-insensitive database name resolution + return findGlueDatabaseName(databaseName) != null; + } + + /** + * Creates a new database in Glue. + * + * @param databaseName the name of the database to create + * @param catalogDatabase the catalog database object containing database metadata + * @param ifNotExists flag indicating whether to ignore the error if the database already exists + * @throws DatabaseAlreadyExistException if the database already exists and ifNotExists is false + * @throws CatalogException if an error occurs while creating the database + */ + @Override + public void createDatabase( + String databaseName, CatalogDatabase catalogDatabase, boolean ifNotExists) + throws DatabaseAlreadyExistException, CatalogException { + Preconditions.checkArgument( + !StringUtils.isNullOrWhitespaceOnly(databaseName), + "databaseName cannot be null or empty"); + Preconditions.checkNotNull(catalogDatabase, "CatalogDatabase cannot be null"); + + // Check for exact case match first + boolean exactExists = databaseExists(databaseName); + if (exactExists && !ifNotExists) { + throw new DatabaseAlreadyExistException(getName(), databaseName); + } + if (exactExists) { + return; // Database exists with exact case, and IF NOT EXISTS is true + } + + // Check for case-insensitive collision (Glue limitation) + String conflictingDatabase = findCaseInsensitiveConflict(databaseName); + if (conflictingDatabase != null) { + String message = + String.format( + "Cannot create database '%s' because it conflicts with existing database '%s'. " + + "AWS Glue stores database names in lowercase, so '%s' and '%s' would both be stored as '%s'.", + databaseName, + conflictingDatabase, + databaseName, + conflictingDatabase, + databaseName.toLowerCase()); + throw new DatabaseAlreadyExistException( + getName(), databaseName, new CatalogException(message)); + } + + // Safe to create - no exact match and no case conflicts + glueDatabaseOperations.createDatabase(databaseName, catalogDatabase); + } + + /** + * Drops an existing database in Glue. + * + * @param databaseName the name of the database to drop + * @param ignoreIfNotExists flag to ignore the exception if the database doesn't exist + * @param cascade flag indicating whether to cascade the operation to drop related objects + * @throws DatabaseNotExistException if the database does not exist and ignoreIfNotExists is + * false + * @throws DatabaseNotEmptyException if the database contains objects and cascade is false + * @throws CatalogException if an error occurs while dropping the database + */ + @Override + public void dropDatabase(String databaseName, boolean ignoreIfNotExists, boolean cascade) + throws DatabaseNotExistException, DatabaseNotEmptyException, CatalogException { + Preconditions.checkArgument( + !StringUtils.isNullOrWhitespaceOnly(databaseName), + "databaseName cannot be null or empty"); + + if (!databaseExists(databaseName)) { + if (!ignoreIfNotExists) { + throw new DatabaseNotExistException(getName(), databaseName); + } + return; // Database doesn't exist and ignoreIfNotExists is true + } + + // Check if database is empty (contains no tables, views, or functions) + boolean isEmpty = isDatabaseEmpty(databaseName); + + if (!isEmpty && !cascade) { + throw new DatabaseNotEmptyException(getName(), databaseName); + } + + if (!isEmpty && cascade) { + // Drop all objects in the database before dropping the database + dropAllObjectsInDatabase(databaseName); + } + + // Drop the database + glueDatabaseOperations.dropGlueDatabase(databaseName); + } + + /** + * Checks if a database is empty (contains no tables, views, or functions). + * + * @param databaseName the name of the database to check + * @return true if the database is empty, false otherwise + * @throws CatalogException if an error occurs while checking the database contents + */ + private boolean isDatabaseEmpty(String databaseName) throws CatalogException { + try { + // Check for tables + List<String> tables = listTables(databaseName); + if (!tables.isEmpty()) { + return false; + } + + // Check for views + List<String> views = listViews(databaseName); + if (!views.isEmpty()) { + return false; + } + + // Check for functions + List<String> functions = listFunctions(databaseName); + if (!functions.isEmpty()) { + return false; + } + + return true; + } catch (DatabaseNotExistException e) { + // This shouldn't happen since we checked existence earlier, but handle it gracefully + throw new CatalogException("Database " + databaseName + " does not exist", e); + } + } + + /** + * Drops all objects (tables, views, functions) in a database. This is used when cascade=true in + * dropDatabase. + * + * @param databaseName the name of the database + * @throws CatalogException if an error occurs while dropping objects + */ + private void dropAllObjectsInDatabase(String databaseName) throws CatalogException { + try { + // Drop all tables + List<String> tables = listTables(databaseName); + for (String tableName : tables) { + ObjectPath tablePath = new ObjectPath(databaseName, tableName); + dropTable(tablePath, true); // Use ifExists=true to avoid exceptions + } + + // Drop all views (views are also stored as tables in Glue, so they should be handled by + // dropTable above) + // But let's be explicit and handle them separately if needed + List<String> views = listViews(databaseName); + for (String viewName : views) { + ObjectPath viewPath = new ObjectPath(databaseName, viewName); + // Views are handled as tables in Glue, so dropTable should work + dropTable(viewPath, true); + } + + // Drop all functions + List<String> functions = listFunctions(databaseName); + for (String functionName : functions) { + ObjectPath functionPath = new ObjectPath(databaseName, functionName); + dropFunction(functionPath, true); // Use ignoreIfNotExists=true to avoid exceptions + } + + LOG.info("Successfully dropped all objects in database: {}", databaseName); + } catch (DatabaseNotExistException e) { + throw new CatalogException("Database " + databaseName + " does not exist", e); + } catch (TableNotExistException | FunctionNotExistException e) { + // This could happen in concurrent scenarios, but we use ifExists/ignoreIfNotExists + // flags + LOG.warn( + "Object was already deleted while cascading drop for database: {}", + databaseName, + e); + } + } + + /** + * Lists all tables in a specified database. + * + * @param databaseName the name of the database + * @return a list of table names in the database + * @throws DatabaseNotExistException if the database does not exist + * @throws CatalogException if an error occurs while listing the tables + */ + @Override + public List<String> listTables(String databaseName) + throws DatabaseNotExistException, CatalogException { + Preconditions.checkArgument( + !StringUtils.isNullOrWhitespaceOnly(databaseName), + "databaseName cannot be null or empty"); + + validateDatabaseExists(databaseName); + + // Use the proper database name resolution + String glueDatabaseName = findGlueDatabaseName(databaseName); + if (glueDatabaseName == null) { + throw new DatabaseNotExistException(getName(), databaseName); + } + + // Return original table names with case preserved + return glueTableOperations.listTablesWithOriginalNames(glueDatabaseName); + } + + /** + * Retrieves a table from the catalog using its object path. + * + * @param objectPath the object path of the table to retrieve + * @return the corresponding CatalogBaseTable for the specified table + * @throws TableNotExistException if the table does not exist + * @throws CatalogException if an error occurs while retrieving the table + */ + @Override + public CatalogBaseTable getTable(ObjectPath objectPath) + throws TableNotExistException, CatalogException { + String originalDatabaseName = objectPath.getDatabaseName(); + String originalTableName = objectPath.getObjectName(); + + // Convert to Glue storage names - Use proper database resolution + String glueDatabaseName = findGlueDatabaseName(originalDatabaseName); + if (glueDatabaseName == null) { + throw new TableNotExistException(getName(), objectPath); + } + + // Use direct lowercase lookup first (like databases), then fall back to complex search + String glueTableName = findGlueTableName(glueDatabaseName, originalTableName); + if (glueTableName == null) { + throw new TableNotExistException(getName(), objectPath); + } + + // Get the table using Glue storage names + Table glueTable = glueTableOperations.getGlueTable(glueDatabaseName, glueTableName); + return getCatalogBaseTableFromGlueTable(glueTable); + } + + /** + * Checks if a table exists in the Glue catalog. + * + * @param objectPath the object path of the table to check + * @return true if the table exists, false otherwise + * @throws CatalogException if an error occurs while checking the table's existence + */ + @Override + public boolean tableExists(ObjectPath objectPath) throws CatalogException { + String originalDatabaseName = objectPath.getDatabaseName(); + String originalTableName = objectPath.getObjectName(); + + // Convert to Glue storage names - Use proper database resolution + String glueDatabaseName = findGlueDatabaseName(originalDatabaseName); + if (glueDatabaseName == null) { + return false; // Database doesn't exist, so table can't exist + } + + // Use efficient table name resolution + String glueTableName = findGlueTableName(glueDatabaseName, originalTableName); + return glueTableName != null; + } + + /** + * Drops a table from the Glue catalog. + * + * @param objectPath the object path of the table to drop + * @param ifExists flag indicating whether to ignore the exception if the table does not exist + * @throws TableNotExistException if the table does not exist and ifExists is false + * @throws CatalogException if an error occurs while dropping the table + */ + @Override + public void dropTable(ObjectPath objectPath, boolean ifExists) + throws TableNotExistException, CatalogException { + String originalDatabaseName = objectPath.getDatabaseName(); + String originalTableName = objectPath.getObjectName(); + + // Convert to Glue storage names - Use proper database resolution + String glueDatabaseName = findGlueDatabaseName(originalDatabaseName); + if (glueDatabaseName == null) { + if (!ifExists) { + throw new TableNotExistException(getName(), objectPath); + } + return; // Database doesn't exist, so table can't exist + } + + // Use efficient table name resolution + String glueTableName = findGlueTableName(glueDatabaseName, originalTableName); + if (glueTableName == null) { + if (!ifExists) { + throw new TableNotExistException(getName(), objectPath); + } + return; // Table doesn't exist, and IF EXISTS is true + } + + // Drop the table using Glue storage names + glueTableOperations.dropTable(glueDatabaseName, glueTableName); + } + + /** + * Creates a table in the Glue catalog. + * + * @param objectPath the object path of the table to create + * @param catalogBaseTable the table definition containing the schema and properties + * @param ifNotExists flag indicating whether to ignore the exception if the table already + * exists + * @throws NullPointerException if objectPath or catalogBaseTable is null + * @throws TableAlreadyExistException if the table already exists and ifNotExists is false + * @throws DatabaseNotExistException if the database does not exist + * @throws CatalogException if an error occurs while creating the table + */ + @Override + public void createTable( + ObjectPath objectPath, CatalogBaseTable catalogBaseTable, boolean ifNotExists) + throws TableAlreadyExistException, DatabaseNotExistException, CatalogException { + + // Validate that required parameters are not null + Preconditions.checkNotNull(objectPath, "ObjectPath cannot be null"); + Preconditions.checkNotNull(catalogBaseTable, "CatalogBaseTable cannot be null"); + + String originalDatabaseName = objectPath.getDatabaseName(); + String originalTableName = objectPath.getObjectName(); + + // Check if the database exists + validateDatabaseExists(originalDatabaseName); + + // Check for exact case match first + boolean exactExists = tableExists(objectPath); + if (exactExists && !ifNotExists) { + throw new TableAlreadyExistException(getName(), objectPath); + } + if (exactExists) { + return; // Table exists with exact case, and IF NOT EXISTS is true + } + + // Check for case-insensitive collision (Glue limitation) + String conflictingTable = findCaseInsensitiveTableConflict(objectPath); + if (conflictingTable != null) { + String message = + String.format( + "Cannot create table '%s.%s' because it conflicts with existing table '%s.%s'. " + + "AWS Glue stores table names in lowercase, so '%s' and '%s' would both be stored as '%s'.", + originalDatabaseName, + originalTableName, + originalDatabaseName, + conflictingTable, + originalTableName, + conflictingTable, + originalTableName.toLowerCase()); + throw new TableAlreadyExistException( + getName(), objectPath, new CatalogException(message)); + } + + // Get common properties + Map<String, String> tableProperties = new HashMap<>(catalogBaseTable.getOptions()); + + try { + // Process based on table type + if (catalogBaseTable.getTableKind() == CatalogBaseTable.TableKind.TABLE) { + createRegularTable(objectPath, (CatalogTable) catalogBaseTable, tableProperties); + } else if (catalogBaseTable.getTableKind() == CatalogBaseTable.TableKind.VIEW) { + createView(objectPath, (CatalogView) catalogBaseTable, tableProperties); + } else { + throw new CatalogException( + "Unsupported table kind: " + catalogBaseTable.getTableKind()); + } + LOG.info( + "Successfully created {}.{} of kind {}", + originalDatabaseName, + originalTableName, + catalogBaseTable.getTableKind()); + } catch (DatabaseNotExistException e) { + // Preserve the typed exception so callers can distinguish a missing database + // from a generic catalog failure. + throw e; + } catch (Exception e) { + throw new CatalogException( + String.format( + "Failed to create table %s.%s: %s", + originalDatabaseName, originalTableName, e.getMessage()), + e); + } + } + + /** + * Lists all views in a specified database. + * + * @param databaseName the name of the database + * @return a list of view names in the database + * @throws DatabaseNotExistException if the database does not exist + * @throws CatalogException if an error occurs while listing the views + */ + @Override + public List<String> listViews(String databaseName) + throws DatabaseNotExistException, CatalogException { + Preconditions.checkArgument( + !StringUtils.isNullOrWhitespaceOnly(databaseName), + "databaseName cannot be null or empty"); + + // Check if the database exists before listing views + validateDatabaseExists(databaseName); + + // Use proper database name resolution + String glueDatabaseName = findGlueDatabaseName(databaseName); + if (glueDatabaseName == null) { + throw new DatabaseNotExistException(getName(), databaseName); + } + + try { + // Get all tables in the database (paginated) + List<Table> allTables = glueTableOperations.getAllGlueTables(glueDatabaseName); + + // Filter tables to only include those that are of type VIEW, and return original names + List<String> viewNames = + allTables.stream() + .filter( + table -> { + String tableType = table.tableType(); + return tableType != null + && tableType.equalsIgnoreCase( + CatalogBaseTable.TableKind.VIEW.name()); + }) + .map(table -> glueTableOperations.getOriginalTableName(table)) + .collect(Collectors.toList()); + + return viewNames; + } catch (Exception e) { + LOG.error("Failed to list views in database {}: {}", databaseName, e.getMessage()); + throw new CatalogException( + String.format( + "Error listing views in database %s: %s", databaseName, e.getMessage()), + e); + } + } + + @Override + public void alterDatabase(String s, CatalogDatabase catalogDatabase, boolean b) + throws DatabaseNotExistException, CatalogException { + throw new UnsupportedOperationException( + "Altering databases is not supported by the Glue Catalog."); + } + + @Override + public void renameTable(ObjectPath objectPath, String s, boolean b) + throws TableNotExistException, TableAlreadyExistException, CatalogException { + throw new UnsupportedOperationException( + "Renaming tables is not supported by the Glue Catalog."); + } + + @Override + public void alterTable(ObjectPath objectPath, CatalogBaseTable catalogBaseTable, boolean b) + throws TableNotExistException, CatalogException { + Preconditions.checkNotNull(objectPath, "ObjectPath cannot be null"); + Preconditions.checkNotNull(catalogBaseTable, "CatalogBaseTable cannot be null"); + + // Resolve Glue storage names with the same case-insensitive resolution as getTable. + String glueDatabaseName = findGlueDatabaseName(objectPath.getDatabaseName()); + String glueTableName = + glueDatabaseName == null + ? null + : findGlueTableName(glueDatabaseName, objectPath.getObjectName()); + if (glueTableName == null) { + if (b) { + return; + } + throw new TableNotExistException(getName(), objectPath); + } + + if (catalogBaseTable.getTableKind() != CatalogBaseTable.TableKind.TABLE) { + throw new UnsupportedOperationException( + "Altering non-TABLE objects is not supported by the Glue Catalog."); + } + + // Preserve the originally declared table name (case) across the alter. + Table existingTable = glueTableOperations.getGlueTable(glueDatabaseName, glueTableName); + String originalTableName = glueTableOperations.getOriginalTableName(existingTable); + + CatalogTable catalogTable = (CatalogTable) catalogBaseTable; + Map<String, String> tableProperties = new HashMap<>(catalogTable.getOptions()); + String tableLocation = glueTableUtils.extractTableLocation(tableProperties, objectPath); + + ResolvedCatalogBaseTable<?> resolvedTable = (ResolvedCatalogBaseTable<?>) catalogTable; + List<String> partitionKeys = catalogTable.getPartitionKeys(); + + List<software.amazon.awssdk.services.glue.model.Column> dataColumns = new ArrayList<>(); + Map<String, software.amazon.awssdk.services.glue.model.Column> partitionColumnsByName = + new HashMap<>(); + for (org.apache.flink.table.catalog.Column flinkColumn : + resolvedTable.getResolvedSchema().getColumns()) { + if (!(flinkColumn instanceof org.apache.flink.table.catalog.Column.PhysicalColumn)) { + // Computed and metadata columns cannot be represented as Glue columns; + // they are persisted as flink.schema.* table parameters instead. + continue; + } + software.amazon.awssdk.services.glue.model.Column glueColumn = + glueTableUtils.mapFlinkColumnToGlueColumn(flinkColumn); + if (partitionKeys.contains(flinkColumn.getName())) { + partitionColumnsByName.put(flinkColumn.getName(), glueColumn); + } else { + dataColumns.add(glueColumn); + } + } + List<software.amazon.awssdk.services.glue.model.Column> partitionColumns = + partitionKeys.stream() + .map(partitionColumnsByName::get) + .filter(Objects::nonNull) + .collect(Collectors.toList()); + + StorageDescriptor storageDescriptor = + glueTableUtils.buildStorageDescriptor(tableProperties, dataColumns, tableLocation); + + // Persist watermarks, primary key, and computed/metadata columns as table + // parameters; Glue columns can only represent physical columns. + GlueFlinkSchemaProperties.serializeNonPhysicalSchema( + resolvedTable.getResolvedSchema(), tableProperties); + + TableInput tableInput = + glueTableOperations.buildTableInput( + originalTableName, + partitionColumns, + catalogTable, + storageDescriptor, + tableProperties); + + glueTableOperations.updateTable(glueDatabaseName, tableInput); + LOG.info( + "Successfully altered {}.{}", + objectPath.getDatabaseName(), + objectPath.getObjectName()); + } + + @Override + public List<CatalogPartitionSpec> listPartitions(ObjectPath objectPath) + throws TableNotExistException, TableNotPartitionedException, CatalogException { + GlueTableRef tableRef = resolvePartitionedTable(objectPath); + List<String> partitionKeys = tableRef.partitionKeys(); + return gluePartitionOperations + .listPartitions(tableRef.databaseName, tableRef.tableName) + .stream() + .map(partition -> toPartitionSpec(partitionKeys, partition.values())) + .collect(Collectors.toList()); + } + + @Override + public List<CatalogPartitionSpec> listPartitions( + ObjectPath objectPath, CatalogPartitionSpec catalogPartitionSpec) + throws TableNotExistException, + TableNotPartitionedException, + PartitionSpecInvalidException, + CatalogException { + GlueTableRef tableRef = resolvePartitionedTable(objectPath); + List<String> partitionKeys = tableRef.partitionKeys(); + + Map<String, String> partialSpec = + catalogPartitionSpec == null + ? Collections.emptyMap() + : catalogPartitionSpec.getPartitionSpec(); + // Flink's Catalog contract: a partial spec referencing unknown partition keys is invalid. + if (!partitionKeys.containsAll(partialSpec.keySet())) { + throw new PartitionSpecInvalidException( + getName(), partitionKeys, objectPath, catalogPartitionSpec); + } + + return gluePartitionOperations + .listPartitions(tableRef.databaseName, tableRef.tableName) + .stream() + .map(partition -> toPartitionSpec(partitionKeys, partition.values())) + .filter( + spec -> + spec.getPartitionSpec() + .entrySet() + .containsAll(partialSpec.entrySet())) + .collect(Collectors.toList()); + } + + @Override + public List<CatalogPartitionSpec> listPartitionsByFilter( + ObjectPath objectPath, List<Expression> list) + throws TableNotExistException, TableNotPartitionedException, CatalogException { + // Expression push-down to Glue partition filters is not implemented. Flink's planner + // catches UnsupportedOperationException and falls back to listPartitions(). + throw new UnsupportedOperationException( + "Listing partitions by filter expression is not supported by the Glue Catalog."); + } + + @Override + public CatalogPartition getPartition( + ObjectPath objectPath, CatalogPartitionSpec catalogPartitionSpec) + throws PartitionNotExistException, CatalogException { + Partition partition = getGluePartitionOrNull(objectPath, catalogPartitionSpec); + if (partition == null) { + throw new PartitionNotExistException(getName(), objectPath, catalogPartitionSpec); + } + + Map<String, String> properties = new HashMap<>(); + if (partition.parameters() != null) { + properties.putAll(partition.parameters()); + } + if (partition.storageDescriptor() != null + && partition.storageDescriptor().location() != null) { + properties.put( + GlueCatalogConstants.PARTITION_LOCATION, + partition.storageDescriptor().location()); + } + return new CatalogPartitionImpl(properties, null); + } + + @Override + public boolean partitionExists(ObjectPath objectPath, CatalogPartitionSpec catalogPartitionSpec) + throws CatalogException { + try { + return getGluePartitionOrNull(objectPath, catalogPartitionSpec) != null; + } catch (PartitionNotExistException e) { + return false; + } + } + + @Override + public void createPartition( + ObjectPath objectPath, + CatalogPartitionSpec catalogPartitionSpec, + CatalogPartition catalogPartition, + boolean ifNotExists) + throws TableNotExistException, + TableNotPartitionedException, + PartitionSpecInvalidException, + PartitionAlreadyExistsException, + CatalogException { + GlueTableRef tableRef = resolvePartitionedTable(objectPath); + List<String> partitionKeys = tableRef.partitionKeys(); + + Map<String, String> spec = catalogPartitionSpec.getPartitionSpec(); + if (!spec.keySet().equals(new HashSet<>(partitionKeys))) { + throw new PartitionSpecInvalidException( + getName(), partitionKeys, objectPath, catalogPartitionSpec); + } + List<String> partitionValues = + partitionKeys.stream().map(spec::get).collect(Collectors.toList()); + + Map<String, String> partitionProperties = + catalogPartition == null + ? new HashMap<>() + : new HashMap<>(catalogPartition.getProperties()); + String location = partitionProperties.remove(GlueCatalogConstants.PARTITION_LOCATION); + StorageDescriptor.Builder sdBuilder = + tableRef.glueTable.storageDescriptor() != null + ? tableRef.glueTable.storageDescriptor().toBuilder() + : StorageDescriptor.builder(); + if (location != null) { + sdBuilder.location(location); + } else if (tableRef.glueTable.storageDescriptor() != null + && tableRef.glueTable.storageDescriptor().location() != null) { + sdBuilder.location( + buildDefaultPartitionLocation( + tableRef.glueTable.storageDescriptor().location(), + partitionKeys, + partitionValues)); + } + + PartitionInput partitionInput = + PartitionInput.builder() + .values(partitionValues) + .storageDescriptor(sdBuilder.build()) + .parameters(partitionProperties) + .build(); + + try { + gluePartitionOperations.createPartition( + tableRef.databaseName, tableRef.tableName, partitionInput); + } catch (software.amazon.awssdk.services.glue.model.AlreadyExistsException e) { + if (!ifNotExists) { + throw new PartitionAlreadyExistsException( + getName(), objectPath, catalogPartitionSpec); + } + } + } + + @Override + public void dropPartition( + ObjectPath objectPath, + CatalogPartitionSpec catalogPartitionSpec, + boolean ignoreIfNotExists) + throws PartitionNotExistException, CatalogException { + try { + GlueTableRef tableRef = resolvePartitionedTable(objectPath); + List<String> partitionValues = + toOrderedPartitionValues(tableRef, objectPath, catalogPartitionSpec); + gluePartitionOperations.dropPartition( + tableRef.databaseName, tableRef.tableName, partitionValues); + } catch (software.amazon.awssdk.services.glue.model.EntityNotFoundException + | PartitionNotExistException + | TableNotExistException + | TableNotPartitionedException e) { + if (!ignoreIfNotExists) { + // Chain the original exception so the real cause (missing table, table not + // partitioned, or missing partition) stays visible for debugging. + throw new PartitionNotExistException( + getName(), objectPath, catalogPartitionSpec, e); + } + } + } + + @Override + public void alterPartition( + ObjectPath objectPath, + CatalogPartitionSpec catalogPartitionSpec, + CatalogPartition catalogPartition, + boolean ignoreIfNotExists) + throws PartitionNotExistException, CatalogException { + try { + GlueTableRef tableRef = resolvePartitionedTable(objectPath); + List<String> partitionValues = + toOrderedPartitionValues(tableRef, objectPath, catalogPartitionSpec); + Partition existing = + gluePartitionOperations.getPartition( + tableRef.databaseName, tableRef.tableName, partitionValues); + if (existing == null) { + throw new PartitionNotExistException(getName(), objectPath, catalogPartitionSpec); + } + + Map<String, String> partitionProperties = + catalogPartition == null + ? new HashMap<>() + : new HashMap<>(catalogPartition.getProperties()); + String location = partitionProperties.remove(GlueCatalogConstants.PARTITION_LOCATION); + StorageDescriptor.Builder sdBuilder = + existing.storageDescriptor() != null + ? existing.storageDescriptor().toBuilder() + : StorageDescriptor.builder(); + if (location != null) { + sdBuilder.location(location); + } + + PartitionInput partitionInput = + PartitionInput.builder() + .values(partitionValues) + .storageDescriptor(sdBuilder.build()) + .parameters(partitionProperties) + .build(); + + gluePartitionOperations.updatePartition( + tableRef.databaseName, tableRef.tableName, partitionValues, partitionInput); + } catch (software.amazon.awssdk.services.glue.model.EntityNotFoundException + | TableNotExistException + | TableNotPartitionedException e) { + if (!ignoreIfNotExists) { + throw new PartitionNotExistException(getName(), objectPath, catalogPartitionSpec); + } + } + } + + /** Resolved Glue storage names plus the fetched Glue table for partition operations. */ + private static final class GlueTableRef { + private final String databaseName; + private final String tableName; + private final Table glueTable; + + private GlueTableRef(String databaseName, String tableName, Table glueTable) { + this.databaseName = databaseName; + this.tableName = tableName; + this.glueTable = glueTable; + } + + private List<String> partitionKeys() { + return org.apache.flink.table.catalog.glue.util.GlueTableUtils.getPartitionKeyNames( + glueTable); + } + } + + /** + * Resolves the Glue storage names for the given path (using the same case-insensitive + * resolution as {@link #getTable(ObjectPath)}) and validates that the table is partitioned. + */ + private GlueTableRef resolvePartitionedTable(ObjectPath objectPath) + throws TableNotExistException, TableNotPartitionedException, CatalogException { + String glueDatabaseName = findGlueDatabaseName(objectPath.getDatabaseName()); + if (glueDatabaseName == null) { + throw new TableNotExistException(getName(), objectPath); + } + String glueTableName = findGlueTableName(glueDatabaseName, objectPath.getObjectName()); + if (glueTableName == null) { + throw new TableNotExistException(getName(), objectPath); + } + Table glueTable = glueTableOperations.getGlueTable(glueDatabaseName, glueTableName); + if (glueTable.partitionKeys() == null || glueTable.partitionKeys().isEmpty()) { + throw new TableNotPartitionedException(getName(), objectPath); + } + return new GlueTableRef(glueDatabaseName, glueTableName, glueTable); + } + + /** + * Resolves a full partition spec into partition values ordered by the table's partition keys. + * An incomplete or mismatched spec identifies a partition that cannot exist, so {@link + * PartitionNotExistException} is thrown (matching the Flink Catalog contract for + * partition-addressing methods that do not declare PartitionSpecInvalidException). + */ + private List<String> toOrderedPartitionValues( + GlueTableRef tableRef, ObjectPath objectPath, CatalogPartitionSpec catalogPartitionSpec) + throws PartitionNotExistException { + List<String> partitionKeys = tableRef.partitionKeys(); + Map<String, String> spec = catalogPartitionSpec.getPartitionSpec(); + if (!spec.keySet().equals(new HashSet<>(partitionKeys))) { + throw new PartitionNotExistException(getName(), objectPath, catalogPartitionSpec); + } + return partitionKeys.stream().map(spec::get).collect(Collectors.toList()); + } + + private Partition getGluePartitionOrNull( + ObjectPath objectPath, CatalogPartitionSpec catalogPartitionSpec) + throws PartitionNotExistException, CatalogException { + try { + GlueTableRef tableRef = resolvePartitionedTable(objectPath); + List<String> partitionValues = + toOrderedPartitionValues(tableRef, objectPath, catalogPartitionSpec); + return gluePartitionOperations.getPartition( + tableRef.databaseName, tableRef.tableName, partitionValues); + } catch (TableNotExistException | TableNotPartitionedException e) { + throw new PartitionNotExistException(getName(), objectPath, catalogPartitionSpec); + } + } + + private static CatalogPartitionSpec toPartitionSpec( + List<String> partitionKeys, List<String> values) { + Map<String, String> spec = new LinkedHashMap<>(); + for (int i = 0; i < partitionKeys.size() && i < values.size(); i++) { + spec.put(partitionKeys.get(i), values.get(i)); + } + return new CatalogPartitionSpec(spec); + } + + private static String buildDefaultPartitionLocation( + String tableLocation, List<String> partitionKeys, List<String> partitionValues) { + StringBuilder location = new StringBuilder(tableLocation); + for (int i = 0; i < partitionKeys.size(); i++) { + location.append('/') + .append(partitionKeys.get(i)) + .append('=') + .append(partitionValues.get(i)); + } + return location.toString(); + } + + /** + * Normalizes an object path according to catalog-specific normalization rules. For functions, + * this ensures consistent case handling in function names. + * + * @param path the object path to normalize + * @return the normalized object path + * @throws NullPointerException if path is null + */ + private ObjectPath normalize(ObjectPath path) { + Preconditions.checkNotNull(path, "ObjectPath cannot be null"); + + return new ObjectPath( + path.getDatabaseName(), FunctionIdentifier.normalizeName(path.getObjectName())); + } + + /** + * Resolves a normalized function path to the path used against Glue: the database part is + * translated from the Flink database name to the Glue storage name (Glue stores database names + * in lowercase). The function name is left as-is because it was already normalized. + * + * @param normalizedPath the normalized function path (Flink database name) + * @return the function path addressed by Glue storage names + * @throws CatalogException if the database cannot be resolved + */ + private ObjectPath toGlueFunctionPath(ObjectPath normalizedPath) throws CatalogException { + String glueDatabaseName = findGlueDatabaseName(normalizedPath.getDatabaseName()); + if (glueDatabaseName == null) { + throw new CatalogException("Database not found: " + normalizedPath.getDatabaseName()); + } + return new ObjectPath(glueDatabaseName, normalizedPath.getObjectName()); + } + + /** + * Lists all functions in a specified database. + * + * @param databaseName the name of the database + * @return a list of function names in the database + * @throws DatabaseNotExistException if the database does not exist + * @throws CatalogException if an error occurs while listing the functions + */ + @Override + public List<String> listFunctions(String databaseName) + throws DatabaseNotExistException, CatalogException { + Preconditions.checkArgument( + !StringUtils.isNullOrWhitespaceOnly(databaseName), + "databaseName cannot be null or empty"); + + validateDatabaseExists(databaseName); + + // Use proper database name resolution (Glue stores database names in lowercase) + String glueDatabaseName = findGlueDatabaseName(databaseName); + if (glueDatabaseName == null) { + throw new DatabaseNotExistException(getName(), databaseName); + } + + try { + List<String> functions = glueFunctionsOperations.listGlueFunctions(glueDatabaseName); + return functions; + } catch (CatalogException e) { + LOG.error("Failed to list functions in database {}: {}", databaseName, e.getMessage()); + throw new CatalogException( + String.format( + "Error listing functions in database %s: %s", + databaseName, e.getMessage()), + e); + } + } + + /** + * Retrieves a function from the catalog. + * + * @param functionPath the object path of the function to retrieve + * @return the corresponding CatalogFunction + * @throws FunctionNotExistException if the function does not exist + * @throws CatalogException if an error occurs while retrieving the function + */ + @Override + public CatalogFunction getFunction(ObjectPath functionPath) + throws FunctionNotExistException, CatalogException { + // Normalize the path for case-insensitive handling + ObjectPath normalizedPath = normalize(functionPath); + + if (!databaseExists(normalizedPath.getDatabaseName())) { + // A function in a non-existent database does not exist: report it per the + // Catalog contract so the planner can fall back to built-in functions instead + // of failing SQL validation (matches Hive/GenericInMemoryCatalog behaviour). + throw new FunctionNotExistException(getName(), normalizedPath); + } + + boolean exists = functionExists(normalizedPath); + + if (!exists) { + throw new FunctionNotExistException(getName(), normalizedPath); + } + + try { + return glueFunctionsOperations.getGlueFunction(toGlueFunctionPath(normalizedPath)); + } catch (CatalogException e) { + throw new CatalogException( + String.format("Failed to get function %s", normalizedPath.getFullName()), e); + } + } + + /** + * Checks if a function exists in the catalog. + * + * @param functionPath the object path of the function to check + * @return true if the function exists, false otherwise + * @throws CatalogException if an error occurs while checking the function's existence + */ + @Override + public boolean functionExists(ObjectPath functionPath) throws CatalogException { + // Normalize the path for case-insensitive handling + ObjectPath normalizedPath = normalize(functionPath); + + if (!databaseExists(normalizedPath.getDatabaseName())) { + return false; + } + + try { + return glueFunctionsOperations.glueFunctionExists(toGlueFunctionPath(normalizedPath)); + } catch (CatalogException e) { + throw new CatalogException( + String.format( + "Failed to check if function %s exists", normalizedPath.getFullName()), + e); + } + } + + /** + * Creates a function in the catalog. + * + * @param functionPath the object path of the function to create + * @param function the function definition + * @param ignoreIfExists flag indicating whether to ignore the exception if the function already + * exists + * @throws FunctionAlreadyExistException if the function already exists and ignoreIfExists is + * false + * @throws DatabaseNotExistException if the database does not exist + * @throws CatalogException if an error occurs while creating the function + */ + @Override + public void createFunction( + ObjectPath functionPath, CatalogFunction function, boolean ignoreIfExists) + throws FunctionAlreadyExistException, DatabaseNotExistException, CatalogException { + + // Normalize the path for case-insensitive handling + ObjectPath normalizedPath = normalize(functionPath); + + validateDatabaseExists(normalizedPath.getDatabaseName()); + + boolean exists = functionExists(normalizedPath); + + if (exists && !ignoreIfExists) { + throw new FunctionAlreadyExistException(getName(), normalizedPath); + } else if (exists) { + return; + } + + try { + glueFunctionsOperations.createGlueFunction( + toGlueFunctionPath(normalizedPath), function); + } catch (CatalogException e) { + throw new CatalogException( + String.format("Failed to create function %s", normalizedPath.getFullName()), e); + } + } + + /** + * Alters a function in the catalog. + * + * @param functionPath the object path of the function to alter + * @param newFunction the new function definition + * @param ignoreIfNotExists flag indicating whether to ignore the exception if the function does + * not exist + * @throws FunctionNotExistException if the function does not exist and ignoreIfNotExists is + * false + * @throws CatalogException if an error occurs while altering the function + */ + @Override + public void alterFunction( + ObjectPath functionPath, CatalogFunction newFunction, boolean ignoreIfNotExists) + throws FunctionNotExistException, CatalogException { + + // Normalize the path for case-insensitive handling + ObjectPath normalizedPath = normalize(functionPath); + + // Check if function exists without throwing an exception first + boolean functionExists = functionExists(normalizedPath); + + if (!functionExists) { + if (ignoreIfNotExists) { + return; + } else { + throw new FunctionNotExistException(getName(), normalizedPath); + } + } + + try { + // Check for type compatibility of function + CatalogFunction existingFunction = getFunction(normalizedPath); + if (existingFunction.getClass() != newFunction.getClass()) { + throw new CatalogException( + String.format( + "Function types don't match. Existing function is '%s' and new function is '%s'.", + existingFunction.getClass().getName(), + newFunction.getClass().getName())); + } + + // Proceed with alteration + glueFunctionsOperations.alterGlueFunction( + toGlueFunctionPath(normalizedPath), newFunction); + } catch (CatalogException e) { + throw new CatalogException( + String.format("Failed to alter function %s", normalizedPath.getFullName()), e); + } + } + + /** + * Drops a function from the catalog. + * + * @param functionPath the object path of the function to drop + * @param ignoreIfNotExists flag indicating whether to ignore the exception if the function does + * not exist + * @throws FunctionNotExistException if the function does not exist and ignoreIfNotExists is + * false + * @throws CatalogException if an error occurs while dropping the function + */ + @Override + public void dropFunction(ObjectPath functionPath, boolean ignoreIfNotExists) + throws FunctionNotExistException, CatalogException { + + // Normalize the path for case-insensitive handling + ObjectPath normalizedPath = normalize(functionPath); + + if (!databaseExists(normalizedPath.getDatabaseName())) { + if (ignoreIfNotExists) { + return; + } + throw new FunctionNotExistException(getName(), normalizedPath); + } + + boolean exists = functionExists(normalizedPath); + + if (!exists) { + if (ignoreIfNotExists) { + return; + } else { + throw new FunctionNotExistException(getName(), normalizedPath); + } + } + + try { + // Function exists, proceed with dropping it + glueFunctionsOperations.dropGlueFunction(toGlueFunctionPath(normalizedPath)); + } catch (CatalogException e) { + throw new CatalogException( + String.format("Failed to drop function %s", normalizedPath.getFullName()), e); + } + } + + @Override + public CatalogTableStatistics getTableStatistics(ObjectPath objectPath) + throws TableNotExistException, CatalogException { + return CatalogTableStatistics.UNKNOWN; + } + + @Override + public CatalogColumnStatistics getTableColumnStatistics(ObjectPath objectPath) + throws TableNotExistException, CatalogException { + return CatalogColumnStatistics.UNKNOWN; + } + + @Override + public CatalogTableStatistics getPartitionStatistics( + ObjectPath objectPath, CatalogPartitionSpec catalogPartitionSpec) + throws PartitionNotExistException, CatalogException { + return CatalogTableStatistics.UNKNOWN; + } + + @Override + public CatalogColumnStatistics getPartitionColumnStatistics( + ObjectPath objectPath, CatalogPartitionSpec catalogPartitionSpec) + throws PartitionNotExistException, CatalogException { + return CatalogColumnStatistics.UNKNOWN; + } + + @Override + public void alterTableStatistics( + ObjectPath objectPath, CatalogTableStatistics catalogTableStatistics, boolean b) + throws TableNotExistException, CatalogException { + throw new UnsupportedOperationException( + "Altering table statistics is not supported by the Glue Catalog."); + } + + @Override + public void alterTableColumnStatistics( + ObjectPath objectPath, CatalogColumnStatistics catalogColumnStatistics, boolean b) + throws TableNotExistException, CatalogException, TablePartitionedException { + throw new UnsupportedOperationException( + "Altering table column statistics is not supported by the Glue Catalog."); + } + + @Override + public void alterPartitionStatistics( + ObjectPath objectPath, + CatalogPartitionSpec catalogPartitionSpec, + CatalogTableStatistics catalogTableStatistics, + boolean b) + throws PartitionNotExistException, CatalogException { + throw new UnsupportedOperationException( + "Altering partition statistics is not supported by the Glue Catalog."); + } + + @Override + public void alterPartitionColumnStatistics( + ObjectPath objectPath, + CatalogPartitionSpec catalogPartitionSpec, + CatalogColumnStatistics catalogColumnStatistics, + boolean b) + throws PartitionNotExistException, CatalogException { + throw new UnsupportedOperationException( + "Altering partition column statistics is not supported by the Glue Catalog."); + } + + // ============================ Private Methods ============================ + /** + * Converts an AWS Glue Table to a Flink CatalogBaseTable, supporting both tables and views. + * + * @param glueTable the AWS Glue table to convert + * @return the corresponding Flink CatalogBaseTable (either CatalogTable or CatalogView) + * @throws CatalogException if the table type is unknown or conversion fails + */ + private CatalogBaseTable getCatalogBaseTableFromGlueTable(Table glueTable) { + + try { + // Parse schema from Glue table structure + Schema schemaInfo = glueTableUtils.getSchemaFromGlueTable(glueTable); + + // Extract partition keys (restoring original case from table parameters) + List<String> partitionKeys = + org.apache.flink.table.catalog.glue.util.GlueTableUtils.getPartitionKeyNames( + glueTable); + + // Collect all properties + Map<String, String> properties = new HashMap<>(); + + // Add table parameters, filtering out internal metadata + if (glueTable.parameters() != null) { + for (Map.Entry<String, String> entry : glueTable.parameters().entrySet()) { + String key = entry.getKey(); + // Filter out our internal metadata parameters + if (!GlueCatalogConstants.ORIGINAL_TABLE_NAME.equals(key) + && !GlueCatalogConstants.ORIGINAL_DATABASE_NAME.equals(key) + && !GlueCatalogConstants.ORIGINAL_PARTITION_KEYS.equals(key) + && !GlueFlinkSchemaProperties.isSchemaParameter(key)) { + properties.put(key, entry.getValue()); + } + } + } + + // Add owner if present + if (glueTable.owner() != null) { + properties.put(GlueCatalogConstants.TABLE_OWNER, glueTable.owner()); + } + + // Add storage parameters if present + if (glueTable.storageDescriptor() != null) { + if (glueTable.storageDescriptor().hasParameters()) { + properties.putAll(glueTable.storageDescriptor().parameters()); + } + + // Add input/output formats if present + if (glueTable.storageDescriptor().inputFormat() != null) { + properties.put( + GlueCatalogConstants.TABLE_INPUT_FORMAT, + glueTable.storageDescriptor().inputFormat()); + } + + if (glueTable.storageDescriptor().outputFormat() != null) { + properties.put( + GlueCatalogConstants.TABLE_OUTPUT_FORMAT, + glueTable.storageDescriptor().outputFormat()); + } + } + + // Check table type and create appropriate catalog object + String tableType = glueTable.tableType(); + if (tableType == null) { + LOG.warn("Table type is null for table {}, defaulting to TABLE", glueTable.name()); + tableType = CatalogBaseTable.TableKind.TABLE.name(); + } + + if (tableType.equalsIgnoreCase(CatalogBaseTable.TableKind.TABLE.name())) { + return CatalogTable.newBuilder() + .schema(schemaInfo) + .comment(glueTable.description()) + .partitionKeys(partitionKeys) + .options(properties) + .build(); + } else if (tableType.equalsIgnoreCase(CatalogBaseTable.TableKind.VIEW.name())) { + String originalQuery = glueTable.viewOriginalText(); + String expandedQuery = glueTable.viewExpandedText(); + + if (originalQuery == null) { + throw new CatalogException( + String.format( + "View '%s' is missing its original query text", + glueTable.name())); + } + + // If expanded query is null, use original query + if (expandedQuery == null) { + expandedQuery = originalQuery; + } + + return CatalogView.of( + schemaInfo, + glueTable.description(), + originalQuery, + expandedQuery, + properties); + } else { + throw new CatalogException( + String.format("Unknown table type: %s from Glue Catalog.", tableType)); + } + } catch (Exception e) { + throw new CatalogException( + String.format( + "Failed to convert Glue table '%s' to Flink table: %s", + glueTable.name(), e.getMessage()), + e); + } + } + + /** + * Creates a regular table in the Glue catalog. + * + * @param objectPath the object path of the table + * @param catalogTable the table definition + * @param tableProperties the table properties + * @throws CatalogException if an error occurs during table creation + */ + private void createRegularTable( + ObjectPath objectPath, CatalogTable catalogTable, Map<String, String> tableProperties) + throws CatalogException, DatabaseNotExistException { + + String databaseName = objectPath.getDatabaseName(); + String tableName = objectPath.getObjectName(); + + // Extract table location + String tableLocation = glueTableUtils.extractTableLocation(tableProperties, objectPath); + + // Resolve the schema and map Flink columns to Glue columns, splitting data columns + // (stored in the storage descriptor) from partition columns (stored on the + // TableInput itself) — Glue models partition keys outside the storage descriptor. + ResolvedCatalogBaseTable<?> resolvedTable = (ResolvedCatalogBaseTable<?>) catalogTable; + List<String> partitionKeys = catalogTable.getPartitionKeys(); + + List<software.amazon.awssdk.services.glue.model.Column> dataColumns = new ArrayList<>(); + Map<String, software.amazon.awssdk.services.glue.model.Column> partitionColumnsByName = + new HashMap<>(); + for (org.apache.flink.table.catalog.Column flinkColumn : + resolvedTable.getResolvedSchema().getColumns()) { + if (!(flinkColumn instanceof org.apache.flink.table.catalog.Column.PhysicalColumn)) { + // Computed and metadata columns cannot be represented as Glue columns; + // they are persisted as flink.schema.* table parameters instead. + continue; + } + software.amazon.awssdk.services.glue.model.Column glueColumn = + glueTableUtils.mapFlinkColumnToGlueColumn(flinkColumn); + if (partitionKeys.contains(flinkColumn.getName())) { + partitionColumnsByName.put(flinkColumn.getName(), glueColumn); + } else { + dataColumns.add(glueColumn); + } + } + + // Preserve the declared partition-key order. + List<software.amazon.awssdk.services.glue.model.Column> partitionColumns = + partitionKeys.stream() + .map(partitionColumnsByName::get) + .filter(Objects::nonNull) + .collect(Collectors.toList()); + + StorageDescriptor storageDescriptor = + glueTableUtils.buildStorageDescriptor(tableProperties, dataColumns, tableLocation); + + // Persist watermarks, primary key, and computed/metadata columns as table + // parameters; Glue columns can only represent physical columns. + GlueFlinkSchemaProperties.serializeNonPhysicalSchema( + resolvedTable.getResolvedSchema(), tableProperties); + + // Pass original table name to preserve case + TableInput tableInput = + glueTableOperations.buildTableInput( + tableName, + partitionColumns, + catalogTable, + storageDescriptor, + tableProperties); + + // Use proper database name resolution + String glueDatabaseName = findGlueDatabaseName(databaseName); + if (glueDatabaseName == null) { + throw new DatabaseNotExistException(getName(), databaseName); + } + glueTableOperations.createTable(glueDatabaseName, tableInput); + } + + /** + * Creates a view in the Glue catalog. + * + * @param objectPath the object path of the view + * @param catalogView the view definition + * @param tableProperties the view properties + * @throws CatalogException if an error occurs during view creation + */ + private void createView( Review Comment: Yes — views now go through the same `buildTableInput` path as tables, so `TIMESTAMP(3)` precision, NOT NULL and comments survive for view columns too (`testViewColumnTypesSurviveRoundTrip`). -- 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]
