如何将S3中CSV文件转为AWS Athena可用SQL表?Glue实施遇障
问题:AWS Glue创建S3 CSV表结构失败,Athena无法正常访问数据
我在AWS S3存储桶中有一个CSV文件,想通过AWS Athena访问,需要用AWS Glue创建关系型数据库模式,但尝试两种方法都出了问题:
- 方法一:用Glue爬虫+自定义CSV分类器(指定列数据类型),但浮点和整数列无法正确保留,要么数据缺失,要么报“字符串无法转换为实数”错误
- 方法二:创建可视化ETL作业并编写Python脚本做列类型转换,但AWS Glue Catalog无法生成数据模式
示例数据
url,brand,color,type,material,item_height,item_depth,item_width,item_weight,condition,price,seller,numsold,rating https://www.ebay.com/itm/235324388478?_trkparms=5079%3A5000006518,Cuisinart,Silver,Blends,Stainless Steel,15,-1,-1,-1,"Certified - Refurbished Certified - Refurbished",149.99,Cuisinart,42,-1
方法一:Glue爬虫配置
使用指定列数据类型的CSV分类器运行爬虫,但处理数值列时出现类型转换或数据丢失错误。
方法二:ETL Python脚本
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 DynamicFrameCollection from awsglue.dynamicframe import DynamicFrame # Script generated for node Custom Transform def MyTransform(glueContext, dfc) -> DynamicFrameCollection: dynf = dfc.select(list(dfc.keys())[0]) df = dynf.toDF() from pyspark.sql.functions import col, regexp_extract from pyspark.sql.types import FloatType, IntegerType from awsglue.dynamicframe import DynamicFrame pattern = r"[\d.]+" columns_to_transform = ['price', 'rating', 'item_width', 'item_weight', 'item_height', 'item_depth', 'numsold'] for col_name in columns_to_transform: # Check if column exists in the DataFrame if col_name in df.columns: # Apply regex pattern to extract numerical values and convert to appropriate types df = df.withColumn(col_name, regexp_extract(col(col_name), pattern, 0).cast(FloatType() if col_name != 'numsold' else IntegerType())) res = DynamicFrame.fromDF(df, glueContext, 'changed') return DynamicFrameCollection({'Custom_Transform_0': res}, glueContext) args = getResolvedOptions(sys.argv, ['JOB_NAME']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) # Script generated for node Amazon S3 AmazonS3_node1713841179360 = glueContext.create_dynamic_frame.from_options(format_options={"quoteChar": "\"", "withHeader": True, "separator": ",", "optimizePerformance": False}, connection_type="s3", format="csv", connection_options={"paths": ["s3://ecommerce-deals"]}, transformation_ctx="AmazonS3_node1713841179360") # Script generated for node Custom Transform CustomTransform_node1713841197747 = MyTransform(glueContext, DynamicFrameCollection({"AmazonS3_node1713841179360": AmazonS3_node1713841179360}, glueContext)) # Script generated for node Select From Collection SelectFromCollection_node1713842428985 = SelectFromCollection.apply(dfc=CustomTransform_node1713841197747, key=list(CustomTransform_node1713841197747.keys())[0], transformation_ctx="SelectFromCollection_node1713842428985") # Script generated for node AWS Glue Data Catalog AWSGlueDataCatalog_node1713841379233 = glueContext.write_dynamic_frame.from_catalog(frame=SelectFromCollection_node1713842428985, database="default", table_name="ebay_csv", additional_options={"enableUpdateCatalog": True, "updateBehavior": "UPDATE_IN_DATABASE"}, transformation_ctx="AWSGlueDataCatalog_node1713841379233") job.commit()
可视化ETL作业结构
作业包含S3数据源节点、自定义转换节点、结果选择节点、Glue Catalog写入节点四个核心环节。
问题排查与解决方案
针对方法一的问题
数值列异常原因:
- 示例数据中
condition列包含换行符,会导致CSV解析错位,后续数值列被错误读取为字符串;-1这类特殊值也可能被分类器识别为非数值类型。 - 爬虫默认采样比例低,若采样行包含格式异常数据,会错误推断列类型。
- 示例数据中
修复方案:
- 预处理CSV文件:将
condition列内的换行符替换为空格,避免解析错位。 - 配置CSV分类器时,明确指定所有数值列的类型,开启
允许列为空,并设置引号字符为"以匹配示例数据格式。 - 将爬虫采样比例调整为100%,确保覆盖所有数据格式场景。
- 预处理CSV文件:将
针对方法二的问题
Catalog未生成模式的原因:
write_dynamic_frame.from_catalog仅能更新已有表的元数据,若目标表ebay_csv不存在则无法自动创建。- 自定义转换后的DataFrame存在类型推断漏洞,导致DynamicFrame元数据不完整。
修复方案:
- 先在Glue Catalog中手动创建
ebay_csv表,或改用write_dynamic_frame.from_options指定Parquet格式存储(推荐,Athena查询性能更优),并开启自动建表配置。 - 优化自定义转换逻辑:
- 替换
regexp_extract为更健壮的数值判断逻辑,避免空字符串转换失败:from pyspark.sql.functions import when for col_name in columns_to_transform: if col_name in df.columns: # 仅保留合法数值,非法值转为null后再转换类型 df = df.withColumn(col_name, when(col(col_name).rlike(r'^-?\d+(\.\d+)?$'), col(col_name)) .otherwise(None) .cast(FloatType() if col_name != 'numsold' else IntegerType()) ) - 显式定义DataFrame的Schema,确保元数据完整:
from pyspark.sql.types import StructType, StructField, StringType, FloatType, IntegerType schema = StructType([ StructField("url", StringType()), StructField("brand", StringType()), StructField("color", StringType()), StructField("type", StringType()), StructField("material", StringType()), StructField("item_height", FloatType()), StructField("item_depth", FloatType()), StructField("item_width", FloatType()), StructField("item_weight", FloatType()), StructField("condition", StringType()), StructField("price", FloatType()), StructField("seller", StringType()), StructField("numsold", IntegerType()), StructField("rating", FloatType()) ]) # 处理condition列的换行符 df = df.withColumn("condition", regexp_replace(col("condition"), "\n", " ")) df = df.select(*schema.fieldNames()) res = DynamicFrame.fromDF(df, glueContext, 'changed')
- 替换
- 确保Glue作业角色拥有S3读写、Glue Catalog操作的完整权限。
- 先在Glue Catalog中手动创建
通用建议
- 优先将CSV转换为Parquet格式存储,既提升Athena查询性能,也能让Glue更准确识别列类型。
- 处理CSV时必须保证格式规范:无内部换行符、分隔符一致、特殊字符用引号正确包裹。
内容的提问来源于stack exchange,提问作者ProgrammerNoob
相关产品推荐
相关产品推荐

