Spark 2.x过滤行为不一致求助:S3 JSON数据源过滤异常
问题:Spark Scala环境下JSON数据源过滤行为不一致
场景
- 从S3存储的JSON数据源读取数据,需过滤
cId字段:要求值长度至少为2,过滤掉null值和单字符(如"0") cId包含纯数字、字母数字混合等格式的值
问题现象
- 将
customerId转换为StringType并重命名为cId后,应用length(cId)>1过滤,返回空数据集 - 用
Seq构造的测试DataFrame(包含"12345","67890","1a2b3c4d")执行完全相同的过滤逻辑,能正常返回所有行
已尝试的操作
- 切换使用
filter()和where()方法,结果无变化 - 关闭
spark.sql.hive.convertMetastoreParquet配置 - 分析执行计划,未发现明显异常
执行计划
== Physical Plan == Format: JSON, Location: InMemoryFileIndex[s3a://s3pathredacted..., PartitionFilters: [], PushedFilters: [IsNotNull(evaluationType), IsNotNull(requestPayload), EqualTo(evaluationType,inbound_acquiring)], ReadSchema: struct<actionRecommended:string,dateInserted:string,evaluationType:string,requestPayload:struct<b... : +- BroadcastExchange HashedRelationBroadcastMode(List(input[0, string, true])) : +- *(1) Project [customerid#116 AS cId#118] : +- *(1) Filter ((length(customerid#116) > 1) && isnotnull(customerid#116)) : +- *(1) FileScan csv [customerid#116] Batched: false, Format: CSV, Location: InMemoryFileIndex[s3://folder1/folder2/folder3], PartitionFilters: [], PushedFilters: [IsNotNull(customerid)], ReadSchema: struct<customerid:string>
测试代码
val s3MLDataSource = "s3Path" var dfJson = spark.read.json(s3MLDataSource) val dfSeq= Seq("12345","67890","1a2b3c4d").toDF("customerId") // 处理JSON数据源并过滤 dfJson = dfJson.select(col("customerId").cast(StringType).as("cId")) dfJson = dfJson.where("length(cId)>1") dfJson.show(false)
执行结果
+--------------+ |cId | +--------------+ | | +--------------+
Session配置
%%configure -f {"executorMemory": "4G","driverMemory": "2G","executorCores": 8, "conf": { "spark.dynamicAllocation.enabled": "true", "spark.dynamicAllocation.shuffleTracking.enabled": "true", "spark.shuffle.service.enabled": "true", "spark.dynamicAllocation.minExecutors": "2", "spark.sql.hive.convertMetastoreParquet": "false", "spark.sql.hive.metastorePartitionPruning": "false", "spark.dynamicAllocation.maxExecutors": "20", "spark.sql.autoBroadcastJoinThreshold": "26214400", "spark.sql.cbo.enabled": "true", "spark.ui.showConsoleProgress": "false", "spark.serializer": "org.apache.spark.serializer.KryoSerializer"}}
解答
核心原因分析
问题大概率出在Spark自动推断JSON字段类型的偏差,以及后续类型转换的隐性错误:
- Spark读取JSON时,若
customerId同时存在数值型(如整数0、123)和字符串型(如"1a2b3"),会自动推断为LongType或DoubleType,而非StringType - 直接对数值型字段执行
cast(StringType),部分值可能转换异常(比如空值、特殊数值),导致后续过滤失效 - 谓词下推(PushedFilters)可能让过滤逻辑在数据源读取阶段就执行,此时字段类型还是推断的数值型,
length()函数无法正确计算数值的长度,导致误过滤
具体解决步骤
1. 手动指定Schema读取JSON
避免Spark自动推断类型出错,提前定义包含customerId为StringType的Schema:
import org.apache.spark.sql.types._ // 根据实际JSON结构补充其他字段 val customSchema = StructType(Seq( StructField("customerId", StringType, nullable = true), StructField("actionRecommended", StringType, nullable = true), StructField("dateInserted", StringType, nullable = true), StructField("evaluationType", StringType, nullable = true), // 其他字段按需添加 )) var dfJson = spark.read.schema(customSchema).json(s3MLDataSource)
2. 调整过滤逻辑,覆盖边界情况
考虑到可能存在空白字符串(而非null),修改过滤条件:
dfJson = dfJson.select(col("customerId").as("cId")) .where("cId IS NOT NULL AND length(trim(cId)) > 1")
3. 禁用谓词下推(临时排查用)
如果怀疑是谓词下推导致的问题,可关闭JSON数据源的谓词下推:
var dfJson = spark.read.option("pushDownPredicate", "false") .schema(customSchema) .json(s3MLDataSource)
4. 检查原始数据类型
先打印原始DataFrame的Schema,确认customerId的实际推断类型:
dfJson.printSchema()
若类型为数值型,可先转换为字符串并处理空值:
dfJson = dfJson.select( when(col("customerId").isNotNull, col("customerId").cast(StringType)) .otherwise(lit(null)).as("cId") ).where("length(cId) > 1")
内容的提问来源于stack exchange,提问作者Gerardsson
相关产品推荐
相关产品推荐

