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

Databricks中Spark转Ray Dataset报错,求raydb-context配置方案

Databricks中Spark转Ray Dataset报错及RayDP上下文配置问题

更新内容

我认为我终于找到了问题所在:由于Databricks初始化了Spark会话,RayDP的会话并未真正配置完成。
请问是否有办法让RayDP上下文在Spark会话中可用?

更新前问题描述

当前在Databricks上运行Spark,并仅在头节点部署了Ray。基础功能似乎正常,但尝试将Spark数据迁移至Ray Dataset时遇到报错:

TypeError                                 Traceback (most recent call last)
<command-2445755691838> in <module>
      5 memory_per_executor = "500M"
      6 # spark = raydp.init_spark(app_name, num_executors, cores_per_executor, memory_per_executor)
----> 7 dataset = ray.data.from_spark(df)

/databricks/python/lib/python3.7/site-packages/ray/data/read_api.py in from_spark(df, parallelism)
   1046     import raydp
   1047 
-> 1048     return raydp.spark.spark_dataframe_to_ray_dataset(df, parallelism)
   1049 
   1050 

/databricks/python/lib/python3.7/site-packages/raydp/spark/dataset.py in spark_dataframe_to_ray_dataset(df, parallelism, _use_owner)
    176         if parallelism != num_part:
    177             df = df.repartition(parallelism)
--> 178     blocks, _ = _save_spark_df_to_object_store(df, False, _use_owner)
    179     return from_arrow_refs(blocks)
    180 

/databricks/python/lib/python3.7/site-packages/raydp/spark/dataset.py in _save_spark_df_to_object_store(df, use_batch, _use_owner)
    150     jvm = df.sql_ctx.sparkSession.sparkContext._jvm
    151     jdf = df._jdf
--> 152     object_store_writer = jvm.org.apache.spark.sql.raydp.ObjectStoreWriter(jdf)
    153     obj_holder_name = df.sql_ctx.sparkSession.sparkContext.appName + RAYDP_OBJ_HOLDER_SUFFIX
    154     if _use_owner is True:

TypeError: 'JavaPackage' object is not callable

使用的基础代码如下:

# loading the data
# CSV options
infer_schema = "true"
first_row_is_header = "true"
delimiter = ","

df = spark.read.format("csv") \
  .option("inferSchema", infer_schema) \
  .option("header", first_row_is_header) \
  .option("sep", delimiter) \
  .load("dbfs:/databricks-datasets/nyctaxi/tripdata/green")

import ray
import sys

# disable stdout
sys.stdout.fileno = lambda: False

# connect to ray cluster on a single instance
ray.init()
ray.cluster_resources()

import raydp

dataset = ray.data.from_spark(df)

环境版本:

pyspark 3.0.1
ray 2.0.0

解决方案

核心问题分析

报错TypeError: 'JavaPackage' object is not callable的本质是:RayDP依赖的Java类(org.apache.spark.sql.raydp.ObjectStoreWriter)没有被加载到Spark的JVM classpath中。Databricks自带的Spark环境默认不包含RayDP的Java依赖,且直接使用Databricks初始化好的Spark会话时,没有通过RayDP完成会话的绑定配置。

让RayDP上下文在现有Spark会话中生效的步骤

  1. 添加RayDP的Java依赖到Spark集群

    • 方式一:集群初始化时配置(推荐)
      在Databricks集群的「高级选项」→「库」中,添加Maven坐标依赖。根据你的Spark 3.0.1(对应Scala 2.12)和Ray 2.0.0,选择兼容的RayDP版本(如0.6.0),Maven坐标为:com.raydp:raydp_2.12:0.6.0。添加后重启集群,确保依赖生效。
    • 方式二:Notebook动态添加JAR
      先将RayDP的JAR包上传到DBFS(如dbfs:/path/to/raydp_2.12-0.6.0.jar),然后在Notebook中执行:
      spark.sparkContext.addJar("dbfs:/path/to/raydp_2.12-0.6.0.jar")
      
  2. 绑定现有Spark会话到RayDP
    不要使用raydp.init_spark创建新会话,而是用raydp.attach_spark_session将Databricks已有的Spark会话绑定到RayDP上下文:

    import ray
    import raydp
    
    # 初始化Ray(确保已连接到集群)
    ray.init()
    
    # 绑定现有Spark会话
    raydp.attach_spark_session(spark)
    
    # 现在可以正常转换数据
    dataset = ray.data.from_spark(df)
    
  3. 版本兼容性验证
    确保RayDP版本与你的Spark、Ray版本匹配:

    • Spark 3.0.x + Ray 2.0.0 对应RayDP 0.5.x ~ 0.6.x版本
    • 避免使用版本差异过大的RayDP,否则会出现兼容性问题

修正后的完整代码示例

# 加载数据
infer_schema = "true"
first_row_is_header = "true"
delimiter = ","

df = spark.read.format("csv") \
  .option("inferSchema", infer_schema) \
  .option("header", first_row_is_header) \
  .option("sep", delimiter) \
  .load("dbfs:/databricks-datasets/nyctaxi/tripdata/green")

import ray
import sys
import raydp

# 禁用stdout(保留原有逻辑)
sys.stdout.fileno = lambda: False

# 连接Ray集群
ray.init()
ray.cluster_resources()

# 绑定现有Spark会话到RayDP
raydp.attach_spark_session(spark)

# 转换为Ray Dataset
dataset = ray.data.from_spark(df)

内容的提问来源于stack exchange,提问作者Bennimi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 13:31:06