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

如何在AWS Glue作业中处理数据类型不一致的嵌套数据?

问题与解决方案建议

问题概述

用boto3抓取云数据存入S3后,数据里的Resource、Actions字段同时存在数组和字符串两种类型。想基于additional_data字段做条件查询与转换,但遇到两个问题:

  • 用Glue爬虫创建Dynamic Frame,或直接从S3读取创建Dynamic Frame时,出现Schema识别错误,且数据计数异常(显示74条,实际仅69条);
  • 用Apache Spark读取数据计数正确,但SQL转换操作繁琐,且PySpark对不一致的嵌套数据类型支持欠佳,想找更合适的处理方案。

相关代码示例

用Glue Catalog创建Dynamic Frame

S3bucket_node1 = glueContext.create_dynamic_frame.from_catalog(
    database="Raghav",
    table_name="Raghav",
    transformation_ctx="S3bucket_node1",
)

直接从S3读取创建Dynamic Frame

S3bucket_node1 = glueContext.create_dynamic_frame.from_options(
    format_options={"multiline": True},
    connection_type="s3",
    format="json",
    connection_options={
        "paths": ["s3://raghav-test-df/raghav3_removedcolon.json"],
        "recurse": True,
    },
    transformation_ctx="S3bucket_node1",
)

用Apache Spark读取数据

sparkDF=spark.read.option("inferSchema", "true").option("multiline", "true").json("s3://abababaa/abaaa.json")

解决方案建议

针对Glue Dynamic Frame的问题

  1. Schema识别错误处理
    • 手动指定Schema:提前定义包含Resource和Actions的StructType,将这两个字段设为ArrayType(StringType()),创建Dynamic Frame时传入schema参数,避免自动推断出错。
    • 使用ResolveChoice转换:对类型冲突的字段强制统一类型,比如把字符串转为单元素数组:
      from awsglue.transforms import ResolveChoice
      resolved_df = ResolveChoice.apply(frame=S3bucket_node1, choice="cast:array<string>", transformation_ctx="resolved_df")
      
  2. 数据计数异常排查
    • 过滤S3隐藏文件:在connection_options中添加excludePattern,过滤.DS_Store、临时文件等非目标数据文件。
    • 检查JSON格式:排查是否存在格式错误的文件,导致部分数据被解析忽略,可尝试单文件验证计数。

针对Spark数据转换复杂的问题

  1. 统一字段类型
    • 自定义UDF将字符串转数组:先把Resource、Actions字段统一为数组类型,再进行后续操作:
      from pyspark.sql.functions import udf
      from pyspark.sql.types import ArrayType, StringType
      
      def normalize_field(value):
          if isinstance(value, str):
              return [value]
          return value if isinstance(value, list) else []
      
      normalize_udf = udf(normalize_field, ArrayType(StringType()))
      sparkDF = sparkDF.withColumn("Resource", normalize_udf("Resource")) \
                       .withColumn("Actions", normalize_udf("Actions"))
      
  2. 用DataFrame API替代SQL
    • 统一类型后,直接用Spark DataFrame API处理嵌套数据,比如基于additional_data过滤:
      filtered_df = sparkDF.filter(sparkDF.additional_data["target_key"] == "desired_value")
      
  3. Glue与Spark结合
    • 将Glue Dynamic Frame转为Spark DataFrame(df = S3bucket_node1.toDF()),既利用Glue的S3集成优势,又能借助Spark灵活的API处理类型不一致问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 12:50:21