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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 04:51:22