如何在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的问题
- 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")
- 手动指定Schema:提前定义包含
- 数据计数异常排查
- 过滤S3隐藏文件:在
connection_options中添加excludePattern,过滤.DS_Store、临时文件等非目标数据文件。 - 检查JSON格式:排查是否存在格式错误的文件,导致部分数据被解析忽略,可尝试单文件验证计数。
- 过滤S3隐藏文件:在
针对Spark数据转换复杂的问题
- 统一字段类型
- 自定义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"))
- 自定义UDF将字符串转数组:先把
- 用DataFrame API替代SQL
- 统一类型后,直接用Spark DataFrame API处理嵌套数据,比如基于
additional_data过滤:filtered_df = sparkDF.filter(sparkDF.additional_data["target_key"] == "desired_value")
- 统一类型后,直接用Spark DataFrame API处理嵌套数据,比如基于
- Glue与Spark结合
- 将Glue Dynamic Frame转为Spark DataFrame(
df = S3bucket_node1.toDF()),既利用Glue的S3集成优势,又能借助Spark灵活的API处理类型不一致问题。
- 将Glue Dynamic Frame转为Spark DataFrame(
内容的提问来源于stack exchange,提问作者Raghav
相关产品推荐
相关产品推荐

