fmorillo7694 commented on code in PR #206:
URL: 
https://github.com/apache/flink-connector-aws/pull/206#discussion_r4142316213


##########
flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/util/GlueTypeConverter.java:
##########
@@ -0,0 +1,356 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.table.catalog.glue.util;
+
+import org.apache.flink.annotation.Internal;
+import org.apache.flink.table.api.DataTypes;
+import 
org.apache.flink.table.catalog.glue.exception.UnsupportedDataTypeMappingException;
+import org.apache.flink.table.types.DataType;
+import org.apache.flink.table.types.logical.ArrayType;
+import org.apache.flink.table.types.logical.DecimalType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.LogicalTypeRoot;
+import org.apache.flink.table.types.logical.MapType;
+import org.apache.flink.table.types.logical.RowType;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
+
+/**
+ * Utility class for converting Flink types to Glue types and vice versa. 
Supports the conversion of
+ * common primitive, array, map, and struct types.
+ */
+@Internal
+public class GlueTypeConverter {
+
+    /** Logger for tracking Glue type conversions. */
+    private static final Logger LOG = 
LoggerFactory.getLogger(GlueTypeConverter.class);
+
+    /** Regular expressions for handling specific Glue types. */
+    private static final Pattern DECIMAL_PATTERN =
+            Pattern.compile("decimal\\((\\d+),(\\d+)\\)", 
Pattern.CASE_INSENSITIVE);
+
+    /**
+     * Parameterized Hive/Glue string types, e.g. {@code varchar(255)} or 
{@code char(10)}, written
+     * by other engines (Athena, Spark, Hive). Flink's Glue catalog stores 
strings as
+     * unparameterized {@code string}, but foreign tables may use the 
parameterized form, so the
+     * read path must accept it. The declared length is preserved as a Flink 
{@code
+     * VARCHAR(n)}/{@code CHAR(n)}.
+     */
+    private static final Pattern VARCHAR_PATTERN =
+            Pattern.compile("varchar\\((\\d+)\\)", Pattern.CASE_INSENSITIVE);
+
+    private static final Pattern CHAR_PATTERN =
+            Pattern.compile("char\\((\\d+)\\)", Pattern.CASE_INSENSITIVE);
+
+    private static final Pattern ARRAY_PATTERN =
+            Pattern.compile("array<(.+)>", Pattern.CASE_INSENSITIVE);
+    private static final Pattern MAP_PATTERN =
+            Pattern.compile("map<(.+),(.+)>", Pattern.CASE_INSENSITIVE);
+    private static final Pattern STRUCT_PATTERN =
+            Pattern.compile("struct<(.+)>", Pattern.CASE_INSENSITIVE);
+
+    /**
+     * Converts a Flink DataType to its corresponding Glue type as a string.
+     *
+     * @param flinkType The Flink DataType to be converted.
+     * @return The Glue type as a string.
+     */
+    public String toGlueDataType(DataType flinkType) {
+        LogicalType logicalType = flinkType.getLogicalType();
+        LogicalTypeRoot typeRoot = logicalType.getTypeRoot();
+
+        // Handle various Flink types and map them to corresponding Glue types
+        switch (typeRoot) {
+            case CHAR:
+            case VARCHAR:
+                return "string";
+            case BOOLEAN:
+                return "boolean";
+            case BINARY:
+            case VARBINARY:
+                return "binary";
+            case DECIMAL:
+                DecimalType decimalType = (DecimalType) logicalType;
+                return String.format(
+                        "decimal(%d,%d)", decimalType.getPrecision(), 
decimalType.getScale());
+            case TINYINT:
+                return "tinyint";
+            case SMALLINT:
+                return "smallint";
+            case INTEGER:
+                return "int";
+            case BIGINT:
+                return "bigint";
+            case FLOAT:
+                return "float";
+            case DOUBLE:
+                return "double";
+            case DATE:
+                return "date";
+            case TIME_WITHOUT_TIME_ZONE:
+                return "string"; // Glue doesn't have a direct time type, use 
string
+            case TIMESTAMP_WITHOUT_TIME_ZONE:
+            case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+                return "timestamp";
+            case ARRAY:
+                ArrayType arrayType = (ArrayType) logicalType;
+                return "array<" + 
toGlueDataType(DataTypes.of(arrayType.getElementType())) + ">";
+            case MAP:
+                MapType mapType = (MapType) logicalType;
+                return String.format(
+                        "map<%s,%s>",
+                        toGlueDataType(DataTypes.of(mapType.getKeyType())),
+                        toGlueDataType(DataTypes.of(mapType.getValueType())));
+            case ROW:
+                RowType rowType = (RowType) logicalType;
+                StringBuilder structBuilder = new StringBuilder("struct<");
+                for (int i = 0; i < rowType.getFieldCount(); i++) {
+                    if (i > 0) {
+                        structBuilder.append(",");
+                    }
+                    // Keep original field name for nested structs
+                    structBuilder
+                            .append(rowType.getFieldNames().get(i))
+                            .append(":")
+                            
.append(toGlueDataType(DataTypes.of(rowType.getChildren().get(i))));
+                }
+                structBuilder.append(">");
+                return structBuilder.toString();
+            default:
+                throw new UnsupportedDataTypeMappingException(
+                        "Flink type not supported by Glue Catalog: " + 
logicalType.getTypeRoot());
+        }
+    }
+
+    /**
+     * Converts a Glue type (as a string) to the corresponding Flink DataType.
+     *
+     * @param glueType The Glue type as a string.
+     * @return The corresponding Flink DataType.
+     * @throws IllegalArgumentException if the Glue type is invalid or unknown.
+     */
+    public DataType toFlinkDataType(String glueType) {
+        if (glueType == null || glueType.trim().isEmpty()) {
+            throw new IllegalArgumentException("Glue type cannot be null or 
empty");
+        }
+
+        // Trim but don't lowercase - we'll handle case-insensitivity per type
+        String trimmedGlueType = glueType.trim();
+
+        // Handle DECIMAL type
+        Matcher decimalMatcher = DECIMAL_PATTERN.matcher(trimmedGlueType);
+        if (decimalMatcher.matches()) {
+            int precision = Integer.parseInt(decimalMatcher.group(1));
+            int scale = Integer.parseInt(decimalMatcher.group(2));
+            return DataTypes.DECIMAL(precision, scale);
+        }
+
+        // Handle parameterized string types written by other engines, e.g. 
varchar(255),
+        // char(10). Flink itself stores strings as unparameterized "string", 
but foreign tables
+        // may carry the parameterized Hive form; preserve the declared length.
+        Matcher varcharMatcher = VARCHAR_PATTERN.matcher(trimmedGlueType);
+        if (varcharMatcher.matches()) {
+            return 
DataTypes.VARCHAR(Integer.parseInt(varcharMatcher.group(1)));
+        }
+        Matcher charMatcher = CHAR_PATTERN.matcher(trimmedGlueType);
+        if (charMatcher.matches()) {
+            return DataTypes.CHAR(Integer.parseInt(charMatcher.group(1)));
+        }
+
+        // Handle ARRAY type - keyword matched case-insensitively, content 
preserved verbatim
+        Matcher arrayMatcher = ARRAY_PATTERN.matcher(trimmedGlueType);
+        if (arrayMatcher.matches()) {
+            // Extract from original string to preserve case in content
+            int contentStart = trimmedGlueType.indexOf('<') + 1;
+            int contentEnd = trimmedGlueType.lastIndexOf('>');
+            String elementType = trimmedGlueType.substring(contentStart, 
contentEnd);
+            return DataTypes.ARRAY(toFlinkDataType(elementType));
+        }
+
+        // Handle MAP type - keyword matched case-insensitively, content 
preserved verbatim
+        Matcher mapMatcher = MAP_PATTERN.matcher(trimmedGlueType);
+        if (mapMatcher.matches()) {
+            // Extract from original string to preserve case in content
+            int contentStart = trimmedGlueType.indexOf('<') + 1;
+            int contentEnd = trimmedGlueType.lastIndexOf('>');
+            String mapContent = trimmedGlueType.substring(contentStart, 
contentEnd);
+
+            // Split key and value types
+            int commaPos = findMapTypeSeparator(mapContent);
+            if (commaPos < 0) {
+                throw new IllegalArgumentException("Invalid map type format: " 
+ glueType);
+            }
+
+            String keyType = mapContent.substring(0, commaPos).trim();
+            String valueType = mapContent.substring(commaPos + 1).trim();
+
+            return DataTypes.MAP(toFlinkDataType(keyType), 
toFlinkDataType(valueType));
+        }
+
+        // Handle STRUCT type - keyword matched case-insensitively, content 
preserved verbatim
+        Matcher structMatcher = STRUCT_PATTERN.matcher(trimmedGlueType);
+        if (structMatcher.matches()) {
+            // Extract from original string to preserve case in field names
+            int contentStart = trimmedGlueType.indexOf('<') + 1;
+            int contentEnd = trimmedGlueType.lastIndexOf('>');
+            String structContent = trimmedGlueType.substring(contentStart, 
contentEnd);
+
+            return parseStructType(structContent);
+        }
+
+        // Handle primitive types (case insensitive)
+        switch (trimmedGlueType.toLowerCase()) {
+            case "string":
+            case "char":
+            case "varchar":
+                return DataTypes.STRING();
+            case "boolean":
+                return DataTypes.BOOLEAN();
+            case "binary":
+                return DataTypes.BYTES();
+            case "tinyint":
+                return DataTypes.TINYINT();
+            case "smallint":
+                return DataTypes.SMALLINT();
+            case "int":
+            case "integer": // Athena DDL alias
+                return DataTypes.INT();
+            case "bigint":
+                return DataTypes.BIGINT();
+            case "float":
+                return DataTypes.FLOAT();
+            case "double":
+                return DataTypes.DOUBLE();
+            case "date":
+                return DataTypes.DATE();
+            case "timestamp":
+                return DataTypes.TIMESTAMP();
+            case "decimal":
+                // Hive/Glue default precision and scale for an 
unparameterized decimal.
+                return DataTypes.DECIMAL(10, 0);
+            default:
+                if (trimmedGlueType.toLowerCase().startsWith("uniontype<")) {
+                    throw new UnsupportedDataTypeMappingException(
+                            "Unsupported Glue type: "
+                                    + glueType
+                                    + " (Hive union types have no Flink 
equivalent)");
+                }
+                throw new UnsupportedDataTypeMappingException("Unsupported 
Glue type: " + glueType);
+        }
+    }
+
+    /**
+     * Helper method to find the comma that separates key and value types in a 
map. Handles nested
+     * types correctly by tracking angle brackets.
+     *
+     * @param mapContent The content of the map type definition.
+     * @return The position of the separator comma, or -1 if not found.
+     */
+    private int findMapTypeSeparator(String mapContent) {

Review Comment:
   Yes, it did. findMapTypeSeparator only tracked angle brackets, so the comma 
inside decimal(10,2) was taken as the key/value separator and the map failed 
with an invalid type. It now tracks parentheses as well (same as 
splitStructFields already did). Fixed in cea9af8 with a parameterized test 
covering map<decimal(10,2),string>, map<string,decimal(10,2)>, both sides 
parameterized, varchar(n)/char(n) keys, and nested array/struct cases; all of 
them fail on the previous code.



##########
flink-catalog-aws/flink-catalog-aws-glue/src/main/java/org/apache/flink/table/catalog/glue/util/GlueTypeConverter.java:
##########
@@ -0,0 +1,356 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.table.catalog.glue.util;
+
+import org.apache.flink.annotation.Internal;
+import org.apache.flink.table.api.DataTypes;
+import 
org.apache.flink.table.catalog.glue.exception.UnsupportedDataTypeMappingException;
+import org.apache.flink.table.types.DataType;
+import org.apache.flink.table.types.logical.ArrayType;
+import org.apache.flink.table.types.logical.DecimalType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.LogicalTypeRoot;
+import org.apache.flink.table.types.logical.MapType;
+import org.apache.flink.table.types.logical.RowType;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
+
+/**
+ * Utility class for converting Flink types to Glue types and vice versa. 
Supports the conversion of
+ * common primitive, array, map, and struct types.
+ */
+@Internal
+public class GlueTypeConverter {
+
+    /** Logger for tracking Glue type conversions. */
+    private static final Logger LOG = 
LoggerFactory.getLogger(GlueTypeConverter.class);
+
+    /** Regular expressions for handling specific Glue types. */
+    private static final Pattern DECIMAL_PATTERN =
+            Pattern.compile("decimal\\((\\d+),(\\d+)\\)", 
Pattern.CASE_INSENSITIVE);
+
+    /**
+     * Parameterized Hive/Glue string types, e.g. {@code varchar(255)} or 
{@code char(10)}, written
+     * by other engines (Athena, Spark, Hive). Flink's Glue catalog stores 
strings as
+     * unparameterized {@code string}, but foreign tables may use the 
parameterized form, so the
+     * read path must accept it. The declared length is preserved as a Flink 
{@code
+     * VARCHAR(n)}/{@code CHAR(n)}.
+     */
+    private static final Pattern VARCHAR_PATTERN =
+            Pattern.compile("varchar\\((\\d+)\\)", Pattern.CASE_INSENSITIVE);
+
+    private static final Pattern CHAR_PATTERN =
+            Pattern.compile("char\\((\\d+)\\)", Pattern.CASE_INSENSITIVE);
+
+    private static final Pattern ARRAY_PATTERN =
+            Pattern.compile("array<(.+)>", Pattern.CASE_INSENSITIVE);
+    private static final Pattern MAP_PATTERN =
+            Pattern.compile("map<(.+),(.+)>", Pattern.CASE_INSENSITIVE);
+    private static final Pattern STRUCT_PATTERN =
+            Pattern.compile("struct<(.+)>", Pattern.CASE_INSENSITIVE);
+
+    /**
+     * Converts a Flink DataType to its corresponding Glue type as a string.
+     *
+     * @param flinkType The Flink DataType to be converted.
+     * @return The Glue type as a string.
+     */
+    public String toGlueDataType(DataType flinkType) {
+        LogicalType logicalType = flinkType.getLogicalType();
+        LogicalTypeRoot typeRoot = logicalType.getTypeRoot();
+
+        // Handle various Flink types and map them to corresponding Glue types
+        switch (typeRoot) {
+            case CHAR:
+            case VARCHAR:
+                return "string";
+            case BOOLEAN:
+                return "boolean";
+            case BINARY:
+            case VARBINARY:
+                return "binary";
+            case DECIMAL:
+                DecimalType decimalType = (DecimalType) logicalType;
+                return String.format(
+                        "decimal(%d,%d)", decimalType.getPrecision(), 
decimalType.getScale());
+            case TINYINT:
+                return "tinyint";
+            case SMALLINT:
+                return "smallint";
+            case INTEGER:
+                return "int";
+            case BIGINT:
+                return "bigint";
+            case FLOAT:
+                return "float";
+            case DOUBLE:
+                return "double";
+            case DATE:
+                return "date";
+            case TIME_WITHOUT_TIME_ZONE:
+                return "string"; // Glue doesn't have a direct time type, use 
string
+            case TIMESTAMP_WITHOUT_TIME_ZONE:
+            case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+                return "timestamp";
+            case ARRAY:
+                ArrayType arrayType = (ArrayType) logicalType;
+                return "array<" + 
toGlueDataType(DataTypes.of(arrayType.getElementType())) + ">";
+            case MAP:
+                MapType mapType = (MapType) logicalType;
+                return String.format(
+                        "map<%s,%s>",
+                        toGlueDataType(DataTypes.of(mapType.getKeyType())),
+                        toGlueDataType(DataTypes.of(mapType.getValueType())));
+            case ROW:
+                RowType rowType = (RowType) logicalType;
+                StringBuilder structBuilder = new StringBuilder("struct<");
+                for (int i = 0; i < rowType.getFieldCount(); i++) {
+                    if (i > 0) {
+                        structBuilder.append(",");
+                    }
+                    // Keep original field name for nested structs
+                    structBuilder
+                            .append(rowType.getFieldNames().get(i))
+                            .append(":")
+                            
.append(toGlueDataType(DataTypes.of(rowType.getChildren().get(i))));
+                }
+                structBuilder.append(">");
+                return structBuilder.toString();
+            default:
+                throw new UnsupportedDataTypeMappingException(
+                        "Flink type not supported by Glue Catalog: " + 
logicalType.getTypeRoot());
+        }
+    }
+
+    /**
+     * Converts a Glue type (as a string) to the corresponding Flink DataType.
+     *
+     * @param glueType The Glue type as a string.
+     * @return The corresponding Flink DataType.
+     * @throws IllegalArgumentException if the Glue type is invalid or unknown.
+     */
+    public DataType toFlinkDataType(String glueType) {
+        if (glueType == null || glueType.trim().isEmpty()) {
+            throw new IllegalArgumentException("Glue type cannot be null or 
empty");
+        }
+
+        // Trim but don't lowercase - we'll handle case-insensitivity per type
+        String trimmedGlueType = glueType.trim();
+
+        // Handle DECIMAL type
+        Matcher decimalMatcher = DECIMAL_PATTERN.matcher(trimmedGlueType);
+        if (decimalMatcher.matches()) {
+            int precision = Integer.parseInt(decimalMatcher.group(1));
+            int scale = Integer.parseInt(decimalMatcher.group(2));
+            return DataTypes.DECIMAL(precision, scale);
+        }
+
+        // Handle parameterized string types written by other engines, e.g. 
varchar(255),
+        // char(10). Flink itself stores strings as unparameterized "string", 
but foreign tables
+        // may carry the parameterized Hive form; preserve the declared length.
+        Matcher varcharMatcher = VARCHAR_PATTERN.matcher(trimmedGlueType);
+        if (varcharMatcher.matches()) {
+            return 
DataTypes.VARCHAR(Integer.parseInt(varcharMatcher.group(1)));
+        }
+        Matcher charMatcher = CHAR_PATTERN.matcher(trimmedGlueType);
+        if (charMatcher.matches()) {
+            return DataTypes.CHAR(Integer.parseInt(charMatcher.group(1)));
+        }
+
+        // Handle ARRAY type - keyword matched case-insensitively, content 
preserved verbatim
+        Matcher arrayMatcher = ARRAY_PATTERN.matcher(trimmedGlueType);
+        if (arrayMatcher.matches()) {
+            // Extract from original string to preserve case in content
+            int contentStart = trimmedGlueType.indexOf('<') + 1;
+            int contentEnd = trimmedGlueType.lastIndexOf('>');
+            String elementType = trimmedGlueType.substring(contentStart, 
contentEnd);
+            return DataTypes.ARRAY(toFlinkDataType(elementType));
+        }
+
+        // Handle MAP type - keyword matched case-insensitively, content 
preserved verbatim
+        Matcher mapMatcher = MAP_PATTERN.matcher(trimmedGlueType);
+        if (mapMatcher.matches()) {
+            // Extract from original string to preserve case in content
+            int contentStart = trimmedGlueType.indexOf('<') + 1;
+            int contentEnd = trimmedGlueType.lastIndexOf('>');
+            String mapContent = trimmedGlueType.substring(contentStart, 
contentEnd);
+
+            // Split key and value types
+            int commaPos = findMapTypeSeparator(mapContent);
+            if (commaPos < 0) {
+                throw new IllegalArgumentException("Invalid map type format: " 
+ glueType);
+            }
+
+            String keyType = mapContent.substring(0, commaPos).trim();
+            String valueType = mapContent.substring(commaPos + 1).trim();
+
+            return DataTypes.MAP(toFlinkDataType(keyType), 
toFlinkDataType(valueType));
+        }
+
+        // Handle STRUCT type - keyword matched case-insensitively, content 
preserved verbatim
+        Matcher structMatcher = STRUCT_PATTERN.matcher(trimmedGlueType);
+        if (structMatcher.matches()) {
+            // Extract from original string to preserve case in field names
+            int contentStart = trimmedGlueType.indexOf('<') + 1;
+            int contentEnd = trimmedGlueType.lastIndexOf('>');
+            String structContent = trimmedGlueType.substring(contentStart, 
contentEnd);
+
+            return parseStructType(structContent);
+        }
+
+        // Handle primitive types (case insensitive)
+        switch (trimmedGlueType.toLowerCase()) {
+            case "string":
+            case "char":
+            case "varchar":
+                return DataTypes.STRING();
+            case "boolean":
+                return DataTypes.BOOLEAN();
+            case "binary":
+                return DataTypes.BYTES();
+            case "tinyint":
+                return DataTypes.TINYINT();
+            case "smallint":
+                return DataTypes.SMALLINT();
+            case "int":
+            case "integer": // Athena DDL alias
+                return DataTypes.INT();
+            case "bigint":
+                return DataTypes.BIGINT();
+            case "float":
+                return DataTypes.FLOAT();
+            case "double":
+                return DataTypes.DOUBLE();
+            case "date":
+                return DataTypes.DATE();
+            case "timestamp":
+                return DataTypes.TIMESTAMP();
+            case "decimal":
+                // Hive/Glue default precision and scale for an 
unparameterized decimal.
+                return DataTypes.DECIMAL(10, 0);
+            default:
+                if (trimmedGlueType.toLowerCase().startsWith("uniontype<")) {
+                    throw new UnsupportedDataTypeMappingException(
+                            "Unsupported Glue type: "
+                                    + glueType
+                                    + " (Hive union types have no Flink 
equivalent)");
+                }
+                throw new UnsupportedDataTypeMappingException("Unsupported 
Glue type: " + glueType);
+        }
+    }
+
+    /**
+     * Helper method to find the comma that separates key and value types in a 
map. Handles nested
+     * types correctly by tracking angle brackets.
+     *
+     * @param mapContent The content of the map type definition.
+     * @return The position of the separator comma, or -1 if not found.
+     */
+    private int findMapTypeSeparator(String mapContent) {
+        int nestedLevel = 0;
+        for (int i = 0; i < mapContent.length(); i++) {
+            char c = mapContent.charAt(i);
+            if (c == '<') {
+                nestedLevel++;
+            } else if (c == '>') {
+                nestedLevel--;
+            } else if (c == ',' && nestedLevel == 0) {
+                return i;
+            }
+        }
+        return -1;
+    }
+
+    /**
+     * Parses a struct type definition and returns the corresponding Flink 
DataType.
+     *
+     * @param structDefinition The struct definition string to parse.
+     * @return The corresponding Flink ROW DataType.
+     */
+    public DataType parseStructType(String structDefinition) {
+        String[] fields = splitStructFields(structDefinition);
+        List<DataTypes.Field> flinkFields = new ArrayList<>();
+
+        for (String field : fields) {
+            // Important: We need to find the colon separator properly,
+            // as field names might contain characters like '<' for nested 
structs
+            int colonPos = field.indexOf(':');
+            if (colonPos < 0) {
+                LOG.warn("Invalid struct field definition (no colon found): 
{}", field);
+                continue;

Review Comment:
   Agreed, it should throw. Skipping the field returned a ROW with fewer fields 
than the Glue table declares and misaligned everything after it, which is worse 
than failing. It now throws UnsupportedDataTypeMappingException naming the 
offending field and the struct. Fixed in cea9af8 with tests for struct<a>, 
struct<a:int,b>, a double comma, and nested occurrences inside array and map.



##########
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(
+                "Listing partitions by filter expression is not supported by 
the Glue Catalog.");
+    }
+
+    @Override
+    public CatalogPartition getPartition(
+            ObjectPath objectPath, CatalogPartitionSpec catalogPartitionSpec)
+            throws PartitionNotExistException, CatalogException {
+        Partition partition = getGluePartitionOrNull(objectPath, 
catalogPartitionSpec);
+        if (partition == null) {
+            throw new PartitionNotExistException(getName(), objectPath, 
catalogPartitionSpec);
+        }
+
+        Map<String, String> properties = new HashMap<>();
+        if (partition.parameters() != null) {
+            properties.putAll(partition.parameters());
+        }
+        if (partition.storageDescriptor() != null
+                && partition.storageDescriptor().location() != null) {
+            properties.put(
+                    GlueCatalogConstants.PARTITION_LOCATION,

Review Comment:
   Good catch, that copy is a real trap. The round-trip itself is intentional 
and matches HiveCatalog, which exposes the storage location as 
hive.location-uri from getPartition and reads it back in createPartition and 
alterPartition, so a partition copied between catalogs keeps its location. What 
was missing is protection against copying properties between partitions of the 
same table. In 00cbd1b createPartition and alterPartition validate an explicit 
location: if it sits under the table location and its trailing segments spell a 
Hive-style key=value chain for this table's partition keys, the values must be 
this partition's, otherwise it is rejected with a message naming the other 
partition's location. Any other location (another bucket, a custom layout) is 
still accepted as a deliberate choice, and a table without a location has 
nothing to check. Tests: testPartitionLocationRoundTrip pins the contract 
(default location, explicit custom location, re-passing a partition's own 
properti
 es to alterPartition), testPartitionRejectsLocationOfAnotherPartition covers 
exactly your scenario for both create and alter. I also documented the 
partition operations and this contract in glue.md (en and zh), which did not 
mention partitions at all before.



-- 
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]

Reply via email to