如何让Spark在过滤前将MySQL数据加载至内存?
解决Spark JDBC谓词下推问题的两种方法
因为Spark的优化器会自动将过滤操作下推至数据源(这里是MySQL),导致无索引的date列过滤拖慢数据库。要让Spark先加载数据再过滤,可通过以下两种方式实现:
方法1:强制缓存数据到Spark内存
在读取完JDBC数据后,调用cache()或persist()方法,触发数据加载到Spark内存,之后再执行过滤操作。这样Spark会先完成全量分区数据的读取,再在内存中做过滤。
修改后的代码:
with SparkSession.builder.appName("Spark App").getOrCreate() as spark: dataframe_mysql = spark.read.format('jdbc').options( url="jdbc:mysql://.../...", driver='com.mysql.cj.jdbc.Driver', dbtable='my_table', user=..., password=..., partitionColumn='id', lowerBound=0, upperBound=10000000, numPartitions=11, fetchsize=1000000, isolationLevel='NONE' ).load() # 强制缓存数据到内存,触发数据读取 dataframe_mysql.cache() # 执行action操作触发缓存执行 dataframe_mysql.count() # 在Spark内存中执行过滤 dataframe_mysql = dataframe_mysql.filter("date > '2022-01-01'") dataframe_mysql.write.parquet('...')
方法2:关闭JDBC谓词下推功能
在JDBC读取选项中添加pushDownPredicate='false',直接禁用Spark将过滤条件下推到数据库的功能,这样数据库只会收到基于id分区的查询,过滤操作由Spark在加载数据后执行。
修改后的代码:
with SparkSession.builder.appName("Spark App").getOrCreate() as spark: dataframe_mysql = spark.read.format('jdbc').options( url="jdbc:mysql://.../...", driver='com.mysql.cj.jdbc.Driver', dbtable='my_table', user=..., password=..., partitionColumn='id', lowerBound=0, upperBound=10000000, numPartitions=11, fetchsize=1000000, isolationLevel='NONE', # 禁用谓词下推 pushDownPredicate='false' ).load() dataframe_mysql = dataframe_mysql.filter("date > '2022-01-01'") dataframe_mysql.write.parquet('...')
注意事项
- 方法1需确保Spark集群有足够内存缓存全量分区数据,避免内存不足引发磁盘溢出。
- 方法2更简洁直接,无需额外缓存操作,适合内存资源有限但需规避数据库端低效过滤的场景。
内容的提问来源于stack exchange,提问作者Rinze
相关产品推荐
相关产品推荐

