【问题标题】:Apache Spark Session: IOException: mkdir of (path) failedApache Spark 会话:IOException:(路径)的 mkdir 失败
【发布时间】:2017-02-07 18:41:52
【问题描述】:

我正在测试 Apache Spark 2.0 的新版本,尝试利用结构化流功能,使用非常简单的代码创建包含流数据的数据集,然后打印创建的数据集。 这是我的代码:

    SparkSession mySession= SparkSession.builder().appName("ProcessData").master("local[*]").config("spark.sql.warehouse.dir","System.getProperty(\"user.dir\")/warehouse").getOrCreate();
    Dataset<Row> measurements=mySession.readStream().format("socket").option("host","localhost").option("port",5556).load();
    StreamingQuery printDataset=measurements.writeStream().format("console").start();
    printDataset.awaitTermination();

问题是我得到一个 IOException: mkdir of (temporary directory) failed。 有人可以帮我解决这个问题吗?非常感谢。

这是显示的完整错误:

Exception in thread "main" java.io.IOException: mkdir of C:/Users/Manuel%20Mourato/AppData/Local/Temp/temporary-891579db-0442-4e1c-8642-d41c7885ab26/offsets failed
at org.apache.hadoop.fs.FileSystem.primitiveMkdir(FileSystem.java:1065)
at org.apache.hadoop.fs.DelegateToFileSystem.mkdir(DelegateToFileSystem.java:176)
at org.apache.hadoop.fs.FilterFs.mkdir(FilterFs.java:197)
at org.apache.hadoop.fs.FileContext$4.next(FileContext.java:730)
at org.apache.hadoop.fs.FileContext$4.next(FileContext.java:726)
at org.apache.hadoop.fs.FSLinkResolver.resolve(FSLinkResolver.java:90)
at org.apache.hadoop.fs.FileContext.mkdir(FileContext.java:733)
at org.apache.spark.sql.execution.streaming.HDFSMetadataLog$FileContextManager.mkdirs(HDFSMetadataLog.scala:281)
at org.apache.spark.sql.execution.streaming.HDFSMetadataLog.<init>(HDFSMetadataLog.scala:57)
at org.apache.spark.sql.execution.streaming.StreamExecution.<init>(StreamExecution.scala:131)
at org.apache.spark.sql.streaming.StreamingQueryManager.startQuery(StreamingQueryManager.scala:251)
at org.apache.spark.sql.streaming.DataStreamWriter.start(DataStreamWriter.scala:287)
at org.apache.spark.sql.streaming.DataStreamWriter.start(DataStreamWriter.scala:231)

【问题讨论】:

    标签: java apache-spark apache-spark-sql spark-structured-streaming


    【解决方案1】:

    可以这样试试吗?

    SparkSession mySession= SparkSession.builder().appName("ProcessData").master("local[*]").config("spark.sql.warehouse.dir",System.getProperty("user.dir") + "/warehouse").getOrCreate();
    

    【讨论】:

    • 我尝试了你的建议,虽然它使代码更“干净”,但它并没有解决我的问题。还是谢谢。
    【解决方案2】:

    为什么在字符串中使用 System.getProperty? 另外,请检查该文件夹是否存在,即:

    val tempDir = System.getProperty("user.dir");
    val path = tempDir + "/warehouse";
    
    SparkSession mySession= SparkSession.builder().appName("ProcessData").master("local[*]").config("spark.sql.warehouse.dir", path).getOrCreate();
    

    另外请检查您是否对该路径有写权限。如果您手动创建仓库目录并设置权限应该很好 - 您将确保一切正常

    编辑:知道了! 首先,您应该检查 AppData/Local/Temp 的写入权限,因为它是标准的临时目录。

    此错误是由OffsetLog 引起的。您可以更改将在其中创建日志的目录by addingoption("checkpointLocation", ...)

    【讨论】:

    • 我已经检查过了,确实存在路径并且我有权在其中写入。但是,我应该提一下:这个 IOException 不是试图在 val path 中创建一个目录,而是试图在另一个文件夹中创建一个目录(.../AppData/Local/Temp)。这正常吗?
    • 不,不是。我将在 Spark 代码中搜索它使用临时目录的位置
    • 非常感谢。我在上面包含了完整的错误,它可能会有所帮助。
    • 谢谢,但它不起作用:/。例如,我按照您所说的添加了option("checkpointLocation", path_to_desktop),并且我选择的任何新路径都会出现相同的错误,所以IOException: mkdir of path_to_desktop/offsets failed 太奇怪了。
    • @manuelmourato 这很奇怪。另一种选择-您可以将此位置设置为名称中没有空格的文件夹吗?也许空间正在制造一些问题
    【解决方案3】:

    确保您还在配置中设置了自己的检查点目录,该目录具有写入权限,一种方法是在相同的应用程序代码中创建检查点目录,例如,

    .config("spark.sql.streaming.checkpointLocation", "C:\\sparkApp\\checkpoints\\")
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2023-01-21
      • 1970-01-01
      • 1970-01-01
      • 2011-07-07
      • 2023-03-04
      • 1970-01-01
      • 2012-11-28
      • 2012-09-24
      相关资源
      最近更新 更多