Apache Spark Structured Streaming流查询启动报错求助
Spark Structured Streaming 写入JSON文件报错问题分析与解决
问题背景
使用Apache Spark 3.2.1(基于Scala 2.12)开发Structured Streaming应用,从Kafka主题读取数据转换后写入JSON文件时遇到两个核心问题:
- 采用
append输出模式启动流查询时,触发Py4JJavaError,错误栈指向目录创建/权限相关操作失败 - 将输出模式改为
complete时,触发AnalysisException,提示Data source json does not support Complete output mode
相关代码片段
from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, col, expr from pyspark.sql.types import StructType, StructField, StringType, LongType, DoubleType, IntegerType, ArrayType if __name__ == "__main__": spark = SparkSession \ .builder \ .appName("File Streaming Demo") \ .master("local[3]") \ .config("spark.streaming.stopGracefullyOnShutdown", "true") \ .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.1") \ .getOrCreate() schema = StructType([ StructField("taxiId", StringType()), StructField("location", StructType([ StructField("lat", StringType()), StructField("lng", StringType()), ])), ]) kafka_df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "localhost:9092") \ .option("subscribe", "map_log_topic") \ .option("startingOffsets", "earliest") \ .load() value_df = kafka_df.select(from_json(col("value").cast("string"), schema).alias("value")) explode_df = value_df.selectExpr("value.taxiId", "value.location.lat", "value.location.lng") invoice_writer_query = explode_df.writeStream \ .format("json") \ .queryName("Map logs writer") \ .outputMode("append") \ .option("path", "output") \ .option("checkpointLocation", "chk-point-dir") \ .trigger(processingTime="1 minute") \ .start() invoice_writer_query.awaitTermination()
错误信息
1. append模式下的Py4JJavaError
py4j.protocol.Py4JJavaError: An error occurred while calling o57.start. : ExitCodeException exitCode=-1073741515: at org.apache.hadoop.util.Shell.runCommand(Shell.java:1007) at org.apache.hadoop.util.Shell.run(Shell.java:900) at org.apache.hadoop.util.Shell$ShellCommandExecutor.execute(Shell.java:1212) at org.apache.hadoop.util.Shell.execCommand(Shell.java:1306) at org.apache.hadoop.util.Shell.execCommand(Shell.java:1288) at org.apache.hadoop.fs.RawLocalFileSystem.setPermission(RawLocalFileSystem.java:978) at org.apache.hadoop.fs.RawLocalFileSystem.mkOneDirWithMode(RawLocalFileSystem.java:660) at org.apache.hadoop.fs.RawLocalFileSystem.mkdirsWithOptionalPermission(RawLocalFileSystem.java:700) at org.apache.hadoop.fs.RawLocalFileSystem.mkdirs(RawLocalFileSystem.java:672) at org.apache.hadoop.fs.RawLocalFileSystem.mkdirsWithOptionalPermission(RawLocalFileSystem.java:699) at org.apache.hadoop.fs.RawLocalFileSystem.mkdirs(RawLocalFileSystem.java:677) at org.apache.hadoop.fs.FileSystem.primitiveMkdir(FileSystem.java:1356) at org.apache.hadoop.fs.DelegateToFileSystem.mkdir(DelegateToFileSystem.java:185) at org.apache.hadoop.fs.FilterFs.mkdir(FilterFs.java:219) at org.apache.hadoop.fs.FileContext$4.next(FileContext.java:809) at org.apache.hadoop.fs.FileContext$4.next(FileContext.java:805) at org.apache.hadoop.fs.FSLinkResolver.resolve(FSLinkResolver.java:90) at org.apache.hadoop.fs.FileContext.mkdir(FileContext.java:812) at org.apache.spark.sql.execution.streaming.AbstractFileContextBasedCheckpointFileManager.mkdirs(CheckpointFileManager.scala:319) at org.apache.spark.sql.execution.streaming.HDFSMetadataLog.<init>(HDFSMetadataLog.scala:67) at org.apache.spark.sql.execution.streaming.CompactibleFileStreamLog.<init>(CompactibleFileStreamLog.scala:48) at org.apache.spark.sql.execution.streaming.FileStreamSinkLog.<init>(FileStreamSinkLog.scala:91) at org.apache.spark.sql.execution.streaming.FileStreamSink.<init>(FileStreamSink.scala:139) at org.apache.spark.sql.execution.datasources.DataSource.createSink(DataSource.scala:322) at org.apache.spark.sql.streaming.DataStreamWriter.createV1Sink(DataStreamWriter.scala:442) at org.apache.spark.sql.streaming.DataStreamWriter.startInternal(DataStreamWriter.scala:407) at org.apache.spark.sql.streaming.DataStreamWriter.start(DataStreamWriter.scala:251) at java.base/jdk.internal.reflect.DirectMethodHandleAccessor.invoke(DirectMethodHandleAccessor.java:103) at java.base/java.lang.reflect.Method.invoke(Method.java:580) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:374) 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.ClientServerConnection.waitForCommands(ClientServerConnection.java:182) at py4j.ClientServerConnection.run(ClientServerConnection.java:106) at java.base/java.lang.Thread.run(Thread.java:1570)
2. complete模式下的AnalysisException
Data source json does not support Complete output mode
问题原因分析
Py4JJavaError(目录/权限报错)
- 错误码
-1073741515对应Windows系统的STATUS_DLL_NOT_FOUND,核心原因是Hadoop Windows原生依赖库缺失:Spark依赖Hadoop的文件系统API管理输出目录和检查点目录,未配置对应DLL文件会导致无法执行创建目录、设置权限等操作 - 若运行环境是Linux/macOS,则可能是运行Spark的用户对
output或chk-point-dir目录无读写权限
- 错误码
JSON数据源不支持Complete输出模式
- Spark的文件类输出源(JSON、CSV、Parquet等)仅支持
append和update模式 complete模式要求每次触发时输出全量数据集,但文件输出是追加式写入逻辑,无法覆盖或生成全量文件,因此天然不支持该模式
- Spark的文件类输出源(JSON、CSV、Parquet等)仅支持
解决建议
针对目录/权限错误的解决方案
Windows环境
- 下载与Spark依赖的Hadoop版本匹配的Windows二进制包(Spark 3.2.1默认依赖Hadoop 3.2.x),解压后将
bin目录下的hadoop.dll、winutils.exe等文件所在路径添加到系统PATH环境变量 - 提前手动创建
output和chk-point-dir目录,确保当前用户拥有读写权限 - 可选:在SparkSession构建时添加配置禁用权限检查
.config("spark.hadoop.fs.permissions.umask-mode", "000")
Linux/macOS环境
- 执行命令为
output和chk-point-dir目录赋予读写权限:chmod -R 755 ./output ./chk-point-dir - 确保运行Spark的用户是目录的所有者,或拥有对应权限
针对输出模式错误的解决方案
- 保持使用
append输出模式,这是文件输出源的标准且唯一适配的增量输出模式 - 若业务需要全量输出,可采用以下方案:
- 定期通过批处理任务将增量JSON文件合并为全量文件
- 改用支持
complete模式的输出源(如内存表、JDBC数据库等)
内容的提问来源于stack exchange,提问作者uday mahajan
相关产品推荐
相关产品推荐

