如何配置AWS Glue作业以使用数据湖表定义的列类型(CSV场景)
如何让AWS Glue作业采用数据湖表定义的列类型读取CSV数据?
我明白你遇到的问题——当用Glue读取CSV格式的数据时,默认的create_dynamic_frame.from_catalog会自动推断Schema,完全忽略你在Glue数据湖表中定义的列类型和顺序,这确实会给CSV这类无自描述格式的数据带来麻烦。下面是两种可靠的解决方法:
方法一:用Spark SQL读取Glue Catalog表(推荐)
Spark SQL会严格遵循Glue Catalog中表的定义(包括列类型、顺序、分隔符等),读取后再转换成DynamicFrame就能完全复用你定义的Schema。修改你的代码如下:
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job from awsglue.dynamicframe import DynamicFrame from pyspark.sql import SparkSession ## @params: [JOB_NAME] args = getResolvedOptions(sys.argv, ['JOB_NAME']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) # 替换原来的create_dynamic_frame代码 # 直接用Spark SQL读取Glue Catalog中的表 medicare_df = spark.sql("SELECT * FROM my_database.my_table") # 转成Glue DynamicFrame medicare_dynamicframe = DynamicFrame.fromDF(medicare_df, glueContext, "medicare_df") medicare_dynamicframe.printSchema() job.commit()
这种方法的优势是代码简洁,不需要额外的Schema处理,完全依赖你在Glue Catalog中已经配置好的表定义,确保列类型和顺序和你预期的一致。
方法二:手动获取Catalog Schema并应用映射
如果你更倾向于使用from_catalog的方式,可以手动从Glue Catalog中获取表的Schema定义,然后通过ApplyMapping强制转换类型:
import sys import boto3 from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job from awsglue.dynamicframe import DynamicFrame from awsglue.client import GlueClient ## @params: [JOB_NAME] args = getResolvedOptions(sys.argv, ['JOB_NAME']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) # 从Catalog读取数据(注意根据CSV是否有表头设置skip.header.line.count) medicare_dynamicframe = glueContext.create_dynamic_frame.from_catalog( database = "my_database", table_name = "my_table", additional_options={"skip.header.line.count": "1"} # 如果CSV有表头请保留,没有则删除 ) # 获取Glue Catalog中的表Schema定义 glue_client = boto3.client('glue') table_details = glue_client.get_table(DatabaseName='my_database', Name='my_table') table_columns = table_details['Table']['StorageDescriptor']['Columns'] # 构建映射规则:(源列名, 目标列名, 目标类型) mapping_rules = [(col['Name'], col['Name'], col['Type']) for col in table_columns] # 应用映射,强制转换为Catalog定义的类型 medicare_dynamicframe = ApplyMapping.apply(frame=medicare_dynamicframe, mappings=mapping_rules) medicare_dynamicframe.printSchema() job.commit()
额外检查项
在尝试以上方法前,先确认你的Glue表定义是否正确:
- 确保
StorageDescriptor中的Columns列表完全匹配你期望的列名、类型和顺序 - 检查
SerdeInfo中的参数是否配置了正确的CSV分隔符(比如field.delim=,) - 如果CSV文件有表头,确保表的参数中设置了
skip.header.line.count=1
内容的提问来源于stack exchange,提问作者Cherry
相关产品推荐
相关产品推荐

