如何通过ODBC连接使用PySpark/Spark创建DataFrame(不使用pandas)
PySpark/Spark通过ODBC创建DataFrame实现方案(无pandas依赖)
前置准备
- 准备对应版本的JDBC-ODBC桥接驱动包,Spark对接ODBC数据源通过JDBC接口桥接实现,全程不依赖pandas库,所有数据读取、转换均在Spark引擎侧完成,直接生成原生Spark分布式DataFrame。
- 提前整理好ODBC连接所需参数:已配置完成的ODBC数据源DSN名称、数据源登录账号、密码、目标读取的表名/自定义查询逻辑。
PySpark版本实现代码
from pyspark.sql import SparkSession # 初始化SparkSession,加载JDBC-ODBC桥接驱动 spark = SparkSession.builder \ .appName("ODBC_Create_DataFrame") \ .config("spark.jars", "/your/local/path/jdbc-odbc-bridge-driver.jar") \ .getOrCreate() # 配置连接参数 conn_props = { "user": "your_data_source_username", "password": "your_data_source_password", "driver": "对应桥接驱动的类名" # 随引入的驱动包调整,部分第三方驱动类名与内置驱动不同 } # ODBC对应的JDBC连接串,直接关联已配置的ODBC DSN odbc_url = "jdbc:odbc:your_pre_configured_dsn_name" # 方式1:直接读取整张表生成DataFrame df_full_table = spark.read.jdbc( url=odbc_url, table="target_db.target_table", properties=conn_props ) # 方式2:通过自定义SQL预筛选后生成DataFrame,减少无效数据拉取 custom_sql = "(SELECT col1, col2, col3 FROM target_db.target_table WHERE dt >= '2024-01-01') as tmp_view" df_filtered = spark.read.jdbc( url=odbc_url, table=custom_sql, properties=conn_props ) # 验证DataFrame可用性 df_filtered.printSchema() df_filtered.show(5, truncate=False)
Scala Spark版本实现代码
import org.apache.spark.sql.SparkSession object ODBCDataFrameDemo { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("ODBC_Create_DataFrame") .config("spark.jars", "/your/local/path/jdbc-odbc-bridge-driver.jar") .getOrCreate() import java.util.Properties val connProps = new Properties() connProps.put("user", "your_data_source_username") connProps.put("password", "your_data_source_password") connProps.put("driver", "对应桥接驱动的类名") val odbcUrl = "jdbc:odbc:your_pre_configured_dsn_name" // 读全表 val dfFullTable = spark.read.jdbc(odbcUrl, "target_db.target_table", connProps) // 按自定义SQL读取 val customSql = "(SELECT col1, col2, col3 FROM target_db.target_table WHERE dt >= '2024-01-01') as tmp_view" val dfFiltered = spark.read.jdbc(odbcUrl, customSql, connProps) dfFiltered.printSchema() dfFiltered.show(5, truncate = false) spark.stop() } }
注意事项
- 上述实现全程未引入pandas依赖,生成的DataFrame为Spark原生分布式数据集,不存在单机内存瓶颈,适配TB级数据读取场景。
- 高版本JDK已移除内置的JDBC-ODBC桥接驱动,自行引入第三方维护的对应桥接包即可,核心读取逻辑无需调整。
- 读取大表时建议在
spark.read.jdbc中追加partitionColumn、lowerBound、upperBound、numPartitions参数,配置多并发分区读取,避免单线程拉取全量数据造成源库压力过大、读取速度过慢的问题。 - 若ODBC数据源需要特殊鉴权(如Kerberos、SSL双向认证),直接在连接参数中追加对应鉴权配置项即可,不需要修改核心读取逻辑。
内容的提问来源于stack exchange,提问作者Sonali Bisht
相关产品推荐
相关产品推荐

