Spark DataFrame过滤MongoDB导入的Struct/String混合类型异常行
解决混合类型字段的异常行过滤问题
针对你遇到的smplChnResp字段同时存在Struct和String类型、导致写入S3失败的问题,可以用以下两种可靠方法过滤异常行:
方法一:直接用Spark DataFrame判断字段类型
Spark提供typeof()函数可以直接获取字段的类型,通过类型判断过滤掉String类型的异常行:
from pyspark.sql.functions import typeof # 仅保留字段类型为struct的正常数据 filtered_df = original_df.filter(typeof("smplChnResp") == "struct")
这种方法无需依赖字段的具体结构,简单直接,不会触发类型转换异常。
方法二:通过Glue DynamicFrame过滤(适合原数据仍为DynamicFrame的场景)
如果还没转换为DataFrame,可以直接在Glue动态帧层面用自定义过滤函数,利用Python类型判断(Glue中Struct对应Python字典):
from awsglue.transforms import Filter def is_valid_struct(rec): # 判断smplChnResp是否为字典(即Struct类型) return isinstance(rec.get("smplChnResp"), dict) # 过滤后再转成DataFrame filtered_dyf = Filter.apply(frame=original_dyf, f=is_valid_struct) filtered_df = filtered_dyf.toDF()
为什么之前的方法会报错?
- 用
startswith()过滤时:因为部分行的smplChnResp是Struct类型,Spark无法对Struct执行字符串操作,因此抛出类型不匹配错误。 - 强制转换为String时:MongoDB中的Struct类型无法直接转成字符串(不是序列化后的JSON字符串),触发类型转换异常。
内容的提问来源于stack exchange,提问作者Ashish Kumar
相关产品推荐
相关产品推荐

