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

如何使用TestContainers创建Apache Spark独立集群用于集成测试?求运行示例

我来帮你搞定这个问题!用TestContainers搭建Apache Spark独立集群做集成测试确实不算太常见的场景,不过只要理清步骤,其实很容易实现。下面我会一步步讲解配置方法,还附上可直接运行的示例代码。

一、先准备依赖

不管你用Maven还是Gradle,都需要添加TestContainers核心依赖,以及Spark的客户端依赖。

Maven 依赖片段

<dependencies>
    <!-- TestContainers 核心 -->
    <dependency>
        <groupId>org.testcontainers</groupId>
        <artifactId>testcontainers</artifactId>
        <version>1.19.7</version>
        <scope>test</scope>
    </dependency>
    <!-- TestContainers Docker Compose 支持 -->
    <dependency>
        <groupId>org.testcontainers</groupId>
        <artifactId>docker-compose</artifactId>
        <version>1.19.7</version>
        <scope>test</scope>
    </dependency>
    <!-- Spark 客户端依赖 -->
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-core_2.12</artifactId>
        <version>3.5.0</version>
        <scope>test</scope>
    </dependency>
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-sql_2.12</artifactId>
        <version>3.5.0</version>
        <scope>test</scope>
    </dependency>
</dependencies>

Gradle 依赖片段

testImplementation 'org.testcontainers:testcontainers:1.19.7'
testImplementation 'org.testcontainers:docker-compose:1.19.7'
testImplementation 'org.apache.spark:spark-core_2.12:3.5.0'
testImplementation 'org.apache.spark:spark-sql_2.12:3.5.0'
二、用Docker Compose定义Spark集群

TestContainers支持Docker Compose,这种方式比手动创建多个容器更简洁。先在src/test/resources下创建docker-compose-spark.yml:

version: '3.8'
services:
  spark-master:
    image: bitnami/spark:3.5.0
    command: bin/spark-class org.apache.spark.deploy.master.Master
    ports:
      - "7077:7077" # Spark master 通信端口
      - "8080:8080" # UI端口(可选,用于调试)
    environment:
      - SPARK_MODE=master
      - SPARK_RPC_AUTHENTICATION_ENABLED=no
      - SPARK_RPC_ENCRYPTION_ENABLED=no
      - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no
      - SPARK_SSL_ENABLED=no

  spark-worker:
    image: bitnami/spark:3.5.0
    command: bin/spark-class org.apache.spark.deploy.worker.Worker spark://spark-master:7077
    depends_on:
      - spark-master
    environment:
      - SPARK_MODE=worker
      - SPARK_MASTER_URL=spark://spark-master:7077
      - SPARK_WORKER_MEMORY=1G
      - SPARK_WORKER_CORES=1
      - SPARK_RPC_AUTHENTICATION_ENABLED=no
      - SPARK_RPC_ENCRYPTION_ENABLED=no
      - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no
      - SPARK_SSL_ENABLED=no

这里选用Bitnami的Spark镜像,它预配置了启动脚本和环境变量,比官方Apache镜像更适合容器化场景。

三、编写集成测试代码

下面是一个JUnit 5的测试类,它会自动启动Spark集群容器,提交简单任务验证集群可用性:

import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.testcontainers.containers.DockerComposeContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;

import java.io.File;

import static org.junit.jupiter.api.Assertions.assertEquals;

@Testcontainers
public class SparkClusterIntegrationTest {

    // 绑定Docker Compose配置文件
    @Container
    private static final DockerComposeContainer<?> sparkCluster = new DockerComposeContainer<>(new File("src/test/resources/docker-compose-spark.yml"))
            .withExposedService("spark-master", 7077)
            .waitingFor("spark-master", org.testcontainers.containers.wait.strategy.Wait.forHttp("/").forPort(8080)) // 等待Master UI就绪
            .waitingFor("spark-worker", org.testcontainers.containers.wait.strategy.Wait.forLogMessage(".*Successfully registered with master.*", 1)); // 等待Worker注册完成

    private static SparkSession sparkSession;

    @BeforeAll
    static void setUp() {
        // 获取本地映射后的Spark Master地址
        String masterUrl = String.format("spark://%s:%d",
                sparkCluster.getServiceHost("spark-master", 7077),
                sparkCluster.getServicePort("spark-master", 7077));

        // 初始化SparkSession,关键是配置driver地址让集群能回调
        sparkSession = SparkSession.builder()
                .appName("SparkClusterTest")
                .master(masterUrl)
                .config("spark.driver.host", getLocalHostAddress())
                .getOrCreate();
    }

    @Test
    void testSimpleSparkJobOnCluster() {
        // 创建测试数据集
        Dataset<Row> df = sparkSession.createDataFrame(
                sparkSession.sparkContext().parallelize(java.util.Arrays.asList(
                        new Person("Alice", 30),
                        new Person("Bob", 25),
                        new Person("Charlie", 35)
                )), Person.class);

        // 执行聚合验证
        long totalCount = df.count();
        assertEquals(3, totalCount);

        Dataset<Row> avgAgeResult = df.groupBy().avg("age");
        double avgAge = avgAgeResult.first().getDouble(0);
        assertEquals(30.0, avgAge, 0.01);
    }

    @AfterAll
    static void tearDown() {
        if (sparkSession != null) {
            sparkSession.stop();
        }
    }

    // 获取本地主机IP,确保集群能访问Driver
    private static String getLocalHostAddress() {
        try {
            return java.net.InetAddress.getLocalHost().getHostAddress();
        } catch (java.net.UnknownHostException e) {
            throw new RuntimeException("Failed to get local host address", e);
        }
    }

    // Spark需要的实体类(带无参构造和Getter)
    public static class Person {
        private String name;
        private int age;

        public Person() {}

        public Person(String name, int age) {
            this.name = name;
            this.age = age;
        }

        public String getName() { return name; }
        public int getAge() { return age; }
    }
}
四、关键注意事项
  • Driver地址配置:必须设置spark.driver.host为本地IP,因为Spark Worker需要连接回Driver传输任务结果,TestContainers容器在内部网络,需要能访问到你的本地机器。
  • 资源调整:如果测试任务较大,记得在docker-compose.yml里给Worker增加SPARK_WORKER_MEMORY和SPARK_WORKER_CORES,避免内存溢出。
  • 等待策略:一定要配置等待条件,确保Master和Worker完全就绪后再提交任务,否则会出现连接失败的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 09:01:21