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会话中生效的步骤
添加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")
- 方式一:集群初始化时配置(推荐)
绑定现有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)版本兼容性验证
确保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
相关产品推荐
相关产品推荐

