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

Spark测试在IntelliJ Maven面板批量运行失败,单独运行正常(winutils问题)

问题解决:Spark测试在Maven全量运行时因winutils报错,单独运行正常

问题现象

在IntelliJ的Maven面板执行test目标运行所有测试时,部分Spark测试抛出如下错误:

terminated with exception: org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Ljava/lang/String;I)Z

但以下场景测试可正常通过:

  • 在IntelliJ中单独运行整个测试类
  • 单独运行失败的测试方法
  • 通过Maven命令行指定单个测试类运行(如mvn test -Dtest=YourTestClass)

测试代码示例:

@Test
public void simpleStreamingListenerMemoryStreaming() {
    SparkContext sc = new SparkContext("local[1]", "test");
    SparkSession ss = new SparkSession(sc);
    StreamingListenerImpl listener = new StreamingListenerImpl(jobConfig);
    ss.streams().addListener(listener);
    ss.sparkContext().env().metricsSystem().registerSource(listener);

    StructField[] kafkaStructFields1 = new StructField[]{
            new StructField("key", DataTypes.StringType, true, Metadata.empty()),
            new StructField("value", DataTypes.StringType, true, Metadata.empty()),
    };

    StructType kafkaStreamSchema = new StructType(kafkaStructFields1);

    MemoryStream<Row> input = new MemoryStream<>(1, ss.sqlContext(), RowEncoder.apply(kafkaStreamSchema));
    Dataset<Row> inputDS = input.toDS();

    StreamingQuery streamingQuery = inputDS
            .writeStream()
            .format("console")
            .queryName("testMemoryStream")
            .outputMode("append")
            .start();

    List<Row> values = new ArrayList<>();
    values.add(RowFactory.create("key1", "value1"));

    input.addData(JavaConverters.asScalaIteratorConverter(values.iterator()).asScala().toSeq());

    streamingQuery.processAllAvailable();
    streamingQuery.stop();
}

即使不手动创建SparkSession,问题依然存在。

可能原因

  1. 多测试并发的资源冲突:Maven全量测试默认并行执行多个测试类,Spark上下文的初始化/销毁可能导致Hadoop native库(winutils)加载异常,或临时目录权限冲突。
  2. 环境变量差异:IntelliJ单独运行测试时可能自动注入了Hadoop相关环境变量(如HADOOP_HOME),但Maven面板运行全量测试时未正确传递。
  3. Spark上下文未彻底清理:前一个测试类的Spark上下文未完全关闭,导致后续测试的Hadoop资源占用。

解决方案

1. 强制Maven测试串行执行

在pom.xml中配置Maven Surefire插件,禁用并行测试,避免资源冲突:

<build>
    <plugins>
        <plugin>
            <groupId>org.apache.maven.plugins</groupId>
            <artifactId>maven-surefire-plugin</artifactId>
            <version>3.0.0-M7</version>
            <configuration>
                <parallel>none</parallel>
                <forkCount>1</forkCount>
                <reuseForks>false</reuseForks>
            </configuration>
        </plugin>
    </plugins>
</build>

2. 配置Hadoop环境变量到Maven测试

在Surefire插件中指定HADOOP_HOME和PATH环境变量,确保winutils可被正确加载:

<plugin>
    <groupId>org.apache.maven.plugins</groupId>
    <artifactId>maven-surefire-plugin</artifactId>
    <version>3.0.0-M7</version>
    <configuration>
        <environmentVariables>
            <HADOOP_HOME>D:\path\to\hadoop-2.8.3</HADOOP_HOME>
            <PATH>%PATH%;${HADOOP_HOME}\bin</PATH>
        </environmentVariables>
    </configuration>
</plugin>

替换D:\path\to\hadoop-2.8.3为你本地的Hadoop解压路径(需包含winutils.exe)。

3. 确保Spark上下文在测试后彻底销毁

在测试类中添加@After方法,强制关闭并清除Spark资源:

@After
public void tearDown() {
    if (ss != null) {
        ss.stop();
        SparkSession.clearActiveSession();
        SparkSession.clearDefaultSession();
    }
    if (sc != null) {
        sc.stop();
    }
}

4. 添加Hadoop native依赖到pom.xml

如果本地未配置Hadoop环境,可直接在pom中添加包含winutils的依赖(版本需与Spark兼容):

<dependency>
    <groupId>org.apache.hadoop</groupId>
    <artifactId>hadoop-common</artifactId>
    <version>2.8.3</version>
    <classifier>winutils</classifier>
</dependency>

内容的提问来源于stack exchange,提问作者maxime G

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 14:23:09