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

本地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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 00:20:55