如何使用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
相关产品推荐
相关产品推荐

