Sqoop与Spark从RDBMS接入数据时过滤逻辑执行位置对比问询
Spark JDBC对接RDBMS的谓词下推执行机制说明
核心结论:Spark执行RDBMS查询时,过滤条件是否下推到数据库侧执行,由代码写法和优化器规则共同决定,两种场景都存在,不能一概而论。
过滤下推到RDBMS执行的场景
只要符合以下任意一种写法,过滤操作都会下发到RDBMS侧执行,仅拉取符合条件的记录,和Sqoop的执行逻辑完全一致:
- 直接将带过滤条件的完整SQL作为子查询传入JDBC读取参数,示例写法:
这种场景下Spark会直接把整个子查询交给RDBMS执行,拉取的结果就是已经过滤好的数据集。val query = "(select id,name from customer where country='US') as tmp" val df = spark.read.jdbc(jdbcUrl, query, connectionProps) - 使用DataFrame DSL操作时,过滤条件为标准SQL可识别的简单表达式,示例写法:
Spark JDBC数据源默认开启谓词下推优化,会自动将// 先读取全表结构 val df = spark.read.jdbc(jdbcUrl, "customer", connectionProps) // 执行列裁剪+过滤 val resDf = df.select("id", "name").filter("country = 'US'")select id,name from customer where country='US'拼接为下发查询语句,交给RDBMS执行。
全量拉取后Spark内存过滤的场景
部分技术文章提到的拉全量再过滤,只会出现在过滤条件无法被RDBMS识别的场景,常见情况包括:
- 过滤条件中使用了Spark自定义函数(UDF)、Spark独有的函数语法,比如小写转换的DSL写法:
这类语法无法被RDBMS解析,Spark只能先拉取全量df.filter(lower(col("country")) === "us")customer表数据到内存,再执行过滤运算。 - 手动关闭了JDBC谓词下推配置(默认
spark.sql.jdbc.pushDownPredicate为true,关闭后所有过滤都不会下推)。
下推效果验证方法
你可以通过执行resDf.explain()查看物理执行计划,只要输出的JDBC Relation部分的查询语句携带了where country='US'条件,就说明下推成功,否则为内存过滤。
内容的提问来源于stack exchange,提问作者user3625945
相关产品推荐
相关产品推荐

