You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将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写入节点四个核心环节。


问题排查与解决方案

针对方法一的问题

  1. 数值列异常原因:

    • 示例数据中condition列包含换行符,会导致CSV解析错位,后续数值列被错误读取为字符串;-1这类特殊值也可能被分类器识别为非数值类型。
    • 爬虫默认采样比例低,若采样行包含格式异常数据,会错误推断列类型。
  2. 修复方案:

    • 预处理CSV文件:将condition列内的换行符替换为空格,避免解析错位。
    • 配置CSV分类器时,明确指定所有数值列的类型,开启允许列为空,并设置引号字符为"以匹配示例数据格式。
    • 将爬虫采样比例调整为100%,确保覆盖所有数据格式场景。

针对方法二的问题

  1. Catalog未生成模式的原因:

    • write_dynamic_frame.from_catalog仅能更新已有表的元数据,若目标表ebay_csv不存在则无法自动创建。
    • 自定义转换后的DataFrame存在类型推断漏洞,导致DynamicFrame元数据不完整。
  2. 修复方案:

    • 先在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操作的完整权限。

通用建议

  • 优先将CSV转换为Parquet格式存储,既提升Athena查询性能,也能让Glue更准确识别列类型。
  • 处理CSV时必须保证格式规范:无内部换行符、分隔符一致、特殊字符用引号正确包裹。

内容的提问来源于stack exchange,提问作者ProgrammerNoob

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 13:02:04