Quarkus应用创建SparkSession遇请求挂起问题排查求助
Quarkus集成Apache Spark时REST接口无限挂起问题
问题现象
开发集成Apache Spark的Quarkus应用,计划部署到Kubernetes集群。Quarkus开发模式启动无报错,但调用使用SparkSession的REST接口时请求无限挂起,无任何报错信息:
- 当
ExampleResource的hello方法返回sparkSessionBean.toString()时运行正常 - 当返回
sparkSessionBean.getSpark().toString()时,请求永久挂起 - 相同代码在独立Java主应用中可正常创建SparkSession并执行任务
环境信息
- Quarkus版本:3.5.3/3.8.1
- Spark版本:3.5.0
- Java版本:17
代码示例
ExampleResource.java
package com.example; import jakarta.inject.Inject; import jakarta.ws.rs.GET; import jakarta.ws.rs.Path; import jakarta.ws.rs.Produces; import jakarta.ws.rs.core.MediaType; @Path("/hello") public class ExampleResource { @Inject SparkSessionBean sparkSessionBean; @GET @Produces(MediaType.TEXT_PLAIN) public String hello() { return sparkSessionBean.getSpark().toString(); } }
SparkSessionBean.java
package com.example; import jakarta.enterprise.context.ApplicationScoped; import org.apache.spark.sql.SparkSession; @ApplicationScoped public class SparkSessionBean { private static final String WAREHOUSE_PATH = "s3a://iceberg"; private SparkSession spark; public SparkSession getSpark() { if (spark == null) { spark = SparkSession.builder() .appName("Java Spark Iceberg Example") .master("local") .config("spark.ui.enabled", "false") //Filesystem (MinIO) config .config("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") .config("fs.s3a.endpoint", "http://localhost:9000") .config("fs.s3a.access.key", "key123") .config("fs.s3a.secret.key", "secret123") .config("fs.s3a.path.style.access", "true") .config("spark.sql.warehouse.dir", WAREHOUSE_PATH) //Nessie catalog config .config("spark.jars.packages", "org.apache.iceberg:iceberg-spark-runtime-3.5_2.13:1.3.0,org.projectnessie.nessie-integrations:nessie-spark-extensions-3.5_2.13:0.77.1") .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions,org.projectnessie.spark.extensions.NessieSparkSessionExtensions") .config("spark.sql.catalog.nessie", "org.apache.iceberg.spark.SparkCatalog") .config("spark.sql.catalog.nessie.catalog-impl", "org.apache.iceberg.nessie.NessieCatalog") .config("spark.sql.catalog.nessie.authentication.type", "NONE") .config("spark.sql.catalog.nessie.uri", "http://localhost:19120/api/v1") .config("spark.sql.catalog.nessie.ref", "main") .config("spark.sql.defaultCatalog", "nessie") .config("spark.sql.catalog.nessie.warehouse", WAREHOUSE_PATH) .getOrCreate(); } return spark; } }
pom.xml
<?xml version="1.0"?> <project xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd" xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"> <modelVersion>4.0.0</modelVersion> <groupId>com.example</groupId> <artifactId>quarkus-spark</artifactId> <version>1.0-SNAPSHOT</version> <properties> <compiler-plugin.version>3.12.1</compiler-plugin.version> <maven.compiler.release>17</maven.compiler.release> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding> <quarkus.platform.artifact-id>quarkus-bom</quarkus.platform.artifact-id> <quarkus.platform.group-id>io.quarkus.platform</quarkus.platform.group-id> <quarkus.platform.version>3.8.1</quarkus.platform.version> <skipITs>true</skipITs> <surefire-plugin.version>3.2.5</surefire-plugin.version> <spark.version>3.5.0</spark.version> <iceberg.version>1.4.3</iceberg.version> <hadoop.version>3.3.4</hadoop.version> </properties> <dependencyManagement> <dependencies> <dependency> <groupId>${quarkus.platform.group-id}</groupId> <artifactId>${quarkus.platform.artifact-id}</artifactId> <version>${quarkus.platform.version}</version> <type>pom</type> <scope>import</scope> </dependency> </dependencies> </dependencyManagement> <dependencies> <!-- QUARKUS --> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-arc</artifactId> </dependency> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-resteasy-reactive</artifactId> </dependency> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-junit5</artifactId> <scope>test</scope> </dependency> <dependency> <groupId>io.rest-assured</groupId> <artifactId>rest-assured</artifactId> <scope>test</scope> </dependency> <!-- SPARK --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.13</artifactId> <version>${spark.version}</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.13</artifactId> <version>${spark.version}</version> </dependency> <!-- ICEBERG --> <dependency> <groupId>org.apache.iceberg</groupId> <artifactId>iceberg-spark-runtime-3.5_2.13</artifactId> <version>${iceberg.version}</version> </dependency> <!-- HADOOP --> <dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-aws</artifactId> <version>${hadoop.version}</version> </dependency> <!-- OTHER --> <dependency> <groupId>org.antlr</groupId> <artifactId>antlr4-runtime</artifactId> <version>4.9.3</version> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>connect-api</artifactId> <version>3.2.3</version> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>${quarkus.platform.group-id}</groupId> <artifactId>quarkus-maven-plugin</artifactId> <version>${quarkus.platform.version}</version> <extensions>true</extensions> <executions> <execution> <goals> <goal>build</goal> <goal>generate-code</goal> <goal>generate-code-tests</goal> </goals> </execution> </executions> </plugin> <plugin> <artifactId>maven-compiler-plugin</artifactId> <version>${compiler-plugin.version}</version> <configuration> <compilerArgs> <arg>-parameters</arg> </compilerArgs> </configuration> </plugin> <plugin> <artifactId>maven-surefire-plugin</artifactId> <version>${surefire-plugin.version}</version> <configuration> <systemPropertyVariables> <java.util.logging.manager>org.jboss.logmanager.LogManager</java.util.logging.manager> <maven.home>${maven.home}</maven.home> </systemPropertyVariables> </configuration> </plugin> <plugin> <artifactId>maven-failsafe-plugin</artifactId> <version>${surefire-plugin.version}</version> <executions> <execution> <goals> <goal>integration-test</goal> <goal>verify</goal> </goals> </execution> </executions> <configuration> <systemPropertyVariables> <native.image.path>${project.build.directory}/${project.build.finalName}-runner </native.image.path> <java.util.logging.manager>org.jboss.logmanager.LogManager</java.util.logging.manager> <maven.home>${maven.home}</maven.home> </systemPropertyVariables> </configuration> </plugin> </plugins> </build> <profiles> <profile> <id>native</id> <activation> <property> <name>native</name> </property> </activation> <properties> <skipITs>false</skipITs> <quarkus.package.type>native</quarkus.package.type> </properties> </profile> </profiles> </project>
问题分析与解决方案
核心原因
Quarkus与Spark的集成确实存在潜在的线程模型和类加载器冲突问题,主要原因包括:
- 线程阻塞:Quarkus的Reactive REST(Resteasy Reactive)使用事件循环线程,Spark初始化过程会阻塞当前线程,导致事件循环卡住,请求无法响应。
- 类加载器隔离:Quarkus的类加载器模型与Spark的类加载机制不兼容,Spark初始化时加载大量依赖可能触发类加载死锁或资源竞争。
- 动态依赖下载阻塞:
spark.jars.packages配置会让Spark动态下载依赖,在Quarkus同步请求线程中执行时,容易引发线程挂起。
解决步骤
切换到阻塞式REST端点
替换Reactive依赖为传统阻塞式JAX-RS实现,避免事件循环线程被阻塞:<!-- 移除Reactive依赖 --> <!-- <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-resteasy-reactive</artifactId> </dependency> --> <!-- 添加阻塞式依赖 --> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-resteasy</artifactId> </dependency>异步初始化SparkSession
在Bean初始化阶段用单独线程异步创建SparkSession,避免阻塞请求线程:@ApplicationScoped public class SparkSessionBean { private static final String WAREHOUSE_PATH = "s3a://iceberg"; private CompletableFuture<SparkSession> sparkFuture; @PostConstruct void init() { sparkFuture = CompletableFuture.supplyAsync(() -> SparkSession.builder() .appName("Java Spark Iceberg Example") .master("local") .config("spark.ui.enabled", "false") // 保留其他原有配置 .getOrCreate() ); } public SparkSession getSpark() throws ExecutionException, InterruptedException { return sparkFuture.get(); } }调整Spark依赖配置
移除spark.jars.packages,将依赖直接添加到pom.xml中,避免动态下载阻塞:<dependency> <groupId>org.apache.iceberg</groupId> <artifactId>iceberg-spark-runtime-3.5_2.13</artifactId> <version>1.3.0</version> </dependency> <dependency> <groupId>org.projectnessie.nessie-integrations</groupId> <artifactId>nessie-spark-extensions-3.5_2.13</artifactId> <version>0.77.1</version> </dependency>类加载器兼容配置
在application.properties中添加配置,让Spark相关类使用Quarkus系统类加载器:quarkus.class-loader.parent-first-artifacts=org.apache.spark,org.apache.hadoop,org.apache.iceberg,org.projectnessie
Kubernetes部署注意事项
- 移除
master("local")配置,改为连接集群中的Spark服务,示例:spark.kubernetes.master=https://kubernetes.default.svc - 确保Pod分配足够的CPU和内存资源运行Spark客户端
- 优先使用JVM模式部署,Spark部分依赖在Quarkus原生镜像中可能存在兼容性问题
内容的提问来源于stack exchange,提问作者PowerfullDeveloper
相关产品推荐
相关产品推荐

