Python

PySpark

qianmoQqianmoQ· 更新于 2026-09-29· 阅读 9 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

PySpark](https://spark.apache.org/docs/latest/api/python/index.html) 是 Apache Spark 的 Python 接口。Kyuubi 可以在 PySpark 中作为 JDBC 数据源使用。

环境要求

PySpark 支持 Python 3.7 及以上版本。

请通过 PyPI 安装 PySpark,其中包含 Spark SQL 以及可选的 Spark 上的 pandas 支持,命令如下:

pip install pyspark 'pyspark[sql]' 'pyspark[pandas_on_spark]'

如需通过 Conda 安装或手动下载,请参阅 PySpark 安装。

准备工作

准备 JDBC 驱动

准备 JDBC 驱动 jar 文件。支持的 Hive 兼容 JDBC 驱动如下:

驱动驱动类名备注
Kyuubi Hive 驱动(文档)org.apache.kyuubi.jdbc.KyuubiHiveDriver需要在 master 分支上编译该驱动,如 KYUUBI devlive-community/knowforge#3484 中所述,Spark JDBC source 所需的功能尚未包含在已发布的版本中。
Hive 驱动(文档)org.apache.hive.jdbc.HiveDriver

请参阅各驱动的文档,准备 JDBC 驱动 jar 文件。

准备 JDBC Hive Dialect 扩展

Spark 需要支持 Hive Dialect,才能正确封装 SQL 并将其发送给 JDBC 驱动。Kyuubi 提供了一个 JDBC 方言扩展,可为 Spark 自动注册 Hive Dialect 支持。请按照 Hive Dialect 支持 中的说明准备插件 jar 文件 kyuubi-extension-spark-jdbc-dialect_-*.jar。

引入 JDBC 驱动和 Hive Dialect 扩展的 jar 包

选择以下任一方式将 jar 文件引入 Spark。

  • 将 JDBC 驱动和 Hive Dialect 的 jar 文件放入 $SPARK_HOME/jars 目录,使其对 PySpark 的 classpath 可见。并将 spark.sql.extensions = org.apache.spark.sql.dialect.KyuubiSparkJdbcDialectExtension 添加到 $SPARK_HOME/conf/spark_defaults.conf.
  • 使用 spark 的启动 shell,在提交应用时通过 --packages 引入 JDBC 驱动,通过 --jars 引入 Hive Dialect 插件
$SPARK_HOME/bin/pyspark --py-files PY_FILES \
--packages org.apache.hive:hive-jdbc:x.y.z \
--jars /path/kyuubi-extension-spark-jdbc-dialect_-*.jar
  • 通过 SparkSession 构建器设置 jar 和配置
from pyspark.sql import SparkSession

spark = SparkSession.builder \
        .config("spark.jars", "/path/hive-jdbc-x.y.z.jar,/path/kyuubi-extension-spark-jdbc-dialect_-*.jar") \
        .config("spark.sql.extensions", "org.apache.spark.sql.dialect.KyuubiSparkJdbcDialectExtension") \
        .getOrCreate()

用法

有关 PySpark JDBC 的用法和选项的更多信息,请参阅 Spark 的 JDBC To Other Databases。

以编程方式用作 JDBC 数据源

# Loading data from Kyuubi via HiveDriver as JDBC datasource
jdbcDF = spark.read \
  .format("jdbc") \
  .options(driver="org.apache.hive.jdbc.HiveDriver",
           url="jdbc:hive2://kyuubi_server_ip:port",
           user="user",
           password="password",
           query="select * from testdb.src_table"
           ) \
  .load()

通过 SQL 用作 JDBC 数据源表

从 Spark 3.2.0 起,CREATE DATASOURCE TABLE 支持通过 SQL 创建 jdbc 数据源。

# create JDBC Datasource table with DDL
spark.sql("""CREATE TABLE kyuubi_table USING JDBC
OPTIONS (
    driver='org.apache.hive.jdbc.HiveDriver',
    url='jdbc:hive2://kyuubi_server_ip:port',
    user='user',
    password='password',
    dbtable='testdb.some_table'
)""")

# read data to dataframe
jdbcDF = spark.sql("SELECT * FROM kyuubi_table")

# write data from dataframe in overwrite mode
df.writeTo("kyuubi_table").overwrite

# write data from query
spark.sql("INSERT INTO kyuubi_table SELECT * FROM some_table")

使用 PySpark 与 Pandas

从 PySpark 3.2.0 起,PySpark 支持 pandas API on Spark,可让你将 pandas 工作负载扩展到分布式环境。

Pandas-on-Spark DataFrame 与 Spark DataFrame 实际上可以互换使用。更多说明请参阅 From/to pandas and PySpark DataFrames。

import pyspark.pandas as ps

psdf = ps.range(10)
sdf = psdf.to_spark().filter("id > 5")
sdf.show()

评论

登录后参与评论

正在加载评论…