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

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字段类型的偏差,以及后续类型转换的隐性错误:

  1. Spark读取JSON时,若customerId同时存在数值型(如整数0、123)和字符串型(如"1a2b3"),会自动推断为LongType或DoubleType,而非StringType
  2. 直接对数值型字段执行cast(StringType),部分值可能转换异常(比如空值、特殊数值),导致后续过滤失效
  3. 谓词下推(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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 19:35:19