如何在PySpark(Databricks)中按条件读取SQL Server表数据
PySpark从SQL Server按条件读取数据到DataFrame
当然可以实现按条件读取,这里有两种常用方式,优先推荐第一种,因为能把过滤逻辑下推到数据库端,避免全表加载:
方法1:将过滤逻辑下推到SQL Server(推荐)
直接在dbtable参数中传入带条件的子查询,这样过滤操作会在SQL Server端完成,只返回符合条件的数据,性能更优,完全对应你给出的SQL语句:
remote_table = (spark.read.format("sqlserver") .option("host", "host_name") .option("user", "user_name") # 原代码中"use_name"应为笔误,修正为"user_name" .option("password", "password") .option("database", "database_name") # 子查询必须添加别名,否则JDBC会报错 .option("dbtable", "(select * from dbo.table_name where time_stamp = cast(getdate() as date)) as filtered_table") .load() )
这种方式的核心是让SQL Server先执行过滤,再把结果返回给Spark,适合大数据量场景,能大幅减少数据传输量。
方法2:PySpark客户端过滤(仅适合小数据量)
如果先读取整张表,再用PySpark的filter()方法过滤,虽然能实现需求,但会先加载全表数据到Spark集群,效率较低:
# 先读取整张表 full_table = (spark.read.format("sqlserver") .option("host", "host_name") .option("user", "user_name") .option("password", "password") .option("database", "database_name") .option("dbtable", "dbo.table_name") .load() ) # 导入PySpark日期函数,过滤当前日期的数据 from pyspark.sql.functions import current_date, col remote_table = full_table.filter(col("time_stamp") == current_date())
注意:current_date()获取的是Spark集群节点的当前日期,如果SQL Server与Spark时区不一致,需要配置时区参数(比如spark.sql.session.timeZone)确保日期匹配。
内容的提问来源于stack exchange,提问作者Abhishek Singh
相关产品推荐
相关产品推荐

