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>"</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
解决方案与分析
一、修复/绕过方法
强制数据集群端执行
小数据集转换为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")禁用计划缓存
异常可能与计划缓存的ProtocolBuffer序列化问题相关,通过会话参数禁用缓存:spark.conf.set("spark.sql.connect.planCache.enabled", "false")使用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")清理特殊数据类型
排查并转换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
相关产品推荐
相关产品推荐

