使用Synapse Spark向Azure Event Hub写入数据时遇Py4JJavaError错误
Synapse PySpark写入Azure Event Hub触发Py4JJavaError问题解决
问题现象
在Synapse Analytics Studio中使用PySpark操作时,可正常读取Azure Event Hub的消息,但执行DataFrame写入操作时,调用save()方法触发Py4JJavaError,核心错误为:
java.lang.NoSuchMethodError: org.apache.spark.sql.AnalysisException.<init>(Ljava/lang/String;Lscala/Option;Lscala/Option;Lscala/Option;Lscala/Option;)V
正常读取Event Hub的代码
import json connectionString = "Endpoint=sb://::hidden::" ehConf = { } ehConf['eventhubs.connectionString'] = sc._jvm.org.apache.spark.eventhubs.EventHubsUtils.encrypt(connectionString) # Create the positions startingEventPosition = { "offset": -1, "seqNo": -1, #not in use "enqueuedTime": None, #not in use "isInclusive": True } ehConf["eventhubs.startingPosition"] = json.dumps(startingEventPosition) df = spark.read.format("eventhubs").options(**ehConf).load() display(df)
写入Event Hub的报错代码
df1 = spark.read.parquet(silver_path) # confirmed to have data via display(df1) df1 \ .select(struct(*[c for c in df1.columns]).alias("body")) \ .write \ .format("eventhubs") \ .options(**ehConf) \ .save()
完整报错信息
Py4JJavaError: An error occurred while calling o4089.save. : java.lang.NoSuchMethodError: org.apache.spark.sql.AnalysisException.<init>(Ljava/lang/String;Lscala/Option;Lscala/Option;Lscala/Option;Lscala/Option;)V at org.apache.spark.sql.eventhubs.EventHubsWriter$.validateQuery(EventHubsWriter.scala:58) at org.apache.spark.sql.eventhubs.EventHubsWriter$.write(EventHubsWriter.scala:70) at org.apache.spark.sql.eventhubs.EventHubsSourceProvider.createRelation(EventHubsSourceProvider.scala:124) at org.apache.spark.sql.execution.datasources.SaveIntoDataSourceCommand.run(SaveIntoDataSourceCommand.scala:47) at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult$lzycompute(commands.scala:75) at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult(commands.scala:73) at org.apache.spark.sql.execution.command.ExecutedCommandExec.executeCollect(commands.scala:84) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.$anonfun$applyOrElse$1(QueryExecution.scala:108) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:111) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:183) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:97) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:779) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:66) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:108) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:104) at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDownWithPruning$1(TreeNode.scala:584) at org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(TreeNode.scala:176) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDownWithPruning(TreeNode.scala:584) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.org$apache$spark$sql$catalyst$plans$logical$AnalysisHelper$$super$transformDownWithPruning(LogicalPlan.scala:31) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning(AnalysisHelper.scala:267) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning$(AnalysisHelper.scala:263) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:31) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:31) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:560) at org.apache.spark.sql.execution.QueryExecution.eagerlyExecuteCommands(QueryExecution.scala:104) at org.apache.spark.sql.execution.QueryExecution.commandExecuted$lzycompute(QueryExecution.scala:88) at org.apache.spark.sql.execution.QueryExecution.commandExecuted(QueryExecution.scala:82) at org.apache.spark.sql.execution.QueryExecution.assertCommandExecuted(QueryExecution.scala:136) at org.apache.spark.sql.DataFrameWriter.runCommand(DataFrameWriter.scala:901) at org.apache.spark.sql.DataFrameWriter.saveToV1Source(DataFrameWriter.scala:415) at org.apache.spark.sql.DataFrameWriter.saveInternal(DataFrameWriter.scala:382) at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:249) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) at py4j.Gateway.invoke(Gateway.java:282) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.GatewayConnection.run(GatewayConnection.java:238) at java.lang.Thread.run(Thread.java:750)
问题原因及解决方案
原因
错误根源是Spark运行时与Event Hub连接器版本不兼容。Synapse Analytics已预安装适配其Spark版本的Event Hub连接器,手动通过pip install azure-eventhub安装的包会与内置版本冲突,引发方法缺失的异常。同时,写入的body字段类型不符合要求——Event Hub要求body为二进制类型,而非struct类型。
解决步骤
- 移除手动安装的包:删除代码中的
### %pip install azure-eventhub语句,使用Synapse内置的Event Hub连接器。 - 修正DataFrame的body字段类型:将struct转为JSON字符串,再编码为二进制类型,符合Event Hub的写入要求。修改后的写入代码如下:
from pyspark.sql.functions import to_json, struct, col, encode df1 = spark.read.parquet(silver_path) # 将所有字段打包为struct,转为JSON字符串后编码为UTF-8二进制 df1 \ .select(encode(to_json(struct(*[col(c) for c in df1.columns])), "UTF-8").alias("body")) \ .write \ .format("eventhubs") \ .options(**ehConf) \ .mode("append") # 指定写入模式,避免重复或覆盖问题 .save()
- 验证权限与配置:确保Event Hub的连接字符串拥有Send权限,且
ehConf中的配置正确无误。
内容的提问来源于stack exchange,提问作者Rob Koch
相关产品推荐
相关产品推荐

