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

Spark Connect调用saveAsTable()时遇InvalidProtocolBufferException的解决问询

问题

在Databricks的Python(Spark Connect)内核中,将Pandas DataFrame转换为Spark DataFrame后,调用saveAsTable()保存为托管表时,持续触发InvalidProtocolBufferException错误。

已尝试的解决措施:

  • 显式定义Schema
  • 将NaN替换为None
  • 将datetime列转换为date类型
  • 改用write.format("delta").save(path)方法

执行代码:

spark_df_merged.write.mode("overwrite").saveAsTable("comm_prod.temp.fn_16_t")

错误栈信息:

JVM stacktrace:
grpc_shaded.com.google.protobuf.InvalidProtocolBufferException
    at grpc_shaded.com.google.protobuf.InvalidProtocolBufferException.invalidTag(InvalidProtocolBufferException.java:110)
    at grpc_shaded.com.google.protobuf.CodedInputStream$StreamDecoder.readTag(CodedInputStream.java:2084)
    at org.apache.spark.connect.proto.LocalRelation.<init>(LocalRelation.java:49)
    at org.apache.spark.connect.proto.LocalRelation.<init>(LocalRelation.java:9)
    at org.apache.spark.connect.proto.LocalRelation$1.parsePartialFrom(LocalRelation.java:660)
    at org.apache.spark.connect.proto.LocalRelation$1.parsePartialFrom(LocalRelation.java:654)
    at grpc_shaded.com.google.protobuf.AbstractParser.parsePartialFrom(AbstractParser.java:192)
    at grpc_shaded.com.google.protobuf.AbstractParser.parseFrom(AbstractParser.java:209)
    at grpc_shaded.com.google.protobuf.AbstractParser.parseFrom(AbstractParser.java:214)
    at grpc_shaded.com.google.protobuf.AbstractParser.parseFrom(AbstractParser.java:25)
    at org.apache.spark.sql.connect.planner.SparkConnectPlanner.$anonfun$transformCachedLocalRelation$1(SparkConnectPlanner.scala:1169)
    at scala.Option.map(Option.scala:230)
    at org.apache.spark.sql.connect.planner.SparkConnectPlanner.transformCachedLocalRelation(SparkConnectPlanner.scala:1163)
    at org.apache.spark.sql.connect.planner.SparkConnectPlanner.$anonfun$transformRelation$1(SparkConnectPlanner.scala:223)
    at org.apache.spark.sql.connect.service.SessionHolder.$anonfun$usePlanCache$3(SessionHolder.scala:529)
    at scala.Option.getOrElse(Option.scala:189)
    at org.apache.spark.sql.connect.service.SessionHolder.usePlanCache(SessionHolder.scala:528)
    at org.apache.spark.sql.connect.planner.SparkConnectPlanner.transformRelation(SparkConnectPlanner.scala:171)
    at org.apache.spark.sql.connect.planner.SparkConnectPlanner.transformRelation(SparkConnectPlanner.scala:158)
    at org.apache.spark.sql.connect.planner.SparkConnectPlanner.transformToDF(SparkConnectPlanner.scala:602)
    at org.apache.spark.sql.connect.planner.SparkConnectPlanner.$anonfun$transformRelation$1(SparkConnectPlanner.scala:216)
    at org.apache.spark.sql.connect.service.SessionHolder.$anonfun$usePlanCache$3(SessionHolder.scala:529)
    at scala.Option.getOrElse(Option.scala:189)
    at org.apache.spark.sql.connect.service.SessionHolder.usePlanCache(SessionHolder.scala:528)
    at org.apache.spark.sql.connect.planner.SparkConnectPlanner.transformRelation(SparkConnectPlanner.scala:171)
    at org.apache.spark.sql.connect.planner.SparkConnectPlanner.transformRelation(SparkConnectPlanner.scala:158)
    at org.apache.spark.sql.connect.planner.SparkConnectPlanner.handleWriteOperation(SparkConnectPlanner.scala:3294)
    at org.apache.spark.sql.connect.planner.SparkConnectPlanner.process(SparkConnectPlanner.scala:2857)
    at org.apache.spark.sql.connect.execution.ExecuteThreadRunner.handleCommand(ExecuteThreadRunner.scala:366)
    at org.apache.spark.sql.connect.execution.ExecuteThreadRunner.$anonfun$executeInternal$1(ExecuteThreadRunner.scala:280)
    at org.apache.spark.sql.connect.execution.ExecuteThreadRunner.$anonfun$executeInternal$1$adapted(ExecuteThreadRunner.scala:210)
    at org.apache.spark.sql.connect.service.SessionHolder.$anonfun$withSession$2(SessionHolder.scala:392)
    at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:1210)
    at org.apache.spark.sql.connect.service.SessionHolder.$anonfun$withSession$1(SessionHolder.scala:392)
    at org.apache.spark.JobArtifactSet$.withActiveJobArtifactState(JobArtifactSet.scala:97)
    at org.apache.spark.sql.artifact.ArtifactManager.$anonfun$withResources$1(ArtifactManager.scala:84)
    at org.apache.spark.util.Utils$.withContextClassLoader(Utils.scala:240)
    at org.apache.spark.sql.artifact.ArtifactManager.withResources(ArtifactManager.scala:83)
    at org.apache.spark.sql.connect.service.SessionHolder.withSession(SessionHolder.scala:391)
    at org.apache.spark.sql.connect.execution.ExecuteThreadRunner.executeInternal(ExecuteThreadRunner.scala:210)
    at org.apache.spark.sql.connect.execution.ExecuteThreadRunner.org$apache$spark$sql$connect$execution$ExecuteThreadRunner$$execute(ExecuteThreadRunner.scala:125)
    at org.apache.spark.sql.connect.execution.ExecuteThreadRunner$ExecutionThread.$anonfun$run$2(ExecuteThreadRunner.scala:592)
    at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
    at com.databricks.unity.UCSEphemeralState$Handle.runWith(UCSEphemeralState.scala:51)
    at com.databricks.unity.HandleImpl.runWith(UCSHandle.scala:104)
    at com.databricks.unity.HandleImpl.$anonfun$runWithAndClose$1(UCSHandle.scala:109)
    at scala.util.Using$.resource(Using.scala:269)
    at com.databricks.unity.HandleImpl.runWithAndClose(UCSHandle.scala:108)
    at org.apache.spark.sql.connect.execution.ExecuteThreadRunner$ExecutionThread.run(ExecuteThreadRunner.scala:592)
