fmorillo7694 commented on code in PR #206: URL: https://github.com/apache/flink-connector-aws/pull/206#discussion_r4125191663
########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/operator/GluePartitionOperator.java: ########## @@ -0,0 +1,209 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.catalog.glue.operator; + +import org.apache.flink.table.catalog.exceptions.CatalogException; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.CreatePartitionRequest; +import software.amazon.awssdk.services.glue.model.DeletePartitionRequest; +import software.amazon.awssdk.services.glue.model.EntityNotFoundException; +import software.amazon.awssdk.services.glue.model.GetPartitionRequest; +import software.amazon.awssdk.services.glue.model.GetPartitionsRequest; +import software.amazon.awssdk.services.glue.model.GetPartitionsResponse; +import software.amazon.awssdk.services.glue.model.GlueException; +import software.amazon.awssdk.services.glue.model.Partition; +import software.amazon.awssdk.services.glue.model.PartitionInput; +import software.amazon.awssdk.services.glue.model.UpdatePartitionRequest; + +import java.util.ArrayList; +import java.util.List; + +/** + * Handles partition operations for the Glue catalog. Provides functionality for listing, + * retrieving, creating, updating and deleting partitions in AWS Glue, with pagination handled + * internally. + */ +public class GluePartitionOperator extends GlueOperator { + + private static final Logger LOG = LoggerFactory.getLogger(GluePartitionOperator.class); + + /** + * Constructor for GluePartitionOperator. + * + * @param glueClient The Glue client to use for partition operations. + * @param catalogName The name of the catalog. + */ + public GluePartitionOperator(GlueClient glueClient, String catalogName) { + super(glueClient, catalogName); + } + + /** + * Lists all partitions of a table, following pagination. + * + * @param databaseName The name of the database containing the table. + * @param tableName The name of the table. + * @return All partitions of the table. + * @throws CatalogException if an error occurs while listing partitions. + */ + public List<Partition> listPartitions(String databaseName, String tableName) { + try { + List<Partition> partitions = new ArrayList<>(); + String nextToken = null; + do { + GetPartitionsRequest.Builder requestBuilder = + GetPartitionsRequest.builder() + .databaseName(databaseName) + .tableName(tableName); + if (nextToken != null) { + requestBuilder.nextToken(nextToken); + } + GetPartitionsResponse response = glueClient.getPartitions(requestBuilder.build()); + if (response.partitions() != null) { + partitions.addAll(response.partitions()); + } + nextToken = response.nextToken(); + } while (nextToken != null); + return partitions; + } catch (GlueException e) { + throw new CatalogException("Error listing partitions: " + e.getMessage(), e); + } Review Comment: Done — `listPartitions` iterates `getPartitionsPaginator(...)`. Note the Glue paginators in SDK 2.40.3 expose response-level iteration only (no per-item `partitions()` stream), so it iterates pages. ########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/operator/GluePartitionOperator.java: ########## @@ -0,0 +1,209 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.catalog.glue.operator; + +import org.apache.flink.table.catalog.exceptions.CatalogException; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.CreatePartitionRequest; +import software.amazon.awssdk.services.glue.model.DeletePartitionRequest; +import software.amazon.awssdk.services.glue.model.EntityNotFoundException; +import software.amazon.awssdk.services.glue.model.GetPartitionRequest; +import software.amazon.awssdk.services.glue.model.GetPartitionsRequest; +import software.amazon.awssdk.services.glue.model.GetPartitionsResponse; +import software.amazon.awssdk.services.glue.model.GlueException; +import software.amazon.awssdk.services.glue.model.Partition; +import software.amazon.awssdk.services.glue.model.PartitionInput; +import software.amazon.awssdk.services.glue.model.UpdatePartitionRequest; + +import java.util.ArrayList; +import java.util.List; + +/** + * Handles partition operations for the Glue catalog. Provides functionality for listing, + * retrieving, creating, updating and deleting partitions in AWS Glue, with pagination handled + * internally. + */ +public class GluePartitionOperator extends GlueOperator { + + private static final Logger LOG = LoggerFactory.getLogger(GluePartitionOperator.class); + + /** + * Constructor for GluePartitionOperator. + * + * @param glueClient The Glue client to use for partition operations. + * @param catalogName The name of the catalog. + */ + public GluePartitionOperator(GlueClient glueClient, String catalogName) { + super(glueClient, catalogName); + } + + /** + * Lists all partitions of a table, following pagination. + * + * @param databaseName The name of the database containing the table. + * @param tableName The name of the table. + * @return All partitions of the table. + * @throws CatalogException if an error occurs while listing partitions. + */ + public List<Partition> listPartitions(String databaseName, String tableName) { + try { + List<Partition> partitions = new ArrayList<>(); + String nextToken = null; + do { + GetPartitionsRequest.Builder requestBuilder = + GetPartitionsRequest.builder() + .databaseName(databaseName) + .tableName(tableName); + if (nextToken != null) { + requestBuilder.nextToken(nextToken); + } + GetPartitionsResponse response = glueClient.getPartitions(requestBuilder.build()); + if (response.partitions() != null) { + partitions.addAll(response.partitions()); + } + nextToken = response.nextToken(); + } while (nextToken != null); + return partitions; + } catch (GlueException e) { + throw new CatalogException("Error listing partitions: " + e.getMessage(), e); + } + } + + /** + * Gets a single partition by its ordered partition values. + * + * @param databaseName The name of the database containing the table. + * @param tableName The name of the table. + * @param partitionValues The partition values, ordered by the table's partition keys. + * @return The partition, or {@code null} if it does not exist. + * @throws CatalogException if an error occurs while getting the partition. + */ + public Partition getPartition( + String databaseName, String tableName, List<String> partitionValues) { + try { + GetPartitionRequest request = + GetPartitionRequest.builder() + .databaseName(databaseName) + .tableName(tableName) + .partitionValues(partitionValues) + .build(); + return glueClient.getPartition(request).partition(); + } catch (EntityNotFoundException e) { + return null; Review Comment: Done — all partition operations funnel through one `translate(...)` that logs at error with the action and target (`listing partitions of db.tbl`, `creating partition [eu] of db.tbl`) and distinguishes `AccessDenied`/`InvalidInput`. ########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/operator/GlueTableOperator.java: ########## @@ -0,0 +1,509 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.catalog.glue.operator; + +import org.apache.flink.table.catalog.CatalogTable; +import org.apache.flink.table.catalog.ObjectPath; +import org.apache.flink.table.catalog.exceptions.CatalogException; +import org.apache.flink.table.catalog.exceptions.TableNotExistException; +import org.apache.flink.table.catalog.glue.util.GlueCatalogConstants; +import org.apache.flink.table.catalog.glue.util.GlueTableUtils; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.AlreadyExistsException; +import software.amazon.awssdk.services.glue.model.Column; +import software.amazon.awssdk.services.glue.model.CreateTableRequest; +import software.amazon.awssdk.services.glue.model.DeleteTableRequest; +import software.amazon.awssdk.services.glue.model.EntityNotFoundException; +import software.amazon.awssdk.services.glue.model.GetTableRequest; +import software.amazon.awssdk.services.glue.model.GetTablesRequest; +import software.amazon.awssdk.services.glue.model.GetTablesResponse; +import software.amazon.awssdk.services.glue.model.GlueException; +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 software.amazon.awssdk.services.glue.model.UpdateTableRequest; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.regex.Pattern; + +/** + * Handles all table-related operations for the Glue catalog. Provides functionality for checking + * existence, listing, creating, getting, and dropping tables in AWS Glue. + */ +public class GlueTableOperator extends GlueOperator { + + /** Logger for logging table operations. */ + private static final Logger LOG = LoggerFactory.getLogger(GlueTableOperator.class); + + /** + * Pattern for validating table names. AWS Glue supports alphanumeric characters and + * underscores. We preserve original case in metadata while storing lowercase in Glue. + */ + private static final Pattern VALID_NAME_PATTERN = Pattern.compile("^[a-zA-Z0-9_]+$"); + + /** + * Constructor for GlueTableOperations. Initializes the Glue client and catalog name. + * + * @param glueClient The Glue client to interact with AWS Glue. + * @param catalogName The name of the catalog. + */ + public GlueTableOperator(GlueClient glueClient, String catalogName) { + super(glueClient, catalogName); + } + + /** + * Validates that a table name contains only allowed characters. AWS Glue supports alphanumeric + * characters and underscores. Case is preserved in metadata while Glue stores lowercase + * internally. + * + * @param tableName The table name to validate + * @throws CatalogException if the table name contains invalid characters + */ + private void validateTableName(String tableName) { + if (tableName == null || tableName.isEmpty()) { + throw new CatalogException("Table name cannot be null or empty"); + } + + if (!VALID_NAME_PATTERN.matcher(tableName).matches()) { + throw new CatalogException( + "Table name can only contain letters, numbers, and underscores. " + + "Original case is preserved in metadata while AWS Glue stores lowercase internally."); + } + } + + /** + * Checks whether a table exists in the Glue catalog by Glue storage names. + * + * @param glueDatabaseName The Glue storage name of the database where the table should exist. + * @param glueTableName The Glue storage name of the table to check. + * @return true if the table exists, false otherwise. + */ + public boolean glueTableExists(String glueDatabaseName, String glueTableName) { + try { + return glueClient + .getTable( + builder -> + builder.databaseName(glueDatabaseName) + .name(glueTableName)) + .table() + != null; + } catch (EntityNotFoundException e) { + return false; + } catch (GlueException e) { + throw new CatalogException( + "Error checking table existence: " + glueDatabaseName + "." + glueTableName, e); + } + } + + /** + * Lists all tables in a given database. Returns the Glue storage names (lowercase). + * + * @param glueDatabaseName The Glue storage name of the database from which to list tables. + * @return A list of table names as stored in Glue (lowercase). + * @throws CatalogException if there is an error fetching the table list. + */ + public List<String> listTables(String glueDatabaseName) { + List<String> tableNames = new ArrayList<>(); + for (Table table : getAllGlueTables(glueDatabaseName)) { + tableNames.add(table.name()); + } + return tableNames; + } + + /** + * Creates a new table in Glue. Stores the original table name in metadata for case + * preservation. + * + * @param databaseName The Glue storage name of the database where the table should be created. + * @param tableInput The input data for creating the table (should include original name in + * parameters). + * @throws CatalogException if there is an error creating the table. + */ + public void createTable(String databaseName, TableInput tableInput) { + try { + // Validate table name from the TableInput + if (tableInput.name() != null) { + validateTableName(tableInput.name()); + } + + // The table name in tableInput should already be the Glue storage name (lowercase) + // The original name should be stored in parameters by the caller + + CreateTableRequest request = + CreateTableRequest.builder() + .databaseName(databaseName) + .tableInput(tableInput) + .build(); + // The SDK throws a typed exception for any service error, so no response + // inspection is needed: reaching the next statement means the call succeeded. + glueClient.createTable(request); + // Log both original and storage names for clarity + String originalTableName = + tableInput.parameters() != null + ? tableInput.parameters().get(GlueCatalogConstants.ORIGINAL_TABLE_NAME) + : tableInput.name(); + LOG.info( + "Created table '{}' in Glue with original name '{}' preserved", + tableInput.name(), + originalTableName); + } catch (AlreadyExistsException e) { + throw new CatalogException("Table already exists: " + e.getMessage(), e); + } catch (GlueException e) { + throw new CatalogException("Error creating table: " + e.getMessage(), e); + } + } + + /** + * Updates an existing table in Glue via UpdateTable. The TableInput's name must be the Glue + * storage name (lowercase) of the existing table. + * + * @param databaseName The Glue storage name of the database containing the table. + * @param tableInput The full replacement definition of the table. + * @throws CatalogException if there is an error updating the table. + */ + public void updateTable(String databaseName, TableInput tableInput) { + try { + UpdateTableRequest request = + UpdateTableRequest.builder() + .databaseName(databaseName) + .tableInput(tableInput) + .build(); + // Rely on SDK typed exceptions rather than inspecting the HTTP response. + glueClient.updateTable(request); + LOG.info("Updated table '{}.{}' in Glue", databaseName, tableInput.name()); + } catch (EntityNotFoundException e) { + throw new CatalogException("Table does not exist: " + e.getMessage(), e); + } catch (GlueException e) { + throw new CatalogException("Error updating table: " + e.getMessage(), e); + } + } + + /** + * Retrieves the details of a specific table from Glue. + * + * @param databaseName The name of the database where the table resides. + * @param tableName The name of the table to retrieve. + * @return The Table object containing the table details. + * @throws TableNotExistException if the table does not exist in the Glue catalog. + * @throws CatalogException if there is an error fetching the table details. + */ + public Table getGlueTable(String databaseName, String tableName) throws TableNotExistException { + try { + GetTableRequest request = + GetTableRequest.builder().databaseName(databaseName).name(tableName).build(); + Table table = glueClient.getTable(request).table(); + if (table == null) { + throw new TableNotExistException( + catalogName, new ObjectPath(databaseName, tableName)); + } + return table; + } catch (EntityNotFoundException e) { + throw new TableNotExistException(catalogName, new ObjectPath(databaseName, tableName)); + } catch (GlueException e) { + throw new CatalogException("Error getting table: " + e.getMessage(), e); + } + } + + /** + * Drops a table from Glue. + * + * @param databaseName The name of the database where the table resides. + * @param tableName The name of the table to drop. + * @throws TableNotExistException if the table does not exist in the Glue catalog. + * @throws CatalogException if there is an error dropping the table. + */ + public void dropTable(String databaseName, String tableName) throws TableNotExistException { + try { + DeleteTableRequest request = + DeleteTableRequest.builder().databaseName(databaseName).name(tableName).build(); + // Rely on SDK typed exceptions rather than inspecting the HTTP response. + glueClient.deleteTable(request); + } catch (EntityNotFoundException e) { + throw new TableNotExistException(catalogName, new ObjectPath(databaseName, tableName)); + } catch (GlueException e) { + throw new CatalogException("Error dropping table: " + e.getMessage(), e); + } + } + + /** + * Converts a Flink catalog table to Glue's TableInput object. + * + * <p>Partition columns are persisted at the {@code TableInput} level (Glue models partition + * keys separately from the storage descriptor columns), and the table comment is persisted as + * the Glue table description. + * + * @param tableName The name of the table. + * @param partitionColumns The Glue columns for the table's partition keys, in the order + * declared by the Flink table (may be empty). + * @param catalogTable The Flink CatalogTable containing the table schema. + * @param storageDescriptor The Glue storage descriptor holding the data (non-partition) + * columns. + * @param properties The properties of the table. + * @return The Glue TableInput object representing the table. + */ + public TableInput buildTableInput( + String tableName, + List<Column> partitionColumns, + CatalogTable catalogTable, + StorageDescriptor storageDescriptor, + Map<String, String> properties) { + + // Validate table name + validateTableName(tableName); + + // Store lowercase name in Glue (Glue requirement) + String glueTableName = toGlueTableName(tableName); + + // Prepare table parameters with original name preservation + Map<String, String> tableParameters = new HashMap<>(); + if (properties != null) { + tableParameters.putAll(properties); + } + + // Store original table name in metadata + tableParameters.put(GlueCatalogConstants.ORIGINAL_TABLE_NAME, tableName); + + // Glue rejects column-level parameters on partition columns, so the originalName + // column parameter cannot be used for them. Strip any parameters and preserve the + // declared partition-key case in an order-preserving table-level parameter instead. + List<Column> sanitizedPartitionColumns = null; + if (partitionColumns != null && !partitionColumns.isEmpty()) { + List<String> originalPartitionKeys = new ArrayList<>(); + sanitizedPartitionColumns = new ArrayList<>(partitionColumns.size()); + boolean anyMixedCase = false; + for (Column partitionColumn : partitionColumns) { + String originalName = GlueTableUtils.getColumnName(partitionColumn); + originalPartitionKeys.add(originalName); + if (!originalName.equals(partitionColumn.name())) { + anyMixedCase = true; + } + sanitizedPartitionColumns.add( + Column.builder() + .name(partitionColumn.name()) + .type(partitionColumn.type()) + .comment(partitionColumn.comment()) + .build()); + } + if (anyMixedCase) { + tableParameters.put( + GlueCatalogConstants.ORIGINAL_PARTITION_KEYS, + String.join(",", originalPartitionKeys)); + } + } + + TableInput.Builder builder = + TableInput.builder() + .name(glueTableName) + .storageDescriptor(storageDescriptor) + .parameters(tableParameters) + .tableType(catalogTable.getTableKind().name()); + + // Persist partition keys at the TableInput level so partition metadata declared + // in DDL survives the round-trip (Glue stores them outside the storage descriptor). + if (sanitizedPartitionColumns != null && !sanitizedPartitionColumns.isEmpty()) { + builder.partitionKeys(sanitizedPartitionColumns); + } + + // Persist the table comment as the Glue description. + if (catalogTable.getComment() != null) { + builder.description(catalogTable.getComment()); + } + + return builder.build(); + } + + /** + * Converts a user-provided table name to the name used for storage in Glue. Glue requires + * lowercase names, so we store in lowercase. + * + * @param tableName The table name as specified by the user + * @return The table name to use for Glue storage (lowercase) + */ + private String toGlueTableName(String tableName) { + return tableName.toLowerCase(); + } + + /** + * Extracts the original table name from a Glue table object. Falls back to the stored name if + * no original name is found. + * + * @param table The Glue table object + * @return The original table name with case preserved + */ + public String getOriginalTableName(Table table) { + if (table.parameters() != null + && table.parameters().containsKey(GlueCatalogConstants.ORIGINAL_TABLE_NAME)) { + return table.parameters().get(GlueCatalogConstants.ORIGINAL_TABLE_NAME); + } + // Fallback to stored name for backward compatibility + return table.name(); + } + + /** + * Finds the Glue storage name for a given original table name. This method handles + * case-insensitive lookups while preserving original case. + * + * @param glueDatabaseName The Glue storage name of the database + * @param originalTableName The original table name to find + * @return The Glue storage name if found, null if not found + * @throws CatalogException if there's an error searching + */ + public String findGlueTableName(String glueDatabaseName, String originalTableName) + throws CatalogException { + try { + // First try the direct lowercase match (most common case). A single GetTable + // call both proves existence and returns the metadata needed to verify the + // stored original name. + String glueTableName = originalTableName.toLowerCase(); + try { + Table table = getGlueTable(glueDatabaseName, glueTableName); + String storedOriginalName = getOriginalTableName(table); + if (storedOriginalName.equals(originalTableName)) { + return glueTableName; + } + } catch (TableNotExistException e) { + // Fall through to the full search below. + } + + // If direct match failed, search all tables for an original-name match. GetTables + // already returns full table metadata, so no per-table GetTable calls are needed. + for (Table table : getAllGlueTables(glueDatabaseName)) { + String storedOriginalName = getOriginalTableName(table); + if (storedOriginalName.equals(originalTableName)) { + return table.name(); // Return the Glue storage name + } + } + + return null; // Table not found + } catch (Exception e) { + throw new CatalogException( + "Error searching for table: " + glueDatabaseName + "." + originalTableName, e); + } + } + + /** + * Lists all tables in a given database with their full metadata, handling pagination. + * + * @param glueDatabaseName The Glue storage name of the database. + * @return All Glue tables in the database. + * @throws CatalogException if there is an error fetching the tables. + */ + public List<Table> getAllGlueTables(String glueDatabaseName) { + try { + List<Table> tables = new ArrayList<>(); + String nextToken = null; + while (true) { + GetTablesRequest.Builder requestBuilder = + GetTablesRequest.builder().databaseName(glueDatabaseName); + if (nextToken != null) { + requestBuilder.nextToken(nextToken); + } + GetTablesResponse response = glueClient.getTables(requestBuilder.build()); + tables.addAll(response.tableList()); + nextToken = response.nextToken(); + if (nextToken == null) { + break; + } + } + return tables; + } catch (GlueException e) { + throw new CatalogException("Error listing tables: " + e.getMessage(), e); + } + } + + /** + * Lists all tables in a given database, returning original names. This is the public version + * that returns original table names with case preserved. + * + * @param glueDatabaseName The Glue storage name of the database from which to list tables. + * @return A list of original table names with case preserved. + * @throws CatalogException if there is an error fetching the table list. + */ + public List<String> listTablesWithOriginalNames(String glueDatabaseName) { + List<String> originalTableNames = new ArrayList<>(); + for (Table table : getAllGlueTables(glueDatabaseName)) { + originalTableNames.add(getOriginalTableName(table)); + } + return originalTableNames; + } + + /** + * Checks whether a table exists by original name. + * + * @param glueDatabaseName The Glue storage name of the database. + * @param originalTableName The original table name to check. + * @return true if the table exists, false otherwise. + */ + public boolean tableExistsByOriginalName(String glueDatabaseName, String originalTableName) { Review Comment: Removed (0 call sites). ########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/operator/GlueTableOperator.java: ########## @@ -0,0 +1,509 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.catalog.glue.operator; + +import org.apache.flink.table.catalog.CatalogTable; +import org.apache.flink.table.catalog.ObjectPath; +import org.apache.flink.table.catalog.exceptions.CatalogException; +import org.apache.flink.table.catalog.exceptions.TableNotExistException; +import org.apache.flink.table.catalog.glue.util.GlueCatalogConstants; +import org.apache.flink.table.catalog.glue.util.GlueTableUtils; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.AlreadyExistsException; +import software.amazon.awssdk.services.glue.model.Column; +import software.amazon.awssdk.services.glue.model.CreateTableRequest; +import software.amazon.awssdk.services.glue.model.DeleteTableRequest; +import software.amazon.awssdk.services.glue.model.EntityNotFoundException; +import software.amazon.awssdk.services.glue.model.GetTableRequest; +import software.amazon.awssdk.services.glue.model.GetTablesRequest; +import software.amazon.awssdk.services.glue.model.GetTablesResponse; +import software.amazon.awssdk.services.glue.model.GlueException; +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 software.amazon.awssdk.services.glue.model.UpdateTableRequest; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.regex.Pattern; + +/** + * Handles all table-related operations for the Glue catalog. Provides functionality for checking + * existence, listing, creating, getting, and dropping tables in AWS Glue. + */ +public class GlueTableOperator extends GlueOperator { + + /** Logger for logging table operations. */ + private static final Logger LOG = LoggerFactory.getLogger(GlueTableOperator.class); + + /** + * Pattern for validating table names. AWS Glue supports alphanumeric characters and + * underscores. We preserve original case in metadata while storing lowercase in Glue. + */ + private static final Pattern VALID_NAME_PATTERN = Pattern.compile("^[a-zA-Z0-9_]+$"); + + /** + * Constructor for GlueTableOperations. Initializes the Glue client and catalog name. + * + * @param glueClient The Glue client to interact with AWS Glue. + * @param catalogName The name of the catalog. + */ + public GlueTableOperator(GlueClient glueClient, String catalogName) { + super(glueClient, catalogName); + } + + /** + * Validates that a table name contains only allowed characters. AWS Glue supports alphanumeric + * characters and underscores. Case is preserved in metadata while Glue stores lowercase + * internally. + * + * @param tableName The table name to validate + * @throws CatalogException if the table name contains invalid characters + */ + private void validateTableName(String tableName) { + if (tableName == null || tableName.isEmpty()) { + throw new CatalogException("Table name cannot be null or empty"); + } + + if (!VALID_NAME_PATTERN.matcher(tableName).matches()) { + throw new CatalogException( + "Table name can only contain letters, numbers, and underscores. " + + "Original case is preserved in metadata while AWS Glue stores lowercase internally."); + } + } + + /** + * Checks whether a table exists in the Glue catalog by Glue storage names. + * + * @param glueDatabaseName The Glue storage name of the database where the table should exist. + * @param glueTableName The Glue storage name of the table to check. + * @return true if the table exists, false otherwise. + */ + public boolean glueTableExists(String glueDatabaseName, String glueTableName) { + try { + return glueClient + .getTable( + builder -> + builder.databaseName(glueDatabaseName) + .name(glueTableName)) + .table() + != null; + } catch (EntityNotFoundException e) { + return false; + } catch (GlueException e) { + throw new CatalogException( + "Error checking table existence: " + glueDatabaseName + "." + glueTableName, e); + } + } + + /** + * Lists all tables in a given database. Returns the Glue storage names (lowercase). + * + * @param glueDatabaseName The Glue storage name of the database from which to list tables. + * @return A list of table names as stored in Glue (lowercase). + * @throws CatalogException if there is an error fetching the table list. + */ + public List<String> listTables(String glueDatabaseName) { + List<String> tableNames = new ArrayList<>(); + for (Table table : getAllGlueTables(glueDatabaseName)) { + tableNames.add(table.name()); + } + return tableNames; + } + + /** + * Creates a new table in Glue. Stores the original table name in metadata for case + * preservation. + * + * @param databaseName The Glue storage name of the database where the table should be created. + * @param tableInput The input data for creating the table (should include original name in + * parameters). + * @throws CatalogException if there is an error creating the table. + */ + public void createTable(String databaseName, TableInput tableInput) { + try { + // Validate table name from the TableInput + if (tableInput.name() != null) { + validateTableName(tableInput.name()); + } + + // The table name in tableInput should already be the Glue storage name (lowercase) + // The original name should be stored in parameters by the caller + + CreateTableRequest request = + CreateTableRequest.builder() + .databaseName(databaseName) + .tableInput(tableInput) + .build(); + // The SDK throws a typed exception for any service error, so no response + // inspection is needed: reaching the next statement means the call succeeded. + glueClient.createTable(request); + // Log both original and storage names for clarity + String originalTableName = + tableInput.parameters() != null + ? tableInput.parameters().get(GlueCatalogConstants.ORIGINAL_TABLE_NAME) + : tableInput.name(); + LOG.info( + "Created table '{}' in Glue with original name '{}' preserved", + tableInput.name(), + originalTableName); + } catch (AlreadyExistsException e) { + throw new CatalogException("Table already exists: " + e.getMessage(), e); + } catch (GlueException e) { + throw new CatalogException("Error creating table: " + e.getMessage(), e); + } + } + + /** + * Updates an existing table in Glue via UpdateTable. The TableInput's name must be the Glue + * storage name (lowercase) of the existing table. + * + * @param databaseName The Glue storage name of the database containing the table. + * @param tableInput The full replacement definition of the table. + * @throws CatalogException if there is an error updating the table. + */ + public void updateTable(String databaseName, TableInput tableInput) { + try { + UpdateTableRequest request = + UpdateTableRequest.builder() + .databaseName(databaseName) + .tableInput(tableInput) + .build(); + // Rely on SDK typed exceptions rather than inspecting the HTTP response. + glueClient.updateTable(request); + LOG.info("Updated table '{}.{}' in Glue", databaseName, tableInput.name()); + } catch (EntityNotFoundException e) { + throw new CatalogException("Table does not exist: " + e.getMessage(), e); + } catch (GlueException e) { + throw new CatalogException("Error updating table: " + e.getMessage(), e); + } + } + + /** + * Retrieves the details of a specific table from Glue. + * + * @param databaseName The name of the database where the table resides. + * @param tableName The name of the table to retrieve. + * @return The Table object containing the table details. + * @throws TableNotExistException if the table does not exist in the Glue catalog. + * @throws CatalogException if there is an error fetching the table details. + */ + public Table getGlueTable(String databaseName, String tableName) throws TableNotExistException { + try { + GetTableRequest request = + GetTableRequest.builder().databaseName(databaseName).name(tableName).build(); + Table table = glueClient.getTable(request).table(); + if (table == null) { + throw new TableNotExistException( + catalogName, new ObjectPath(databaseName, tableName)); + } + return table; + } catch (EntityNotFoundException e) { + throw new TableNotExistException(catalogName, new ObjectPath(databaseName, tableName)); + } catch (GlueException e) { + throw new CatalogException("Error getting table: " + e.getMessage(), e); + } + } + + /** + * Drops a table from Glue. + * + * @param databaseName The name of the database where the table resides. + * @param tableName The name of the table to drop. + * @throws TableNotExistException if the table does not exist in the Glue catalog. + * @throws CatalogException if there is an error dropping the table. + */ + public void dropTable(String databaseName, String tableName) throws TableNotExistException { + try { + DeleteTableRequest request = + DeleteTableRequest.builder().databaseName(databaseName).name(tableName).build(); + // Rely on SDK typed exceptions rather than inspecting the HTTP response. + glueClient.deleteTable(request); + } catch (EntityNotFoundException e) { + throw new TableNotExistException(catalogName, new ObjectPath(databaseName, tableName)); + } catch (GlueException e) { + throw new CatalogException("Error dropping table: " + e.getMessage(), e); + } + } + + /** + * Converts a Flink catalog table to Glue's TableInput object. + * + * <p>Partition columns are persisted at the {@code TableInput} level (Glue models partition + * keys separately from the storage descriptor columns), and the table comment is persisted as + * the Glue table description. + * + * @param tableName The name of the table. + * @param partitionColumns The Glue columns for the table's partition keys, in the order + * declared by the Flink table (may be empty). + * @param catalogTable The Flink CatalogTable containing the table schema. + * @param storageDescriptor The Glue storage descriptor holding the data (non-partition) + * columns. + * @param properties The properties of the table. + * @return The Glue TableInput object representing the table. + */ + public TableInput buildTableInput( + String tableName, + List<Column> partitionColumns, + CatalogTable catalogTable, + StorageDescriptor storageDescriptor, + Map<String, String> properties) { + + // Validate table name + validateTableName(tableName); + + // Store lowercase name in Glue (Glue requirement) + String glueTableName = toGlueTableName(tableName); + + // Prepare table parameters with original name preservation + Map<String, String> tableParameters = new HashMap<>(); + if (properties != null) { + tableParameters.putAll(properties); + } + + // Store original table name in metadata + tableParameters.put(GlueCatalogConstants.ORIGINAL_TABLE_NAME, tableName); + + // Glue rejects column-level parameters on partition columns, so the originalName + // column parameter cannot be used for them. Strip any parameters and preserve the + // declared partition-key case in an order-preserving table-level parameter instead. + List<Column> sanitizedPartitionColumns = null; + if (partitionColumns != null && !partitionColumns.isEmpty()) { + List<String> originalPartitionKeys = new ArrayList<>(); + sanitizedPartitionColumns = new ArrayList<>(partitionColumns.size()); + boolean anyMixedCase = false; + for (Column partitionColumn : partitionColumns) { + String originalName = GlueTableUtils.getColumnName(partitionColumn); + originalPartitionKeys.add(originalName); + if (!originalName.equals(partitionColumn.name())) { + anyMixedCase = true; + } + sanitizedPartitionColumns.add( + Column.builder() + .name(partitionColumn.name()) + .type(partitionColumn.type()) + .comment(partitionColumn.comment()) + .build()); + } + if (anyMixedCase) { + tableParameters.put( + GlueCatalogConstants.ORIGINAL_PARTITION_KEYS, + String.join(",", originalPartitionKeys)); + } + } + + TableInput.Builder builder = + TableInput.builder() + .name(glueTableName) + .storageDescriptor(storageDescriptor) + .parameters(tableParameters) + .tableType(catalogTable.getTableKind().name()); + + // Persist partition keys at the TableInput level so partition metadata declared + // in DDL survives the round-trip (Glue stores them outside the storage descriptor). + if (sanitizedPartitionColumns != null && !sanitizedPartitionColumns.isEmpty()) { + builder.partitionKeys(sanitizedPartitionColumns); + } + + // Persist the table comment as the Glue description. + if (catalogTable.getComment() != null) { + builder.description(catalogTable.getComment()); + } + + return builder.build(); + } + + /** + * Converts a user-provided table name to the name used for storage in Glue. Glue requires + * lowercase names, so we store in lowercase. + * + * @param tableName The table name as specified by the user + * @return The table name to use for Glue storage (lowercase) + */ + private String toGlueTableName(String tableName) { + return tableName.toLowerCase(); + } + + /** + * Extracts the original table name from a Glue table object. Falls back to the stored name if + * no original name is found. + * + * @param table The Glue table object + * @return The original table name with case preserved + */ + public String getOriginalTableName(Table table) { + if (table.parameters() != null + && table.parameters().containsKey(GlueCatalogConstants.ORIGINAL_TABLE_NAME)) { + return table.parameters().get(GlueCatalogConstants.ORIGINAL_TABLE_NAME); + } + // Fallback to stored name for backward compatibility + return table.name(); + } + + /** + * Finds the Glue storage name for a given original table name. This method handles + * case-insensitive lookups while preserving original case. + * + * @param glueDatabaseName The Glue storage name of the database + * @param originalTableName The original table name to find + * @return The Glue storage name if found, null if not found + * @throws CatalogException if there's an error searching + */ + public String findGlueTableName(String glueDatabaseName, String originalTableName) + throws CatalogException { + try { + // First try the direct lowercase match (most common case). A single GetTable + // call both proves existence and returns the metadata needed to verify the + // stored original name. + String glueTableName = originalTableName.toLowerCase(); + try { + Table table = getGlueTable(glueDatabaseName, glueTableName); + String storedOriginalName = getOriginalTableName(table); + if (storedOriginalName.equals(originalTableName)) { + return glueTableName; + } + } catch (TableNotExistException e) { + // Fall through to the full search below. + } + + // If direct match failed, search all tables for an original-name match. GetTables + // already returns full table metadata, so no per-table GetTable calls are needed. + for (Table table : getAllGlueTables(glueDatabaseName)) { + String storedOriginalName = getOriginalTableName(table); + if (storedOriginalName.equals(originalTableName)) { + return table.name(); // Return the Glue storage name + } + } + + return null; // Table not found + } catch (Exception e) { + throw new CatalogException( + "Error searching for table: " + glueDatabaseName + "." + originalTableName, e); + } + } + + /** + * Lists all tables in a given database with their full metadata, handling pagination. + * + * @param glueDatabaseName The Glue storage name of the database. + * @return All Glue tables in the database. + * @throws CatalogException if there is an error fetching the tables. + */ + public List<Table> getAllGlueTables(String glueDatabaseName) { + try { + List<Table> tables = new ArrayList<>(); + String nextToken = null; + while (true) { + GetTablesRequest.Builder requestBuilder = + GetTablesRequest.builder().databaseName(glueDatabaseName); + if (nextToken != null) { + requestBuilder.nextToken(nextToken); + } + GetTablesResponse response = glueClient.getTables(requestBuilder.build()); + tables.addAll(response.tableList()); + nextToken = response.nextToken(); + if (nextToken == null) { + break; + } + } + return tables; + } catch (GlueException e) { + throw new CatalogException("Error listing tables: " + e.getMessage(), e); + } + } + + /** + * Lists all tables in a given database, returning original names. This is the public version + * that returns original table names with case preserved. + * + * @param glueDatabaseName The Glue storage name of the database from which to list tables. + * @return A list of original table names with case preserved. + * @throws CatalogException if there is an error fetching the table list. + */ + public List<String> listTablesWithOriginalNames(String glueDatabaseName) { + List<String> originalTableNames = new ArrayList<>(); + for (Table table : getAllGlueTables(glueDatabaseName)) { + originalTableNames.add(getOriginalTableName(table)); + } + return originalTableNames; + } + + /** + * Checks whether a table exists by original name. + * + * @param glueDatabaseName The Glue storage name of the database. + * @param originalTableName The original table name to check. + * @return true if the table exists, false otherwise. + */ + public boolean tableExistsByOriginalName(String glueDatabaseName, String originalTableName) { + try { + String glueTableName = findGlueTableName(glueDatabaseName, originalTableName); + return glueTableName != null; + } catch (CatalogException e) { + LOG.warn( + "Error checking table existence for: {}.{}", + glueDatabaseName, + originalTableName, + e); + return false; + } + } + + /** + * Retrieves a table by original name. + * + * @param glueDatabaseName The Glue storage name of the database. + * @param originalTableName The original table name. + * @return The Table object containing the table details. + * @throws TableNotExistException if the table does not exist. + * @throws CatalogException if there is an error fetching the table details. + */ + public Table getTableByOriginalName(String glueDatabaseName, String originalTableName) Review Comment: Removed (0 call sites). ########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/operator/GlueTableOperator.java: ########## @@ -0,0 +1,509 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.catalog.glue.operator; + +import org.apache.flink.table.catalog.CatalogTable; +import org.apache.flink.table.catalog.ObjectPath; +import org.apache.flink.table.catalog.exceptions.CatalogException; +import org.apache.flink.table.catalog.exceptions.TableNotExistException; +import org.apache.flink.table.catalog.glue.util.GlueCatalogConstants; +import org.apache.flink.table.catalog.glue.util.GlueTableUtils; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.AlreadyExistsException; +import software.amazon.awssdk.services.glue.model.Column; +import software.amazon.awssdk.services.glue.model.CreateTableRequest; +import software.amazon.awssdk.services.glue.model.DeleteTableRequest; +import software.amazon.awssdk.services.glue.model.EntityNotFoundException; +import software.amazon.awssdk.services.glue.model.GetTableRequest; +import software.amazon.awssdk.services.glue.model.GetTablesRequest; +import software.amazon.awssdk.services.glue.model.GetTablesResponse; +import software.amazon.awssdk.services.glue.model.GlueException; +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 software.amazon.awssdk.services.glue.model.UpdateTableRequest; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.regex.Pattern; + +/** + * Handles all table-related operations for the Glue catalog. Provides functionality for checking + * existence, listing, creating, getting, and dropping tables in AWS Glue. + */ +public class GlueTableOperator extends GlueOperator { + + /** Logger for logging table operations. */ + private static final Logger LOG = LoggerFactory.getLogger(GlueTableOperator.class); + + /** + * Pattern for validating table names. AWS Glue supports alphanumeric characters and + * underscores. We preserve original case in metadata while storing lowercase in Glue. + */ + private static final Pattern VALID_NAME_PATTERN = Pattern.compile("^[a-zA-Z0-9_]+$"); + + /** + * Constructor for GlueTableOperations. Initializes the Glue client and catalog name. + * + * @param glueClient The Glue client to interact with AWS Glue. + * @param catalogName The name of the catalog. + */ + public GlueTableOperator(GlueClient glueClient, String catalogName) { + super(glueClient, catalogName); + } + + /** + * Validates that a table name contains only allowed characters. AWS Glue supports alphanumeric + * characters and underscores. Case is preserved in metadata while Glue stores lowercase + * internally. + * + * @param tableName The table name to validate + * @throws CatalogException if the table name contains invalid characters + */ + private void validateTableName(String tableName) { + if (tableName == null || tableName.isEmpty()) { + throw new CatalogException("Table name cannot be null or empty"); + } + + if (!VALID_NAME_PATTERN.matcher(tableName).matches()) { + throw new CatalogException( + "Table name can only contain letters, numbers, and underscores. " + + "Original case is preserved in metadata while AWS Glue stores lowercase internally."); + } + } + + /** + * Checks whether a table exists in the Glue catalog by Glue storage names. + * + * @param glueDatabaseName The Glue storage name of the database where the table should exist. + * @param glueTableName The Glue storage name of the table to check. + * @return true if the table exists, false otherwise. + */ + public boolean glueTableExists(String glueDatabaseName, String glueTableName) { + try { + return glueClient + .getTable( + builder -> + builder.databaseName(glueDatabaseName) + .name(glueTableName)) + .table() + != null; + } catch (EntityNotFoundException e) { + return false; + } catch (GlueException e) { + throw new CatalogException( + "Error checking table existence: " + glueDatabaseName + "." + glueTableName, e); + } + } + + /** + * Lists all tables in a given database. Returns the Glue storage names (lowercase). + * + * @param glueDatabaseName The Glue storage name of the database from which to list tables. + * @return A list of table names as stored in Glue (lowercase). + * @throws CatalogException if there is an error fetching the table list. + */ + public List<String> listTables(String glueDatabaseName) { + List<String> tableNames = new ArrayList<>(); + for (Table table : getAllGlueTables(glueDatabaseName)) { + tableNames.add(table.name()); + } + return tableNames; + } + + /** + * Creates a new table in Glue. Stores the original table name in metadata for case + * preservation. + * + * @param databaseName The Glue storage name of the database where the table should be created. + * @param tableInput The input data for creating the table (should include original name in + * parameters). + * @throws CatalogException if there is an error creating the table. + */ + public void createTable(String databaseName, TableInput tableInput) { + try { + // Validate table name from the TableInput + if (tableInput.name() != null) { + validateTableName(tableInput.name()); + } + + // The table name in tableInput should already be the Glue storage name (lowercase) + // The original name should be stored in parameters by the caller + + CreateTableRequest request = + CreateTableRequest.builder() + .databaseName(databaseName) + .tableInput(tableInput) + .build(); + // The SDK throws a typed exception for any service error, so no response + // inspection is needed: reaching the next statement means the call succeeded. + glueClient.createTable(request); + // Log both original and storage names for clarity + String originalTableName = + tableInput.parameters() != null + ? tableInput.parameters().get(GlueCatalogConstants.ORIGINAL_TABLE_NAME) + : tableInput.name(); + LOG.info( + "Created table '{}' in Glue with original name '{}' preserved", + tableInput.name(), + originalTableName); + } catch (AlreadyExistsException e) { + throw new CatalogException("Table already exists: " + e.getMessage(), e); + } catch (GlueException e) { + throw new CatalogException("Error creating table: " + e.getMessage(), e); + } + } + + /** + * Updates an existing table in Glue via UpdateTable. The TableInput's name must be the Glue + * storage name (lowercase) of the existing table. + * + * @param databaseName The Glue storage name of the database containing the table. + * @param tableInput The full replacement definition of the table. + * @throws CatalogException if there is an error updating the table. + */ + public void updateTable(String databaseName, TableInput tableInput) { + try { + UpdateTableRequest request = + UpdateTableRequest.builder() + .databaseName(databaseName) + .tableInput(tableInput) + .build(); + // Rely on SDK typed exceptions rather than inspecting the HTTP response. + glueClient.updateTable(request); + LOG.info("Updated table '{}.{}' in Glue", databaseName, tableInput.name()); + } catch (EntityNotFoundException e) { + throw new CatalogException("Table does not exist: " + e.getMessage(), e); + } catch (GlueException e) { + throw new CatalogException("Error updating table: " + e.getMessage(), e); + } + } + + /** + * Retrieves the details of a specific table from Glue. + * + * @param databaseName The name of the database where the table resides. + * @param tableName The name of the table to retrieve. + * @return The Table object containing the table details. + * @throws TableNotExistException if the table does not exist in the Glue catalog. + * @throws CatalogException if there is an error fetching the table details. + */ + public Table getGlueTable(String databaseName, String tableName) throws TableNotExistException { + try { + GetTableRequest request = + GetTableRequest.builder().databaseName(databaseName).name(tableName).build(); + Table table = glueClient.getTable(request).table(); + if (table == null) { + throw new TableNotExistException( + catalogName, new ObjectPath(databaseName, tableName)); + } + return table; + } catch (EntityNotFoundException e) { + throw new TableNotExistException(catalogName, new ObjectPath(databaseName, tableName)); + } catch (GlueException e) { + throw new CatalogException("Error getting table: " + e.getMessage(), e); + } + } + + /** + * Drops a table from Glue. + * + * @param databaseName The name of the database where the table resides. + * @param tableName The name of the table to drop. + * @throws TableNotExistException if the table does not exist in the Glue catalog. + * @throws CatalogException if there is an error dropping the table. + */ + public void dropTable(String databaseName, String tableName) throws TableNotExistException { + try { + DeleteTableRequest request = + DeleteTableRequest.builder().databaseName(databaseName).name(tableName).build(); + // Rely on SDK typed exceptions rather than inspecting the HTTP response. + glueClient.deleteTable(request); + } catch (EntityNotFoundException e) { + throw new TableNotExistException(catalogName, new ObjectPath(databaseName, tableName)); + } catch (GlueException e) { + throw new CatalogException("Error dropping table: " + e.getMessage(), e); + } + } + + /** + * Converts a Flink catalog table to Glue's TableInput object. + * + * <p>Partition columns are persisted at the {@code TableInput} level (Glue models partition + * keys separately from the storage descriptor columns), and the table comment is persisted as + * the Glue table description. + * + * @param tableName The name of the table. + * @param partitionColumns The Glue columns for the table's partition keys, in the order + * declared by the Flink table (may be empty). + * @param catalogTable The Flink CatalogTable containing the table schema. + * @param storageDescriptor The Glue storage descriptor holding the data (non-partition) + * columns. + * @param properties The properties of the table. + * @return The Glue TableInput object representing the table. + */ + public TableInput buildTableInput( + String tableName, + List<Column> partitionColumns, + CatalogTable catalogTable, + StorageDescriptor storageDescriptor, + Map<String, String> properties) { + + // Validate table name + validateTableName(tableName); + + // Store lowercase name in Glue (Glue requirement) + String glueTableName = toGlueTableName(tableName); + + // Prepare table parameters with original name preservation + Map<String, String> tableParameters = new HashMap<>(); + if (properties != null) { + tableParameters.putAll(properties); + } + + // Store original table name in metadata + tableParameters.put(GlueCatalogConstants.ORIGINAL_TABLE_NAME, tableName); + + // Glue rejects column-level parameters on partition columns, so the originalName + // column parameter cannot be used for them. Strip any parameters and preserve the + // declared partition-key case in an order-preserving table-level parameter instead. + List<Column> sanitizedPartitionColumns = null; + if (partitionColumns != null && !partitionColumns.isEmpty()) { + List<String> originalPartitionKeys = new ArrayList<>(); + sanitizedPartitionColumns = new ArrayList<>(partitionColumns.size()); + boolean anyMixedCase = false; + for (Column partitionColumn : partitionColumns) { + String originalName = GlueTableUtils.getColumnName(partitionColumn); + originalPartitionKeys.add(originalName); + if (!originalName.equals(partitionColumn.name())) { + anyMixedCase = true; + } + sanitizedPartitionColumns.add( + Column.builder() + .name(partitionColumn.name()) + .type(partitionColumn.type()) + .comment(partitionColumn.comment()) + .build()); + } + if (anyMixedCase) { + tableParameters.put( + GlueCatalogConstants.ORIGINAL_PARTITION_KEYS, + String.join(",", originalPartitionKeys)); + } + } + + TableInput.Builder builder = + TableInput.builder() + .name(glueTableName) + .storageDescriptor(storageDescriptor) + .parameters(tableParameters) + .tableType(catalogTable.getTableKind().name()); + + // Persist partition keys at the TableInput level so partition metadata declared + // in DDL survives the round-trip (Glue stores them outside the storage descriptor). + if (sanitizedPartitionColumns != null && !sanitizedPartitionColumns.isEmpty()) { + builder.partitionKeys(sanitizedPartitionColumns); + } + + // Persist the table comment as the Glue description. + if (catalogTable.getComment() != null) { + builder.description(catalogTable.getComment()); + } + + return builder.build(); + } + + /** + * Converts a user-provided table name to the name used for storage in Glue. Glue requires + * lowercase names, so we store in lowercase. + * + * @param tableName The table name as specified by the user + * @return The table name to use for Glue storage (lowercase) + */ + private String toGlueTableName(String tableName) { + return tableName.toLowerCase(); + } + + /** + * Extracts the original table name from a Glue table object. Falls back to the stored name if + * no original name is found. + * + * @param table The Glue table object + * @return The original table name with case preserved + */ + public String getOriginalTableName(Table table) { + if (table.parameters() != null + && table.parameters().containsKey(GlueCatalogConstants.ORIGINAL_TABLE_NAME)) { + return table.parameters().get(GlueCatalogConstants.ORIGINAL_TABLE_NAME); + } + // Fallback to stored name for backward compatibility + return table.name(); + } + + /** + * Finds the Glue storage name for a given original table name. This method handles + * case-insensitive lookups while preserving original case. + * + * @param glueDatabaseName The Glue storage name of the database + * @param originalTableName The original table name to find + * @return The Glue storage name if found, null if not found + * @throws CatalogException if there's an error searching + */ + public String findGlueTableName(String glueDatabaseName, String originalTableName) + throws CatalogException { + try { + // First try the direct lowercase match (most common case). A single GetTable + // call both proves existence and returns the metadata needed to verify the + // stored original name. + String glueTableName = originalTableName.toLowerCase(); + try { + Table table = getGlueTable(glueDatabaseName, glueTableName); + String storedOriginalName = getOriginalTableName(table); + if (storedOriginalName.equals(originalTableName)) { + return glueTableName; + } + } catch (TableNotExistException e) { + // Fall through to the full search below. + } + + // If direct match failed, search all tables for an original-name match. GetTables + // already returns full table metadata, so no per-table GetTable calls are needed. + for (Table table : getAllGlueTables(glueDatabaseName)) { + String storedOriginalName = getOriginalTableName(table); + if (storedOriginalName.equals(originalTableName)) { + return table.name(); // Return the Glue storage name + } + } + + return null; // Table not found + } catch (Exception e) { + throw new CatalogException( + "Error searching for table: " + glueDatabaseName + "." + originalTableName, e); + } + } + + /** + * Lists all tables in a given database with their full metadata, handling pagination. + * + * @param glueDatabaseName The Glue storage name of the database. + * @return All Glue tables in the database. + * @throws CatalogException if there is an error fetching the tables. + */ + public List<Table> getAllGlueTables(String glueDatabaseName) { + try { + List<Table> tables = new ArrayList<>(); + String nextToken = null; + while (true) { + GetTablesRequest.Builder requestBuilder = + GetTablesRequest.builder().databaseName(glueDatabaseName); + if (nextToken != null) { + requestBuilder.nextToken(nextToken); + } + GetTablesResponse response = glueClient.getTables(requestBuilder.build()); + tables.addAll(response.tableList()); + nextToken = response.nextToken(); + if (nextToken == null) { + break; + } + } + return tables; + } catch (GlueException e) { + throw new CatalogException("Error listing tables: " + e.getMessage(), e); + } + } + + /** + * Lists all tables in a given database, returning original names. This is the public version + * that returns original table names with case preserved. + * + * @param glueDatabaseName The Glue storage name of the database from which to list tables. + * @return A list of original table names with case preserved. + * @throws CatalogException if there is an error fetching the table list. + */ + public List<String> listTablesWithOriginalNames(String glueDatabaseName) { + List<String> originalTableNames = new ArrayList<>(); + for (Table table : getAllGlueTables(glueDatabaseName)) { + originalTableNames.add(getOriginalTableName(table)); + } + return originalTableNames; + } + + /** + * Checks whether a table exists by original name. + * + * @param glueDatabaseName The Glue storage name of the database. + * @param originalTableName The original table name to check. + * @return true if the table exists, false otherwise. + */ + public boolean tableExistsByOriginalName(String glueDatabaseName, String originalTableName) { + try { + String glueTableName = findGlueTableName(glueDatabaseName, originalTableName); + return glueTableName != null; + } catch (CatalogException e) { + LOG.warn( + "Error checking table existence for: {}.{}", + glueDatabaseName, + originalTableName, + e); + return false; + } + } + + /** + * Retrieves a table by original name. + * + * @param glueDatabaseName The Glue storage name of the database. + * @param originalTableName The original table name. + * @return The Table object containing the table details. + * @throws TableNotExistException if the table does not exist. + * @throws CatalogException if there is an error fetching the table details. + */ + public Table getTableByOriginalName(String glueDatabaseName, String originalTableName) + throws TableNotExistException { + String glueTableName = findGlueTableName(glueDatabaseName, originalTableName); + if (glueTableName == null) { + throw new TableNotExistException( + catalogName, new ObjectPath(glueDatabaseName, originalTableName)); + } + return getGlueTable(glueDatabaseName, glueTableName); + } + + /** + * Drops a table by original name. + * + * @param glueDatabaseName The Glue storage name of the database. + * @param originalTableName The original table name. + * @throws TableNotExistException if the table does not exist. + * @throws CatalogException if there is an error dropping the table. + */ + public void dropTableByOriginalName(String glueDatabaseName, String originalTableName) Review Comment: Removed (0 call sites). ########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/operator/GlueTableOperator.java: ########## @@ -0,0 +1,509 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.catalog.glue.operator; + +import org.apache.flink.table.catalog.CatalogTable; +import org.apache.flink.table.catalog.ObjectPath; +import org.apache.flink.table.catalog.exceptions.CatalogException; +import org.apache.flink.table.catalog.exceptions.TableNotExistException; +import org.apache.flink.table.catalog.glue.util.GlueCatalogConstants; +import org.apache.flink.table.catalog.glue.util.GlueTableUtils; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.AlreadyExistsException; +import software.amazon.awssdk.services.glue.model.Column; +import software.amazon.awssdk.services.glue.model.CreateTableRequest; +import software.amazon.awssdk.services.glue.model.DeleteTableRequest; +import software.amazon.awssdk.services.glue.model.EntityNotFoundException; +import software.amazon.awssdk.services.glue.model.GetTableRequest; +import software.amazon.awssdk.services.glue.model.GetTablesRequest; +import software.amazon.awssdk.services.glue.model.GetTablesResponse; +import software.amazon.awssdk.services.glue.model.GlueException; +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 software.amazon.awssdk.services.glue.model.UpdateTableRequest; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.regex.Pattern; + +/** + * Handles all table-related operations for the Glue catalog. Provides functionality for checking + * existence, listing, creating, getting, and dropping tables in AWS Glue. + */ +public class GlueTableOperator extends GlueOperator { + + /** Logger for logging table operations. */ + private static final Logger LOG = LoggerFactory.getLogger(GlueTableOperator.class); + + /** + * Pattern for validating table names. AWS Glue supports alphanumeric characters and + * underscores. We preserve original case in metadata while storing lowercase in Glue. + */ + private static final Pattern VALID_NAME_PATTERN = Pattern.compile("^[a-zA-Z0-9_]+$"); + + /** + * Constructor for GlueTableOperations. Initializes the Glue client and catalog name. + * + * @param glueClient The Glue client to interact with AWS Glue. + * @param catalogName The name of the catalog. + */ + public GlueTableOperator(GlueClient glueClient, String catalogName) { + super(glueClient, catalogName); + } + + /** + * Validates that a table name contains only allowed characters. AWS Glue supports alphanumeric + * characters and underscores. Case is preserved in metadata while Glue stores lowercase + * internally. + * + * @param tableName The table name to validate + * @throws CatalogException if the table name contains invalid characters + */ + private void validateTableName(String tableName) { + if (tableName == null || tableName.isEmpty()) { + throw new CatalogException("Table name cannot be null or empty"); + } + + if (!VALID_NAME_PATTERN.matcher(tableName).matches()) { + throw new CatalogException( + "Table name can only contain letters, numbers, and underscores. " + + "Original case is preserved in metadata while AWS Glue stores lowercase internally."); + } + } + + /** + * Checks whether a table exists in the Glue catalog by Glue storage names. + * + * @param glueDatabaseName The Glue storage name of the database where the table should exist. + * @param glueTableName The Glue storage name of the table to check. + * @return true if the table exists, false otherwise. + */ + public boolean glueTableExists(String glueDatabaseName, String glueTableName) { + try { + return glueClient + .getTable( + builder -> + builder.databaseName(glueDatabaseName) + .name(glueTableName)) + .table() + != null; + } catch (EntityNotFoundException e) { + return false; + } catch (GlueException e) { + throw new CatalogException( + "Error checking table existence: " + glueDatabaseName + "." + glueTableName, e); + } + } + + /** + * Lists all tables in a given database. Returns the Glue storage names (lowercase). + * + * @param glueDatabaseName The Glue storage name of the database from which to list tables. + * @return A list of table names as stored in Glue (lowercase). + * @throws CatalogException if there is an error fetching the table list. + */ + public List<String> listTables(String glueDatabaseName) { + List<String> tableNames = new ArrayList<>(); + for (Table table : getAllGlueTables(glueDatabaseName)) { + tableNames.add(table.name()); + } + return tableNames; + } + + /** + * Creates a new table in Glue. Stores the original table name in metadata for case + * preservation. + * + * @param databaseName The Glue storage name of the database where the table should be created. + * @param tableInput The input data for creating the table (should include original name in + * parameters). + * @throws CatalogException if there is an error creating the table. + */ + public void createTable(String databaseName, TableInput tableInput) { + try { + // Validate table name from the TableInput + if (tableInput.name() != null) { + validateTableName(tableInput.name()); + } + + // The table name in tableInput should already be the Glue storage name (lowercase) + // The original name should be stored in parameters by the caller + + CreateTableRequest request = + CreateTableRequest.builder() + .databaseName(databaseName) + .tableInput(tableInput) + .build(); + // The SDK throws a typed exception for any service error, so no response + // inspection is needed: reaching the next statement means the call succeeded. + glueClient.createTable(request); + // Log both original and storage names for clarity + String originalTableName = + tableInput.parameters() != null + ? tableInput.parameters().get(GlueCatalogConstants.ORIGINAL_TABLE_NAME) + : tableInput.name(); + LOG.info( + "Created table '{}' in Glue with original name '{}' preserved", + tableInput.name(), + originalTableName); + } catch (AlreadyExistsException e) { + throw new CatalogException("Table already exists: " + e.getMessage(), e); + } catch (GlueException e) { + throw new CatalogException("Error creating table: " + e.getMessage(), e); + } + } + + /** + * Updates an existing table in Glue via UpdateTable. The TableInput's name must be the Glue + * storage name (lowercase) of the existing table. + * + * @param databaseName The Glue storage name of the database containing the table. + * @param tableInput The full replacement definition of the table. + * @throws CatalogException if there is an error updating the table. + */ + public void updateTable(String databaseName, TableInput tableInput) { + try { + UpdateTableRequest request = + UpdateTableRequest.builder() + .databaseName(databaseName) + .tableInput(tableInput) + .build(); + // Rely on SDK typed exceptions rather than inspecting the HTTP response. + glueClient.updateTable(request); + LOG.info("Updated table '{}.{}' in Glue", databaseName, tableInput.name()); + } catch (EntityNotFoundException e) { + throw new CatalogException("Table does not exist: " + e.getMessage(), e); + } catch (GlueException e) { + throw new CatalogException("Error updating table: " + e.getMessage(), e); + } + } + + /** + * Retrieves the details of a specific table from Glue. + * + * @param databaseName The name of the database where the table resides. + * @param tableName The name of the table to retrieve. + * @return The Table object containing the table details. + * @throws TableNotExistException if the table does not exist in the Glue catalog. + * @throws CatalogException if there is an error fetching the table details. + */ + public Table getGlueTable(String databaseName, String tableName) throws TableNotExistException { + try { + GetTableRequest request = + GetTableRequest.builder().databaseName(databaseName).name(tableName).build(); + Table table = glueClient.getTable(request).table(); + if (table == null) { + throw new TableNotExistException( + catalogName, new ObjectPath(databaseName, tableName)); + } + return table; + } catch (EntityNotFoundException e) { + throw new TableNotExistException(catalogName, new ObjectPath(databaseName, tableName)); + } catch (GlueException e) { + throw new CatalogException("Error getting table: " + e.getMessage(), e); + } + } + + /** + * Drops a table from Glue. + * + * @param databaseName The name of the database where the table resides. + * @param tableName The name of the table to drop. + * @throws TableNotExistException if the table does not exist in the Glue catalog. + * @throws CatalogException if there is an error dropping the table. + */ + public void dropTable(String databaseName, String tableName) throws TableNotExistException { + try { + DeleteTableRequest request = + DeleteTableRequest.builder().databaseName(databaseName).name(tableName).build(); + // Rely on SDK typed exceptions rather than inspecting the HTTP response. + glueClient.deleteTable(request); + } catch (EntityNotFoundException e) { + throw new TableNotExistException(catalogName, new ObjectPath(databaseName, tableName)); + } catch (GlueException e) { + throw new CatalogException("Error dropping table: " + e.getMessage(), e); + } + } + + /** + * Converts a Flink catalog table to Glue's TableInput object. + * + * <p>Partition columns are persisted at the {@code TableInput} level (Glue models partition + * keys separately from the storage descriptor columns), and the table comment is persisted as + * the Glue table description. + * + * @param tableName The name of the table. + * @param partitionColumns The Glue columns for the table's partition keys, in the order + * declared by the Flink table (may be empty). + * @param catalogTable The Flink CatalogTable containing the table schema. + * @param storageDescriptor The Glue storage descriptor holding the data (non-partition) + * columns. + * @param properties The properties of the table. + * @return The Glue TableInput object representing the table. + */ + public TableInput buildTableInput( + String tableName, + List<Column> partitionColumns, + CatalogTable catalogTable, + StorageDescriptor storageDescriptor, + Map<String, String> properties) { + + // Validate table name + validateTableName(tableName); + + // Store lowercase name in Glue (Glue requirement) + String glueTableName = toGlueTableName(tableName); + + // Prepare table parameters with original name preservation + Map<String, String> tableParameters = new HashMap<>(); + if (properties != null) { + tableParameters.putAll(properties); + } + + // Store original table name in metadata + tableParameters.put(GlueCatalogConstants.ORIGINAL_TABLE_NAME, tableName); + + // Glue rejects column-level parameters on partition columns, so the originalName + // column parameter cannot be used for them. Strip any parameters and preserve the + // declared partition-key case in an order-preserving table-level parameter instead. + List<Column> sanitizedPartitionColumns = null; + if (partitionColumns != null && !partitionColumns.isEmpty()) { + List<String> originalPartitionKeys = new ArrayList<>(); + sanitizedPartitionColumns = new ArrayList<>(partitionColumns.size()); + boolean anyMixedCase = false; + for (Column partitionColumn : partitionColumns) { + String originalName = GlueTableUtils.getColumnName(partitionColumn); + originalPartitionKeys.add(originalName); + if (!originalName.equals(partitionColumn.name())) { + anyMixedCase = true; + } + sanitizedPartitionColumns.add( + Column.builder() + .name(partitionColumn.name()) + .type(partitionColumn.type()) + .comment(partitionColumn.comment()) + .build()); + } + if (anyMixedCase) { + tableParameters.put( + GlueCatalogConstants.ORIGINAL_PARTITION_KEYS, + String.join(",", originalPartitionKeys)); + } + } + + TableInput.Builder builder = + TableInput.builder() + .name(glueTableName) + .storageDescriptor(storageDescriptor) + .parameters(tableParameters) + .tableType(catalogTable.getTableKind().name()); + + // Persist partition keys at the TableInput level so partition metadata declared + // in DDL survives the round-trip (Glue stores them outside the storage descriptor). + if (sanitizedPartitionColumns != null && !sanitizedPartitionColumns.isEmpty()) { + builder.partitionKeys(sanitizedPartitionColumns); + } + + // Persist the table comment as the Glue description. + if (catalogTable.getComment() != null) { + builder.description(catalogTable.getComment()); + } + + return builder.build(); + } + + /** + * Converts a user-provided table name to the name used for storage in Glue. Glue requires + * lowercase names, so we store in lowercase. + * + * @param tableName The table name as specified by the user + * @return The table name to use for Glue storage (lowercase) + */ + private String toGlueTableName(String tableName) { + return tableName.toLowerCase(); + } + + /** + * Extracts the original table name from a Glue table object. Falls back to the stored name if + * no original name is found. + * + * @param table The Glue table object + * @return The original table name with case preserved + */ + public String getOriginalTableName(Table table) { + if (table.parameters() != null + && table.parameters().containsKey(GlueCatalogConstants.ORIGINAL_TABLE_NAME)) { + return table.parameters().get(GlueCatalogConstants.ORIGINAL_TABLE_NAME); + } + // Fallback to stored name for backward compatibility + return table.name(); + } + + /** + * Finds the Glue storage name for a given original table name. This method handles + * case-insensitive lookups while preserving original case. + * + * @param glueDatabaseName The Glue storage name of the database + * @param originalTableName The original table name to find + * @return The Glue storage name if found, null if not found + * @throws CatalogException if there's an error searching + */ + public String findGlueTableName(String glueDatabaseName, String originalTableName) + throws CatalogException { + try { + // First try the direct lowercase match (most common case). A single GetTable + // call both proves existence and returns the metadata needed to verify the + // stored original name. + String glueTableName = originalTableName.toLowerCase(); + try { + Table table = getGlueTable(glueDatabaseName, glueTableName); + String storedOriginalName = getOriginalTableName(table); + if (storedOriginalName.equals(originalTableName)) { + return glueTableName; + } + } catch (TableNotExistException e) { + // Fall through to the full search below. + } + + // If direct match failed, search all tables for an original-name match. GetTables + // already returns full table metadata, so no per-table GetTable calls are needed. + for (Table table : getAllGlueTables(glueDatabaseName)) { + String storedOriginalName = getOriginalTableName(table); + if (storedOriginalName.equals(originalTableName)) { + return table.name(); // Return the Glue storage name + } + } Review Comment: No — for the same reason as on the database side: storage name is always `lowercase(declared)`, so if the lowercase `GetTable` misses, nothing else can hit. The fallback (and the "verify original name" step) are gone; table resolution is a single `GetTable` (`getGlueTableOrNull`). ########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/operator/GlueTableOperator.java: ########## @@ -0,0 +1,509 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.catalog.glue.operator; + +import org.apache.flink.table.catalog.CatalogTable; +import org.apache.flink.table.catalog.ObjectPath; +import org.apache.flink.table.catalog.exceptions.CatalogException; +import org.apache.flink.table.catalog.exceptions.TableNotExistException; +import org.apache.flink.table.catalog.glue.util.GlueCatalogConstants; +import org.apache.flink.table.catalog.glue.util.GlueTableUtils; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.AlreadyExistsException; +import software.amazon.awssdk.services.glue.model.Column; +import software.amazon.awssdk.services.glue.model.CreateTableRequest; +import software.amazon.awssdk.services.glue.model.DeleteTableRequest; +import software.amazon.awssdk.services.glue.model.EntityNotFoundException; +import software.amazon.awssdk.services.glue.model.GetTableRequest; +import software.amazon.awssdk.services.glue.model.GetTablesRequest; +import software.amazon.awssdk.services.glue.model.GetTablesResponse; +import software.amazon.awssdk.services.glue.model.GlueException; +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 software.amazon.awssdk.services.glue.model.UpdateTableRequest; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.regex.Pattern; + +/** + * Handles all table-related operations for the Glue catalog. Provides functionality for checking + * existence, listing, creating, getting, and dropping tables in AWS Glue. + */ +public class GlueTableOperator extends GlueOperator { + + /** Logger for logging table operations. */ + private static final Logger LOG = LoggerFactory.getLogger(GlueTableOperator.class); + + /** + * Pattern for validating table names. AWS Glue supports alphanumeric characters and + * underscores. We preserve original case in metadata while storing lowercase in Glue. + */ + private static final Pattern VALID_NAME_PATTERN = Pattern.compile("^[a-zA-Z0-9_]+$"); + + /** + * Constructor for GlueTableOperations. Initializes the Glue client and catalog name. + * + * @param glueClient The Glue client to interact with AWS Glue. + * @param catalogName The name of the catalog. + */ + public GlueTableOperator(GlueClient glueClient, String catalogName) { + super(glueClient, catalogName); + } + + /** + * Validates that a table name contains only allowed characters. AWS Glue supports alphanumeric + * characters and underscores. Case is preserved in metadata while Glue stores lowercase + * internally. + * + * @param tableName The table name to validate + * @throws CatalogException if the table name contains invalid characters + */ + private void validateTableName(String tableName) { + if (tableName == null || tableName.isEmpty()) { + throw new CatalogException("Table name cannot be null or empty"); + } + + if (!VALID_NAME_PATTERN.matcher(tableName).matches()) { + throw new CatalogException( + "Table name can only contain letters, numbers, and underscores. " + + "Original case is preserved in metadata while AWS Glue stores lowercase internally."); + } + } + + /** + * Checks whether a table exists in the Glue catalog by Glue storage names. + * + * @param glueDatabaseName The Glue storage name of the database where the table should exist. + * @param glueTableName The Glue storage name of the table to check. + * @return true if the table exists, false otherwise. + */ + public boolean glueTableExists(String glueDatabaseName, String glueTableName) { + try { + return glueClient + .getTable( + builder -> + builder.databaseName(glueDatabaseName) + .name(glueTableName)) + .table() + != null; + } catch (EntityNotFoundException e) { + return false; + } catch (GlueException e) { + throw new CatalogException( + "Error checking table existence: " + glueDatabaseName + "." + glueTableName, e); + } + } + + /** + * Lists all tables in a given database. Returns the Glue storage names (lowercase). + * + * @param glueDatabaseName The Glue storage name of the database from which to list tables. + * @return A list of table names as stored in Glue (lowercase). + * @throws CatalogException if there is an error fetching the table list. + */ + public List<String> listTables(String glueDatabaseName) { + List<String> tableNames = new ArrayList<>(); + for (Table table : getAllGlueTables(glueDatabaseName)) { + tableNames.add(table.name()); + } + return tableNames; + } + + /** + * Creates a new table in Glue. Stores the original table name in metadata for case + * preservation. + * + * @param databaseName The Glue storage name of the database where the table should be created. + * @param tableInput The input data for creating the table (should include original name in + * parameters). + * @throws CatalogException if there is an error creating the table. + */ + public void createTable(String databaseName, TableInput tableInput) { + try { + // Validate table name from the TableInput + if (tableInput.name() != null) { + validateTableName(tableInput.name()); + } + + // The table name in tableInput should already be the Glue storage name (lowercase) + // The original name should be stored in parameters by the caller + + CreateTableRequest request = + CreateTableRequest.builder() + .databaseName(databaseName) + .tableInput(tableInput) + .build(); + // The SDK throws a typed exception for any service error, so no response + // inspection is needed: reaching the next statement means the call succeeded. + glueClient.createTable(request); + // Log both original and storage names for clarity + String originalTableName = + tableInput.parameters() != null + ? tableInput.parameters().get(GlueCatalogConstants.ORIGINAL_TABLE_NAME) + : tableInput.name(); + LOG.info( + "Created table '{}' in Glue with original name '{}' preserved", + tableInput.name(), + originalTableName); + } catch (AlreadyExistsException e) { + throw new CatalogException("Table already exists: " + e.getMessage(), e); + } catch (GlueException e) { + throw new CatalogException("Error creating table: " + e.getMessage(), e); + } + } + + /** + * Updates an existing table in Glue via UpdateTable. The TableInput's name must be the Glue + * storage name (lowercase) of the existing table. + * + * @param databaseName The Glue storage name of the database containing the table. + * @param tableInput The full replacement definition of the table. + * @throws CatalogException if there is an error updating the table. + */ + public void updateTable(String databaseName, TableInput tableInput) { + try { + UpdateTableRequest request = + UpdateTableRequest.builder() + .databaseName(databaseName) + .tableInput(tableInput) + .build(); + // Rely on SDK typed exceptions rather than inspecting the HTTP response. + glueClient.updateTable(request); + LOG.info("Updated table '{}.{}' in Glue", databaseName, tableInput.name()); + } catch (EntityNotFoundException e) { + throw new CatalogException("Table does not exist: " + e.getMessage(), e); + } catch (GlueException e) { + throw new CatalogException("Error updating table: " + e.getMessage(), e); + } + } + + /** + * Retrieves the details of a specific table from Glue. + * + * @param databaseName The name of the database where the table resides. + * @param tableName The name of the table to retrieve. + * @return The Table object containing the table details. + * @throws TableNotExistException if the table does not exist in the Glue catalog. + * @throws CatalogException if there is an error fetching the table details. + */ + public Table getGlueTable(String databaseName, String tableName) throws TableNotExistException { + try { + GetTableRequest request = + GetTableRequest.builder().databaseName(databaseName).name(tableName).build(); + Table table = glueClient.getTable(request).table(); + if (table == null) { + throw new TableNotExistException( + catalogName, new ObjectPath(databaseName, tableName)); + } + return table; + } catch (EntityNotFoundException e) { + throw new TableNotExistException(catalogName, new ObjectPath(databaseName, tableName)); + } catch (GlueException e) { + throw new CatalogException("Error getting table: " + e.getMessage(), e); + } + } + + /** + * Drops a table from Glue. + * + * @param databaseName The name of the database where the table resides. + * @param tableName The name of the table to drop. + * @throws TableNotExistException if the table does not exist in the Glue catalog. + * @throws CatalogException if there is an error dropping the table. + */ + public void dropTable(String databaseName, String tableName) throws TableNotExistException { + try { + DeleteTableRequest request = + DeleteTableRequest.builder().databaseName(databaseName).name(tableName).build(); + // Rely on SDK typed exceptions rather than inspecting the HTTP response. + glueClient.deleteTable(request); + } catch (EntityNotFoundException e) { + throw new TableNotExistException(catalogName, new ObjectPath(databaseName, tableName)); + } catch (GlueException e) { + throw new CatalogException("Error dropping table: " + e.getMessage(), e); + } + } + + /** + * Converts a Flink catalog table to Glue's TableInput object. + * + * <p>Partition columns are persisted at the {@code TableInput} level (Glue models partition + * keys separately from the storage descriptor columns), and the table comment is persisted as + * the Glue table description. + * + * @param tableName The name of the table. + * @param partitionColumns The Glue columns for the table's partition keys, in the order + * declared by the Flink table (may be empty). + * @param catalogTable The Flink CatalogTable containing the table schema. + * @param storageDescriptor The Glue storage descriptor holding the data (non-partition) + * columns. + * @param properties The properties of the table. + * @return The Glue TableInput object representing the table. + */ + public TableInput buildTableInput( + String tableName, + List<Column> partitionColumns, + CatalogTable catalogTable, + StorageDescriptor storageDescriptor, + Map<String, String> properties) { + + // Validate table name + validateTableName(tableName); + + // Store lowercase name in Glue (Glue requirement) + String glueTableName = toGlueTableName(tableName); + + // Prepare table parameters with original name preservation + Map<String, String> tableParameters = new HashMap<>(); + if (properties != null) { + tableParameters.putAll(properties); + } + + // Store original table name in metadata + tableParameters.put(GlueCatalogConstants.ORIGINAL_TABLE_NAME, tableName); + + // Glue rejects column-level parameters on partition columns, so the originalName + // column parameter cannot be used for them. Strip any parameters and preserve the + // declared partition-key case in an order-preserving table-level parameter instead. + List<Column> sanitizedPartitionColumns = null; + if (partitionColumns != null && !partitionColumns.isEmpty()) { + List<String> originalPartitionKeys = new ArrayList<>(); + sanitizedPartitionColumns = new ArrayList<>(partitionColumns.size()); + boolean anyMixedCase = false; + for (Column partitionColumn : partitionColumns) { + String originalName = GlueTableUtils.getColumnName(partitionColumn); + originalPartitionKeys.add(originalName); + if (!originalName.equals(partitionColumn.name())) { + anyMixedCase = true; + } + sanitizedPartitionColumns.add( + Column.builder() + .name(partitionColumn.name()) + .type(partitionColumn.type()) + .comment(partitionColumn.comment()) + .build()); + } + if (anyMixedCase) { + tableParameters.put( + GlueCatalogConstants.ORIGINAL_PARTITION_KEYS, + String.join(",", originalPartitionKeys)); + } + } + + TableInput.Builder builder = + TableInput.builder() + .name(glueTableName) + .storageDescriptor(storageDescriptor) + .parameters(tableParameters) + .tableType(catalogTable.getTableKind().name()); + + // Persist partition keys at the TableInput level so partition metadata declared + // in DDL survives the round-trip (Glue stores them outside the storage descriptor). + if (sanitizedPartitionColumns != null && !sanitizedPartitionColumns.isEmpty()) { + builder.partitionKeys(sanitizedPartitionColumns); + } + + // Persist the table comment as the Glue description. + if (catalogTable.getComment() != null) { + builder.description(catalogTable.getComment()); + } + + return builder.build(); + } + + /** + * Converts a user-provided table name to the name used for storage in Glue. Glue requires + * lowercase names, so we store in lowercase. + * + * @param tableName The table name as specified by the user + * @return The table name to use for Glue storage (lowercase) + */ + private String toGlueTableName(String tableName) { + return tableName.toLowerCase(); + } + + /** + * Extracts the original table name from a Glue table object. Falls back to the stored name if + * no original name is found. + * + * @param table The Glue table object + * @return The original table name with case preserved + */ + public String getOriginalTableName(Table table) { + if (table.parameters() != null + && table.parameters().containsKey(GlueCatalogConstants.ORIGINAL_TABLE_NAME)) { + return table.parameters().get(GlueCatalogConstants.ORIGINAL_TABLE_NAME); + } + // Fallback to stored name for backward compatibility + return table.name(); + } + + /** + * Finds the Glue storage name for a given original table name. This method handles + * case-insensitive lookups while preserving original case. + * + * @param glueDatabaseName The Glue storage name of the database + * @param originalTableName The original table name to find + * @return The Glue storage name if found, null if not found + * @throws CatalogException if there's an error searching + */ + public String findGlueTableName(String glueDatabaseName, String originalTableName) + throws CatalogException { + try { + // First try the direct lowercase match (most common case). A single GetTable + // call both proves existence and returns the metadata needed to verify the + // stored original name. + String glueTableName = originalTableName.toLowerCase(); + try { + Table table = getGlueTable(glueDatabaseName, glueTableName); + String storedOriginalName = getOriginalTableName(table); + if (storedOriginalName.equals(originalTableName)) { Review Comment: It no longer matters: the comparison itself is gone. Resolution is by the lowercase storage name only, and the declared (case-preserved) name is only read back for display via `getOriginalTableName`. -- 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]
