拓冰建站拓冰建站
首页 / 资讯中心 / 正文

从 Koalas 迁移到 pandas API on Spark:PySpark 迁移指南与 API 变化详解

从 Koalas 迁移到 pandas API on SparkPySpark 迁移指南与 API 变化详解【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark导读本文对应 Apache Spark 仓库中 python/docs/source/migration_guide/koalas_to_pyspark.rst 官方迁移指南系统梳理从独立的databricks.koalasKoalas包迁移到 PySpark 内置的 pandas API on Sparkpyspark.pandas所需的全部改动包括包导入路径、DataFrame/Series 访问器命名、PySpark DataFrame 转换方法重命名以及版本号获取方式的变更。读完本文你将能够快速定位并替换旧 Koalas 代码中的每一处 API并理解这些变更背后的源码实现与版本演进依据。背景Koalas 与 pandas API on Spark 的关系Koalas 是 Databricks 早期发起的开源项目目标是在 Apache Spark 之上提供近似 pandas 的 DataFrame API。随着该能力被合入 Apache Spark 主干即 pandas API on Spark模块名为pyspark.pandasKoalas 作为独立发行包的使命结束官方迁移指南明确要求用户从databricks.koalas切换到 PySpark 内置模块。这一迁移涉及以下五类核心变更变更点Koalas旧pandas API on Spark新包导入import databricks.koalas as ksimport pyspark.pandas as psDataFrame 访问器DataFrame.koalasDataFrame.pandas_on_sparkPySpark DataFrame 转换DataFrame.to_koalasDataFrame.pandas_apiPySpark DataFrame 转换曾用名DataFrame.to_pandas_on_sparkDataFrame.pandas_api版本号databricks.koalas.__version__pyspark.__version__其中DataFrame.koalas、DataFrame.to_koalas、DataFrame.to_pandas_on_spark三个旧接口已在 Spark 4.0 中移除。一、包导入路径变更databricks.koalas→pyspark.pandas迁移的第一步是修改模块导入语句。原指南给出的最小改动如下# import databricks.koalas as ks import pyspark.pandas as ps只需将databricks.koalas替换为pyspark.pandas并同步将代码中所有ks前缀如ks.DataFrame、ks.Series、ks.read_csv改为ps前缀。pandas API on Spark 的顶层命名空间由 python/pyspark/pandas/namespace.py 提供其 API 形状与 Koalas 保持一致因此绝大多数业务代码可以在不改动逻辑的前提下完成迁移。二、访问器重命名DataFrame.koalas→DataFrame.pandas_on_spark在 Koalas 中用户可以通过DataFrame.koalas访问 pandas-on-Spark 专属的扩展方法在 pandas API on Spark 中这一访问器被重命名为DataFrame.pandas_on_sparkSeries 对应Series.pandas_on_spark旧访问器DataFrame.koalas已在 Spark 4.0 中移除。从当前仓库源码可以确认新访问器的注册方式。在 python/pyspark/pandas/frame.py 中pandas_on_spark CachedAccessor(pandas_on_spark, PandasOnSparkFrameMethods)对应地python/pyspark/pandas/series.py 为 Series 注册了同样的访问器pandas_on_spark CachedAccessor(pandas_on_spark, PandasOnSparkSeriesMethods)通过pandas_on_spark访问的典型方法访问器对应的实现类位于 python/pyspark/pandas/accessors.pyPandasOnSparkFrameMethods与PandasOnSparkSeriesMethods这些方法以 pandas chunk 为单位在 Spark 上执行原生 pandas 逻辑是迁移后最常用的扩展能力attach_id_column(id_type, column)为 DataFrame 附加一列作为行标识id_type支持sequence、distributed-sequence与distributed。注意sequence使用无分区条件的 Spark Window会把数据集中到单个分区官方注释明确提示在大数据集上避免使用。apply_batch(func, *args, **kwargs)以每个 pandas chunk 作为输入调用自定义函数逐块应用可用于行/列级操作例如对pdf.query(A 1)这类 pandas 查询逻辑逐批执行。transform_batch(func, *args, **kwargs)以每个 pandas chunk 为输入做逐块转换例如lambda c: c 1这类逐 Series 处理逻辑。一个典型的迁移示例原 Koalas 写法 → 新写法import pyspark.pandas as ps pdf ps.DataFrame({A: [1, 2, 3], B: [4, 5, 6]}) # DataFrame.koalas.apply_batch → DataFrame.pandas_on_spark.apply_batch result pdf.pandas_on_spark.apply_batch(lambda pdf: pdf.query(A 1))三、PySpark DataFrame 转换方法to_koalas/to_pandas_on_spark→pandas_api在 Koalas 时代PySpark 的DataFrame通过 monkey-patch 注入了DataFrame.to_koalas方法用于将 Spark DataFrame 转换为 pandas-on-Spark DataFrame随后一度引入中间名DataFrame.to_pandas_on_spark。这两个旧方法现在统一更名为DataFrame.pandas_api并已在 Spark 4.0 中移除旧名称。源码中的实现依据当前仓库中pandas_api是一个通过dispatch_df_method分发到各执行端Classic / Connect的方法其公共 API 定义在 python/pyspark/sql/dataframe.pydispatch_df_method def pandas_api( self, index_col: Optional[Union[str, List[str]]] None ) - PandasOnSparkDataFrame: Converts the existing DataFrame into a pandas-on-Spark DataFrame. .. versionadded:: 3.2.0 .. versionchanged:: 3.5.0 Supports Spark Connect. ... 分发到各执行端的实现分别位于 python/pyspark/sql/classic/dataframe.py经典执行模式与 python/pyspark/sql/connect/dataframe.pySpark Connect 客户端。使用方法与关键语义转换的基本用法from pyspark.sql import SparkSession spark SparkSession.builder.getOrCreate() sdf spark.createDataFrame( [(14, Tom), (23, Alice), (16, Bob)], [age, name]) # 旧写法sdf.to_koalas() 或 sdf.to_pandas_on_spark() psdf sdf.pandas_api() print(psdf) # age name # 0 14 Tom # 1 23 Alice # 2 16 Bob # 可指定索引列 psdf_with_index sdf.pandas_api(index_colage) print(psdf_with_index) # name # age # 14 Tom # 23 Alice # 16 Bob使用pandas_api时需要特别注意索引语义官方 docstring 明确说明如果 pandas-on-Spark DataFrame 先被转换为 Spark DataFrame、再转换回 pandas-on-Spark索引信息会丢失原始索引会退化为普通列。因此跨类型转换时应显式通过index_col指定期望保留为索引的列。此外pandas_api仅在环境中安装并可用 Pandas 时才生效。仓库中的对应测试位于 python/pyspark/sql/tests/test_dataframe.py验证了默认转换与index_colCol1两种调用路径。Spark Connect 兼容性pandas_api的版本注释显示该方法自 Spark 3.2.0 引入自 3.5.0 起支持 Spark Connect。也就是说在 Spark Connect 模式下客户端一侧的 DataFrame 同样可以调用pandas_api进入 pandas-on-Spark 工作流这覆盖了现代部署形态下的迁移需求。四、版本号获取方式databricks.koalas.__version__→pyspark.__version__旧代码中通过databricks.koalas.__version__获取 Koalas 版本号的做法已被移除。迁移后应统一使用 PySpark 自身的版本属性import pyspark print(pyspark.__version__)该属性定义于 python/pyspark/version.py__version__: str 5.0.0.dev0注意pyspark.__version__反映的是当前 PySpark 发行版或开发版的版本号而不是某个独立 pandas API 模块的版本——因为 pandas API on Spark 已随 PySpark 一起发布不再存在独立的 Koalas 版本线。例如在开发版分支上该值会是类似5.0.0.dev0的 SNAPSHOT 版本字符串。五、迁移检查清单与相关文档完成迁移后建议按以下清单逐项核对代码库导入语句全文搜索databricks.koalas替换为pyspark.pandas同时将ks.前缀替换为ps.。访问器调用搜索\.koalas\b或DataFrame.koalas替换为pandas_on_spark注意同时检查 Series 场景。转换方法搜索to_koalas、to_pandas_on_spark统一替换为pandas_api并根据需要补充index_col参数以保留索引语义。版本号搜索koalas.__version__替换为pyspark.__version__。依赖与运行前提确认环境已安装 Pandaspandas_api与部分访问器方法依赖它并确认所用 Spark 版本满足接口引入要求pandas_api需 Spark 3.2.0Spark Connect 场景需 3.5.0。如需了解更多迁移上下文可继续阅读仓库内相关文档python/docs/source/migration_guide/index.rst迁移指南总览指明从 Koalas 迁移到 PySpark 的入口python/docs/source/migration_guide/pyspark_upgrade.rstSpark 4.0 升级说明明确记录了DataFrame.koalas、DataFrame.to_koalas的移除及替代 APIpython/pyspark/pandas/accessors.pypandas_on_spark访问器的完整方法实现attach_id_column、apply_batch、transform_batch等python/pyspark/sql/dataframe.pypandas_api的公共 API 定义与示例。总结从 Koalas 迁移到 pandas API on Spark 本质上是一次命名空间的收敛功能本身已随 PySpark 内置而无需重新安装依赖需要修改的只是包名、访问器名、转换方法名与版本号获取方式四处调用点。结合本文给出的源码路径你可以逐一验证迁移后的每个调用是否指向新接口并在 Spark 4.0 及后续版本中安全运行。【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门