本地PySpark操作Delta文件报错:idWithoutTopologyInfo为空
PySpark操作Delta文件重启后抛出NullPointerException
问题背景
本地使用PySpark操作Delta文件,初始运行基础Spark代码无异常,通过pip安装delta-spark后成功将CSV数据写入Delta表。但重启电脑后运行相同代码时,抛出NullPointerException,核心报错信息为:Cannot invoke "org.apache.spark.storage.BlockManagerId.executorId()" because "idWithoutTopologyInfo" is null。尝试在SparkSession配置中添加.master("local[*]")后,问题仍未解决。
相关代码
基础Spark测试代码
from pyspark.sql import SparkSession from delta import configure_spark_with_delta_pip spark = SparkSession.builder.appName("Spark Application Starting").getOrCreate() print("Started Spark Context: ", spark.sparkContext) print("spark_packages :- ", spark.sparkContext.getConf().get("spark.jars.packages")) my_df = spark.sql("select 1 as id") print(my_df.show()) spark.stop()
Delta操作代码
from pyspark.sql import SparkSession from delta import configure_spark_with_delta_pip import pyspark.sql.functions as F spark_app_builder = (SparkSession.builder .appName("Employee Analysis") .config("spark.jars.packages", "io.delta:delta-core_2.12:3.2.0") .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") ) spark = configure_spark_with_delta_pip(spark_app_builder).getOrCreate() csv_df = spark.read.format("csv").option("header", True).option("inferSchema", True).load(r"..\8AM\employees.csv") csv_df = csv_df.withColumnsRenamed({"First Name": "FirstName", 'Start Date' : "StartDate", "Last Login Time" : "LastLoginTime", "Bonus %" : "Bonus_Percent", "Senior Management" : "SeniorManagement" } ).withColumn("StartDate", F.to_date(F.col("StartDate"), 'M/d/yyyy')) delta_table_path = r"..\8AM\employees.delta" print(delta_table_path) csv_df.write.mode("overwrite").format("delta").save(delta_table_path) print(" Saved to Delta File") spark.stop()
添加master配置后的代码片段
spark_app_builder = (SparkSession.builder.master("local[*]") .appName("Employee Analysis") .config("spark.jars.packages", "io.delta:delta-core_2.12:3.2.0") .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") ) spark = configure_spark_with_delta_pip(spark_app_builder).getOrCreate()
完整报错日志
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel). 24/09/01 23:00:54 ERROR Inbox: Ignoring error java.lang.NullPointerException: **Cannot invoke "org.apache.spark.storage.BlockManagerId.executorId()" because "idWithoutTopologyInfo"** is null at org.apache.spark.storage.BlockManagerMasterEndpoint.org$apache$spark$storage$BlockManagerMasterEndpoint$$register(BlockManagerMasterEndpoint.scala:677) at org.apache.spark.storage.BlockManagerMasterEndpoint$$anonfun$receiveAndReply$1.applyOrElse(BlockManagerMasterEndpoint.scala:133) at org.apache.spark.rpc.netty.Inbox.$anonfun$process$1(Inbox.scala:103) at org.apache.spark.rpc.netty.Inbox.safelyCall(Inbox.scala:213) at org.apache.spark.rpc.netty.Inbox.process(Inbox.scala:100) at org.apache.spark.rpc.netty.MessageLoop.org$apache$spark$rpc$netty$MessageLoop$$receiveLoop(MessageLoop.scala:75) at org.apache.spark.rpc.netty.MessageLoop$$anon$1.run(MessageLoop.scala:41) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:840) 24/09/01 23:00:54 WARN Executor: Issue communicating with driver in heartbeater org.apache.spark.SparkException: Exception thrown in awaitResult: at org.apache.spark.util.SparkThreadUtils$.awaitResult(SparkThreadUtils.scala:56) at org.apache.spark.util.ThreadUtils$.awaitResult(ThreadUtils.scala:310) at org.apache.spark.rpc.RpcTimeout.awaitResult(RpcTimeout.scala:75) at org.apache.spark.rpc.RpcEndpointRef.askSync(RpcEndpointRef.scala:101) at org.apache.spark.rpc.RpcEndpointRef.askSync(RpcEndpointRef.scala:85) at org.apache.spark.storage.BlockManagerMaster.registerBlockManager(BlockManagerMaster.scala:80) at org.apache.spark.storage.BlockManager.reregister(BlockManager.scala:642) at org.apache.spark.executor.Executor.reportHeartBeat(Executor.scala:1223) at org.apache.spark.executor.Executor.$anonfun$heartbeater$1(Executor.scala:295) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:1928) at org.apache.spark.Heartbeater$$anon$1.run(Heartbeater.scala:46) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.runAndReset(FutureTask.java:305) at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:305) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:840) Caused by: java.lang.NullPointerException: Cannot invoke "org.apache.spark.storage.BlockManagerId.executorId()" because "idWithoutTopologyInfo" is null
解决建议
- 版本兼容性检查:确认
delta-core_2.12:3.2.0与本地Spark版本匹配,Delta 3.2.0对应Spark 3.3.x版本,版本不匹配会引发底层通信异常。 - 清理临时文件:删除Spark本地临时目录(Windows默认在
%TEMP%\spark-*,Linux/macOS在/tmp/spark-*),避免重启后残留旧缓存文件干扰。 - 指定Driver主机地址:在SparkSession配置中添加
spark.driver.host参数固定本地IP,避免自动解析出错:spark_app_builder = (SparkSession.builder.master("local[*]") .appName("Employee Analysis") .config("spark.driver.host", "127.0.0.1") .config("spark.jars.packages", "io.delta:delta-core_2.12:3.2.0") .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") ) - 重装依赖:卸载后重新安装匹配Spark版本的delta-spark:
pip uninstall delta-spark -y pip install delta-spark==2.4.0 # 适配Spark 3.3.x,可根据实际Spark版本调整 - 检查Java环境:确认Java版本符合Spark要求(Spark 3.3.x支持Java 8/11),且
JAVA_HOME环境变量配置正确。
内容的提问来源于stack exchange,提问作者Amit Kumar Panda
相关产品推荐
相关产品推荐

