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

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

问题原因分析

  1. Py4JJavaError(目录/权限报错)

    • 错误码-1073741515对应Windows系统的STATUS_DLL_NOT_FOUND,核心原因是Hadoop Windows原生依赖库缺失:Spark依赖Hadoop的文件系统API管理输出目录和检查点目录,未配置对应DLL文件会导致无法执行创建目录、设置权限等操作
    • 若运行环境是Linux/macOS,则可能是运行Spark的用户对output或chk-point-dir目录无读写权限
  2. JSON数据源不支持Complete输出模式

    • Spark的文件类输出源(JSON、CSV、Parquet等)仅支持append和update模式
    • complete模式要求每次触发时输出全量数据集,但文件输出是追加式写入逻辑,无法覆盖或生成全量文件,因此天然不支持该模式

解决建议

针对目录/权限错误的解决方案

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 04:59:52