File <command-5558600306453652>, line 5
      1 # Write to CSV on DBFS
      2 spark_df_merged.write \
      3     .mode("overwrite") \
      4     .option("header", True) \
----> 5     .csv('abfss://structured@westusanalyticsadls.dfs.core.windows.net/temp/fn_16_t</a></span><span>'</span>)
File /databricks/spark/python/pyspark/sql/connect/client/core.py:2155, in SparkConnectClient._handle_rpc_error(self, rpc_error)
   2140                 raise Exception(
   2141                     "Python versions in the Spark Connect client and server are different. "
   2142                     "To execute user-defined functions, client and server should have the "
   (...)
   2151                     "https://docs.databricks.com/en/release-notes/serverless.html" target="_blank" rel="noopener noreferrer">https://docs.databricks.com/en/release-notes/serverless.html</a>.</span><span>&quot;</span>
   2152                 )
   2153             # END-EDGE
-> 2155             raise convert_exception(
   2156                 info,
   2157                 status.message,
   2158                 self._fetch_enriched_error(info),
   2159                 self._display_server_stack_trace(),
   2160             ) from None
   2162     raise SparkConnectGrpcException(status.message) from None

解决方案与分析

一、修复/绕过方法

  1. 强制数据集群端执行
    小数据集转换为Spark DataFrame后,Spark Connect可能将其缓存为LocalRelation引发序列化异常。添加shuffle操作强制数据传输到集群:

    spark_df_merged = spark_df_merged.repartition(1)  # 根据数据量调整分区数
    spark_df_merged.write.mode("overwrite").saveAsTable("comm_prod.temp.fn_16_t")
    
  2. 禁用计划缓存
    异常可能与计划缓存的ProtocolBuffer序列化问题相关,通过会话参数禁用缓存:

    spark.conf.set("spark.sql.connect.planCache.enabled", "false")
    
  3. 使用Databricks Pandas API直接写入
    绕过Spark Connect的序列化流程,直接用Databricks Pandas API处理:

    import databricks.pandas as dp
    dp_df = dp.from_pandas(pandas_df)
    dp_df.write.mode("overwrite").saveAsTable("comm_prod.temp.fn_16_t")
    
  4. 清理特殊数据类型
    排查并转换Pandas中的特殊类型(如category、嵌套结构)为基础类型:

    pandas_df = pandas_df.astype({"category_col": "string", "nested_col": "string"})
    spark_df_merged = spark.createDataFrame(pandas_df)
    

二、Spark Connect已知问题与推荐模式

  • 已知Bug:Spark Connect早期版本(Spark 3.3.x、Databricks Runtime 11.x及更早)中,LocalRelation的ProtocolBuffer序列化存在兼容性问题,升级到Databricks Runtime 13.x+或Spark 3.4.x+可解决大部分此类问题。
  • 推荐使用模式:
    • 大数据集优先通过Spark直接读取数据源,避免Pandas转换带来的客户端-集群传输瓶颈。
    • 必须从Pandas转换时,先写入临时存储再由Spark读取:
      # 写入临时Parquet
      pandas_df.to_parquet("/dbfs/temp/temp_data.parquet")
      # Spark读取后写入托管表
      spark.read.parquet("/dbfs/temp/temp_data.parquet").write.mode("overwrite").saveAsTable("comm_prod.temp.fn_16_t")
      
    • 小数据集场景优先使用Databricks Pandas API,稳定性优于Spark Connect。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 14:54:52