fmorillo7694 commented on code in PR #206: URL: https://github.com/apache/flink-connector-aws/pull/206#discussion_r4130487463
########## flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/GlueCatalog.java: ########## @@ -0,0 +1,1303 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.catalog.glue; + +import org.apache.flink.annotation.PublicEvolving; +import org.apache.flink.annotation.VisibleForTesting; +import org.apache.flink.connector.aws.config.AWSConfigConstants; +import org.apache.flink.connector.aws.util.AWSClientUtil; +import org.apache.flink.connector.aws.util.AWSGeneralUtil; +import org.apache.flink.table.api.Schema; +import org.apache.flink.table.catalog.AbstractCatalog; +import org.apache.flink.table.catalog.CatalogBaseTable; +import org.apache.flink.table.catalog.CatalogDatabase; +import org.apache.flink.table.catalog.CatalogFunction; +import org.apache.flink.table.catalog.CatalogPartition; +import org.apache.flink.table.catalog.CatalogPartitionImpl; +import org.apache.flink.table.catalog.CatalogPartitionSpec; +import org.apache.flink.table.catalog.CatalogTable; +import org.apache.flink.table.catalog.CatalogView; +import org.apache.flink.table.catalog.Column; +import org.apache.flink.table.catalog.FunctionLanguage; +import org.apache.flink.table.catalog.ObjectPath; +import org.apache.flink.table.catalog.ResolvedCatalogBaseTable; +import org.apache.flink.table.catalog.ResolvedSchema; +import org.apache.flink.table.catalog.exceptions.CatalogException; +import org.apache.flink.table.catalog.exceptions.DatabaseAlreadyExistException; +import org.apache.flink.table.catalog.exceptions.DatabaseNotEmptyException; +import org.apache.flink.table.catalog.exceptions.DatabaseNotExistException; +import org.apache.flink.table.catalog.exceptions.FunctionAlreadyExistException; +import org.apache.flink.table.catalog.exceptions.FunctionNotExistException; +import org.apache.flink.table.catalog.exceptions.PartitionAlreadyExistsException; +import org.apache.flink.table.catalog.exceptions.PartitionNotExistException; +import org.apache.flink.table.catalog.exceptions.PartitionSpecInvalidException; +import org.apache.flink.table.catalog.exceptions.TableAlreadyExistException; +import org.apache.flink.table.catalog.exceptions.TableNotExistException; +import org.apache.flink.table.catalog.exceptions.TableNotPartitionedException; +import org.apache.flink.table.catalog.exceptions.TablePartitionedException; +import org.apache.flink.table.catalog.glue.exception.UnsupportedDataTypeMappingException; +import org.apache.flink.table.catalog.glue.operator.GlueDatabaseOperator; +import org.apache.flink.table.catalog.glue.operator.GlueFunctionOperator; +import org.apache.flink.table.catalog.glue.operator.GluePartitionOperator; +import org.apache.flink.table.catalog.glue.operator.GlueTableOperator; +import org.apache.flink.table.catalog.glue.util.GlueCatalogConstants; +import org.apache.flink.table.catalog.glue.util.GlueFlinkSchemaProperties; +import org.apache.flink.table.catalog.glue.util.GlueFunctionsUtil; +import org.apache.flink.table.catalog.glue.util.GlueTableUtils; +import org.apache.flink.table.catalog.glue.util.GlueTypeConverter; +import org.apache.flink.table.catalog.stats.CatalogColumnStatistics; +import org.apache.flink.table.catalog.stats.CatalogTableStatistics; +import org.apache.flink.table.expressions.Expression; +import org.apache.flink.table.functions.FunctionIdentifier; +import org.apache.flink.util.Preconditions; +import org.apache.flink.util.StringUtils; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import software.amazon.awssdk.http.apache.ApacheHttpClient; +import software.amazon.awssdk.regions.Region; +import software.amazon.awssdk.services.glue.GlueClient; +import software.amazon.awssdk.services.glue.model.Database; +import software.amazon.awssdk.services.glue.model.Partition; +import software.amazon.awssdk.services.glue.model.PartitionInput; +import software.amazon.awssdk.services.glue.model.StorageDescriptor; +import software.amazon.awssdk.services.glue.model.Table; +import software.amazon.awssdk.services.glue.model.TableInput; +import software.amazon.awssdk.services.glue.model.UserDefinedFunction; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Properties; +import java.util.stream.Collectors; + +/** + * A Flink {@link org.apache.flink.table.catalog.Catalog} backed by the AWS Glue Data Catalog. + * + * <p>Databases, tables, views, partitions and functions are stored as their Glue counterparts. + * Since Glue stores object names in lowercase, this catalog stores every object under its lowercase + * name and records the declared name in a Glue parameter, so identifiers are presented to Flink + * exactly as declared. Because Glue storage names are unique, every lookup resolves with a single + * {@code Get*} call on the lowercase name; the catalog never scans Glue to resolve a name. + * + * <p>Parts of a Flink schema that Glue columns cannot represent (computed and metadata columns, + * watermarks, primary keys, exact types such as {@code TIMESTAMP(3)}, NOT NULL constraints) are + * persisted as {@code flink.schema.*} table parameters and restored on read, so tables created by + * this catalog round-trip exactly. Tables created by other engines (Athena, Glue crawlers, Spark, + * Hive) are exposed with the schema Glue holds for them. + */ +@PublicEvolving +public class GlueCatalog extends AbstractCatalog { + + private static final Logger LOG = LoggerFactory.getLogger(GlueCatalog.class); + + /** Glue table types that Flink exposes as views. */ + private static final List<String> GLUE_VIEW_TABLE_TYPES = + Collections.unmodifiableList( + java.util.Arrays.asList( + CatalogBaseTable.TableKind.VIEW.name(), + "VIRTUAL_VIEW", + "MATERIALIZED_VIEW")); + + /** Characters Hive escapes in partition path segments, plus control characters. */ + private static final String PATH_UNSAFE_CHARS = "\"#%'*/:=?\\\u007F{[]^"; + + private GlueClient glueClient; + private GlueTypeConverter glueTypeConverter; + private GlueDatabaseOperator glueDatabaseOperations; + private GlueTableOperator glueTableOperations; + private GlueFunctionOperator glueFunctionsOperations; + private GluePartitionOperator gluePartitionOperations; + private GlueTableUtils glueTableUtils; + + /** + * Constructs a GlueCatalog with a provided Glue client. + * + * @param name the name of the catalog + * @param defaultDatabase the default database for the catalog + * @param region the AWS region to be used for Glue operations + * @param glueClient the Glue client to use; when null a default client for the region is built + */ + @VisibleForTesting + GlueCatalog(String name, String defaultDatabase, String region, GlueClient glueClient) { + super(name, defaultDatabase); + Preconditions.checkNotNull(region, "region cannot be null"); + Preconditions.checkArgument(!region.trim().isEmpty(), "region cannot be empty"); + + if (glueClient != null) { + setup(glueClient); + } else { + setup(GlueClient.builder().region(Region.of(region)).build()); + } + } + + /** + * Constructs a GlueCatalog with default client configuration. + * + * @param name the name of the catalog + * @param defaultDatabase the default database for the catalog + * @param region the AWS region to be used for Glue operations + */ + public GlueCatalog(String name, String defaultDatabase, String region) { + this(name, defaultDatabase, region, new Properties()); + } + + /** + * Constructs a GlueCatalog whose Glue client is built from the given AWS client properties, + * using the same client-creation path as the other AWS connectors ({@link AWSClientUtil}). This + * makes the standard {@code aws.*} settings available to the catalog - for example {@code + * aws.credentials.provider} to select a credential mode, {@code aws.endpoint} to point at a + * Glue-compatible endpoint, and the {@code aws.http-client.*} options. + * + * @param name the name of the catalog + * @param defaultDatabase the default database for the catalog + * @param region the AWS region to be used for Glue operations + * @param glueClientProperties AWS client properties, keyed by {@link AWSConfigConstants} + */ + public GlueCatalog( + String name, String defaultDatabase, String region, Properties glueClientProperties) { + super(name, defaultDatabase); + Preconditions.checkNotNull(region, "region cannot be null"); + Preconditions.checkArgument(!region.trim().isEmpty(), "region cannot be empty"); + Preconditions.checkNotNull(glueClientProperties, "glueClientProperties cannot be null"); + + Properties clientProperties = new Properties(); + clientProperties.putAll(glueClientProperties); + // The explicit region argument wins over any aws.region property. + clientProperties.setProperty(AWSConfigConstants.AWS_REGION, region); + AWSGeneralUtil.validateAwsConfiguration(clientProperties); + + GlueClient client = + AWSClientUtil.createAwsSyncClient( + clientProperties, + AWSGeneralUtil.createSyncHttpClient( + clientProperties, ApacheHttpClient.builder()), + GlueClient.builder(), + GlueCatalogConstants.BASE_GLUE_USER_AGENT_PREFIX_FORMAT, + GlueCatalogConstants.GLUE_CLIENT_USER_AGENT_PREFIX); + setup(client); + } + + private void setup(GlueClient glueClient) { + this.glueClient = glueClient; + this.glueTypeConverter = new GlueTypeConverter(); + this.glueTableUtils = new GlueTableUtils(glueTypeConverter); + this.glueDatabaseOperations = new GlueDatabaseOperator(glueClient, getName()); + this.glueTableOperations = new GlueTableOperator(glueClient, getName()); + this.glueFunctionsOperations = new GlueFunctionOperator(glueClient, getName()); + this.gluePartitionOperations = new GluePartitionOperator(glueClient, getName()); + } + + @Override + public void open() throws CatalogException { + LOG.info("Opening GlueCatalog '{}'", getName()); + } + + @Override + public void close() throws CatalogException { + if (glueClient != null) { + LOG.info("Closing GlueCatalog '{}'", getName()); + glueClient.close(); + } + } + + // ------------------------------------------------------------------------------------------ + // Databases + // ------------------------------------------------------------------------------------------ + + @Override + public List<String> listDatabases() throws CatalogException { + return glueDatabaseOperations.listDatabases(); + } + + @Override + public CatalogDatabase getDatabase(String databaseName) + throws DatabaseNotExistException, CatalogException { + checkDatabaseName(databaseName); + return glueDatabaseOperations.getDatabase(databaseName); + } + + @Override + public boolean databaseExists(String databaseName) throws CatalogException { + checkDatabaseName(databaseName); + return glueDatabaseOperations.glueDatabaseExists(databaseName); + } + + @Override + public void createDatabase( + String databaseName, CatalogDatabase catalogDatabase, boolean ifNotExists) + throws DatabaseAlreadyExistException, CatalogException { + checkDatabaseName(databaseName); + Preconditions.checkNotNull(catalogDatabase, "CatalogDatabase cannot be null"); + + // One GetDatabase both answers "does it exist" and tells us how it was declared, so the + // error can explain a case-only clash (Glue stores names in lowercase). + Database existing = glueDatabaseOperations.getGlueDatabaseOrNull(databaseName); + if (existing != null) { + if (ifNotExists) { + return; + } + throw databaseAlreadyExists( + databaseName, glueDatabaseOperations.getOriginalDatabaseName(existing)); + } + glueDatabaseOperations.createDatabase(databaseName, catalogDatabase); + } + + @Override + public void dropDatabase(String databaseName, boolean ignoreIfNotExists, boolean cascade) + throws DatabaseNotExistException, DatabaseNotEmptyException, CatalogException { + checkDatabaseName(databaseName); + + String glueDatabaseName = glueDatabaseOperations.findGlueDatabaseName(databaseName); + if (glueDatabaseName == null) { + if (ignoreIfNotExists) { + return; + } + throw new DatabaseNotExistException(getName(), databaseName); + } + + // GetTables already returns views (they are Glue tables), so one listing covers both. + List<Table> tables = glueTableOperations.getAllGlueTables(glueDatabaseName); + List<String> functions = glueFunctionsOperations.listGlueFunctions(glueDatabaseName); + if (!tables.isEmpty() || !functions.isEmpty()) { + if (!cascade) { + throw new DatabaseNotEmptyException(getName(), databaseName); + } + for (Table table : tables) { + try { + glueTableOperations.dropTable(glueDatabaseName, table.name()); + } catch (TableNotExistException e) { + LOG.debug("Table {} vanished during cascading drop", table.name()); + } + } + for (String function : functions) { + try { + glueFunctionsOperations.dropGlueFunction( + new ObjectPath(glueDatabaseName, function)); + } catch (FunctionNotExistException e) { + LOG.debug("Function {} vanished during cascading drop", function); + } + } + } + glueDatabaseOperations.dropGlueDatabase(databaseName); + } + + @Override + public void alterDatabase( + String databaseName, CatalogDatabase catalogDatabase, boolean ignoreIfNotExists) + throws DatabaseNotExistException, CatalogException { + throw new UnsupportedOperationException( + "Altering databases is not supported by the Glue Catalog."); + } + + // ------------------------------------------------------------------------------------------ + // Tables and views + // ------------------------------------------------------------------------------------------ + + @Override + public List<String> listTables(String databaseName) + throws DatabaseNotExistException, CatalogException { + return glueTableOperations.listTables(requireGlueDatabaseName(databaseName)); + } + + @Override + public List<String> listViews(String databaseName) + throws DatabaseNotExistException, CatalogException { + String glueDatabaseName = requireGlueDatabaseName(databaseName); + return glueTableOperations.getAllGlueTables(glueDatabaseName).stream() + .filter(table -> resolveTableKind(table) == CatalogBaseTable.TableKind.VIEW) + .map(glueTableOperations::getOriginalTableName) + .collect(Collectors.toList()); + } + + @Override + public CatalogBaseTable getTable(ObjectPath objectPath) + throws TableNotExistException, CatalogException { + Table glueTable = getGlueTableOrNull(objectPath); + if (glueTable == null) { + throw new TableNotExistException(getName(), objectPath); + } + return toCatalogBaseTable(objectPath, glueTable); + } + + @Override + public boolean tableExists(ObjectPath objectPath) throws CatalogException { + return getGlueTableOrNull(objectPath) != null; + } + + @Override + public void dropTable(ObjectPath objectPath, boolean ifExists) + throws TableNotExistException, CatalogException { + Preconditions.checkNotNull(objectPath, "ObjectPath cannot be null"); + String glueDatabaseName = + glueDatabaseOperations.findGlueDatabaseName(objectPath.getDatabaseName()); + if (glueDatabaseName == null) { + if (ifExists) { + return; + } + throw new TableNotExistException(getName(), objectPath); + } + try { + glueTableOperations.dropTable(glueDatabaseName, objectPath.getObjectName()); + } catch (TableNotExistException e) { + if (!ifExists) { + throw new TableNotExistException(getName(), objectPath, e); + } + } + } + + @Override + public void createTable( + ObjectPath objectPath, CatalogBaseTable catalogBaseTable, boolean ifNotExists) + throws TableAlreadyExistException, DatabaseNotExistException, CatalogException { + Preconditions.checkNotNull(objectPath, "ObjectPath cannot be null"); + Preconditions.checkNotNull(catalogBaseTable, "CatalogBaseTable cannot be null"); + + String glueDatabaseName = requireGlueDatabaseName(objectPath.getDatabaseName()); + + // One GetTable both answers "does it exist" and tells us how it was declared, so the + // error can explain a case-only clash (Glue stores names in lowercase). + Table existing = + glueTableOperations.getGlueTableOrNull( + glueDatabaseName, objectPath.getObjectName()); + if (existing != null) { + if (ifNotExists) { + return; + } + throw tableAlreadyExists( + objectPath, glueTableOperations.getOriginalTableName(existing)); + } + + ResolvedSchema resolvedSchema = requireResolvedSchema(objectPath, catalogBaseTable); + Map<String, String> options = new HashMap<>(catalogBaseTable.getOptions()); + rejectReservedOptions(objectPath, options); + + TableInput tableInput; + switch (catalogBaseTable.getTableKind()) { + case TABLE: + CatalogTable catalogTable = (CatalogTable) catalogBaseTable; + tableInput = + buildTableInput( + objectPath.getObjectName(), + CatalogBaseTable.TableKind.TABLE, + catalogTable.getComment(), + resolvedSchema, + catalogTable.getPartitionKeys(), + options, + glueTableUtils.resolveTableLocation( + options, + GlueDatabaseOperator.toGlueDatabaseName( + objectPath.getDatabaseName()), + GlueTableOperator.toGlueTableName( + objectPath.getObjectName()), + !catalogTable.getPartitionKeys().isEmpty())); + break; + case VIEW: + CatalogView catalogView = (CatalogView) catalogBaseTable; + tableInput = + buildTableInput( + objectPath.getObjectName(), + CatalogBaseTable.TableKind.VIEW, + catalogView.getComment(), + resolvedSchema, + Collections.emptyList(), + options, + null) + .toBuilder() + .viewOriginalText(catalogView.getOriginalQuery()) + .viewExpandedText(catalogView.getExpandedQuery()) + .build(); + break; + default: + throw new CatalogException( + String.format( + "Cannot create %s: table kind %s is not supported by the Glue Catalog", + objectPath.getFullName(), catalogBaseTable.getTableKind())); + } + + try { + glueTableOperations.createTable(glueDatabaseName, tableInput); + } catch (TableAlreadyExistException e) { + // Created concurrently between our existence check and the create. + if (!ifNotExists) { + throw new TableAlreadyExistException(getName(), objectPath, e); + } + } + LOG.info( + "Created {} {} in Glue", catalogBaseTable.getTableKind(), objectPath.getFullName()); + } + + @Override + public void renameTable(ObjectPath objectPath, String newTableName, boolean ignoreIfNotExists) + throws TableNotExistException, TableAlreadyExistException, CatalogException { + throw new UnsupportedOperationException( + "Renaming tables is not supported by the Glue Catalog."); + } + + @Override + public void alterTable( + ObjectPath objectPath, CatalogBaseTable newTable, boolean ignoreIfNotExists) + throws TableNotExistException, CatalogException { + Preconditions.checkNotNull(objectPath, "ObjectPath cannot be null"); + Preconditions.checkNotNull(newTable, "CatalogBaseTable cannot be null"); + + Table existing = getGlueTableOrNull(objectPath); + if (existing == null) { + if (ignoreIfNotExists) { + return; + } + throw new TableNotExistException(getName(), objectPath); + } + + if (newTable.getTableKind() != CatalogBaseTable.TableKind.TABLE) { + throw new UnsupportedOperationException( + "Altering views is not supported by the Glue Catalog."); + } + if (resolveTableKind(existing) != CatalogBaseTable.TableKind.TABLE) { + throw new CatalogException( + String.format( + "Cannot alter %s as a table: the Glue object is a %s", + objectPath.getFullName(), existing.tableType())); + } + + CatalogTable catalogTable = (CatalogTable) newTable; + List<String> existingPartitionKeys = GlueTableUtils.getPartitionKeyNames(existing); + if (!existingPartitionKeys.equals(catalogTable.getPartitionKeys())) { + throw new CatalogException( + String.format( + "Cannot alter %s: changing the partition keys (%s -> %s) is not " + + "supported because existing partitions would become " + + "unreachable", + objectPath.getFullName(), + existingPartitionKeys, + catalogTable.getPartitionKeys())); + } + + ResolvedSchema resolvedSchema = requireResolvedSchema(objectPath, newTable); + Map<String, String> options = new HashMap<>(catalogTable.getOptions()); + rejectReservedOptions(objectPath, options); + + String glueDatabaseName = + GlueDatabaseOperator.toGlueDatabaseName(objectPath.getDatabaseName()); + TableInput tableInput = + buildTableInput( + // Preserve the originally declared table name (case) across the alter. + glueTableOperations.getOriginalTableName(existing), + CatalogBaseTable.TableKind.TABLE, + catalogTable.getComment(), + resolvedSchema, + catalogTable.getPartitionKeys(), + options, + glueTableUtils.resolveTableLocation( + options, + glueDatabaseName, + existing.name(), + !catalogTable.getPartitionKeys().isEmpty())); + glueTableOperations.updateTable(glueDatabaseName, tableInput); + LOG.info("Altered table {} in Glue", objectPath.getFullName()); + } + + // ------------------------------------------------------------------------------------------ + // Partitions + // ------------------------------------------------------------------------------------------ + + @Override + public List<CatalogPartitionSpec> listPartitions(ObjectPath objectPath) + throws TableNotExistException, TableNotPartitionedException, CatalogException { + GlueTableRef tableRef = resolvePartitionedTable(objectPath); + List<String> partitionKeys = tableRef.partitionKeys(); + return gluePartitionOperations + .listPartitions(tableRef.databaseName, tableRef.tableName) + .stream() + .map(partition -> toPartitionSpec(partitionKeys, partition.values())) + .collect(Collectors.toList()); + } + + @Override + public List<CatalogPartitionSpec> listPartitions( + ObjectPath objectPath, CatalogPartitionSpec catalogPartitionSpec) + throws TableNotExistException, + TableNotPartitionedException, + PartitionSpecInvalidException, + CatalogException { + GlueTableRef tableRef = resolvePartitionedTable(objectPath); + List<String> partitionKeys = tableRef.partitionKeys(); + + Map<String, String> partialSpec = + catalogPartitionSpec == null + ? Collections.emptyMap() + : catalogPartitionSpec.getPartitionSpec(); + // Flink's Catalog contract: a partial spec referencing unknown partition keys is invalid. + if (!partitionKeys.containsAll(partialSpec.keySet())) { + throw new PartitionSpecInvalidException( + getName(), partitionKeys, objectPath, catalogPartitionSpec); + } + + return gluePartitionOperations + .listPartitions(tableRef.databaseName, tableRef.tableName) + .stream() + .map(partition -> toPartitionSpec(partitionKeys, partition.values())) + .filter( + spec -> + spec.getPartitionSpec() + .entrySet() + .containsAll(partialSpec.entrySet())) + .collect(Collectors.toList()); + } + + @Override + public List<CatalogPartitionSpec> listPartitionsByFilter( + ObjectPath objectPath, List<Expression> filters) + throws TableNotExistException, TableNotPartitionedException, CatalogException { + // Expression push-down to Glue partition filters is not implemented. Flink's planner + // catches UnsupportedOperationException and falls back to listPartitions(). + throw new UnsupportedOperationException( Review Comment: Expected: `listPartitionsByFilter` is not implemented, and the planner catches the `UnsupportedOperationException` and falls back to `listPartitions`, then evaluates the partition predicate on the Flink side. Queries over partitioned tables stay correct but list every partition from Glue, which counts against the Glue API quotas. Documented in the Limitations section in 79cdcff. While doing that I checked every `UnsupportedOperationException` in main code against the docs: six in code, two were documented. The section now lists all of them (alterDatabase, renameTable, alterTable on a view or changing partition keys, the statistics API, and this one). -- 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]
