Spring Boot Spark部署Minikube遇SerializedLambda赋值异常求助
Spark on Minikube:Executor序列化异常求助
我在Minikube Kubernetes集群上运行一个简单的Spring Boot Spark应用,本地用SparkSession.builder().master("local")运行正常,但部署到Minikube后,Driver Pod能成功启动Executor Pod,Executor Pod却抛出如下异常:
ERROR Executor: Exception in task 0.1 in stage 0.0 (TID 1) cannot assign instance of java.lang.invoke.SerializedLambda to field org.apache.spark.sql.execution.MapPartitionsExec.func of type scala.Function1 in instance of org.apache.spark.sql.execution.MapPartitionsExec
WordcountController代码
所有逻辑集中在控制器中:
@RestController public class WordCountController implements Serializable { @PostMapping("/wordcount") public ResponseEntity<String> handleFileUpload(@RequestParam("file") MultipartFile file) throws IOException { String hostIp; try { hostIp = InetAddress.getLocalHost().getHostAddress(); } catch (UnknownHostException e) { throw new RuntimeException(e); } SparkConf conf = new SparkConf(); conf.setAppName("count.words.in.file") .setMaster("k8s://https://kubernetes.default.svc:443") .setJars(new String[]{"/app/wordcount.jar"}) .set("spark.driver.host", hostIp) .set("spark.driver.port", "8080") .set("spark.kubernetes.namespace", "default") .set("spark.kubernetes.container.image", "spark:3.3.2h.1") .set("spark.executor.cores", "2") .set("spark.executor.memory", "1g") .set("spark.kubernetes.authenticate.executor.serviceAccountName", "spark") .set("spark.kubernetes.dynamicAllocation.deleteGracePeriod", "20") .set("spark.cores.max", "4") .set("spark.executor.instances", "2"); SparkSession spark = SparkSession.builder() .config(conf) .getOrCreate(); byte[] byteArray = file.getBytes(); String contents = new String(byteArray, StandardCharsets.UTF_8); Dataset<String> text = spark.createDataset(Arrays.asList(contents), Encoders.STRING()); Dataset<String> wordsDataset = text.flatMap((FlatMapFunction<String, String>) line -> { List<String> words = new ArrayList<>(); for (String word : line.split(" ")) { words.add(word); } return words.iterator(); }, Encoders.STRING()); // Count the number of occurrences of each word Dataset<Row> wordCounts = wordsDataset.groupBy("value") .agg(count("*").as("count")) .orderBy(desc("count")); // Convert the word count results to a List of Rows List<Row> wordCountsList = wordCounts.collectAsList(); StringBuilder resultStringBuffer = new StringBuilder(); // Build the final string representation for (Row row : wordCountsList) { resultStringBuffer.append(row.getString(0)).append(": ").append(row.getLong(1)).append("\n"); } return ResponseEntity.ok(resultStringBuffer.toString()); } }
Maven依赖配置(pom.xml)
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.7.8</version> <relativePath/> <!-- lookup parent from repository --> </parent> <groupId>com.example</groupId> <artifactId>wordcount</artifactId> <version>0.0.1-SNAPSHOT</version> <name>wordcount</name> <description>wordcount</description> <properties> <java.version>11</java.version> <spark.version>3.3.2</spark.version> <scala.version>2.12</scala.version> </properties> <dependencyManagement> <dependencies> <!--Spark java.lang.NoClassDefFoundError: org/codehaus/janino/InternalCompilerException--> <dependency> <groupId>org.codehaus.janino</groupId> <artifactId>commons-compiler</artifactId> <version>3.0.8</version> </dependency> <dependency> <groupId>org.codehaus.janino</groupId> <artifactId>janino</artifactId> <version>3.0.8</version> </dependency> </dependencies> </dependencyManagement> <dependencies> <dependency> <groupId>org.codehaus.janino</groupId> <artifactId>commons-compiler</artifactId> <version>3.0.8</version> </dependency> <dependency> <groupId>org.codehaus.janino</groupId> <artifactId>janino</artifactId> <version>3.0.8</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_${scala.version}</artifactId> <version>${spark.version}</version> </dependency> <dependency> <!-- Spark dependency --> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_${scala.version}</artifactId> <version>${spark.version}</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-kubernetes_${scala.version}</artifactId> <version>${spark.version}</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <scope>test</scope> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-maven-plugin</artifactId> </plugin> </plugins> </build> </project>
Driver镜像Dockerfile
用于打包Spring Boot应用并部署到Minikube:
# Use an existing image as the base image FROM openjdk:11-jdk # Set the working directory WORKDIR /app # Copy the compiled JAR file to the image COPY target/wordcount-0.0.1-SNAPSHOT.jar /app/wordcount.jar RUN useradd -u 185 sparkuser # Set the entrypoint command to run the JAR file ENTRYPOINT ["java", "-jar", "wordcount.jar"]
Executor镜像说明
spark.kubernetes.container.image使用的是本地spark-3.3.2-bin-hadoop3安装包附带的Dockerfile构建的镜像,版本与Spring Boot应用依赖的Spark一致,已加载到Minikube。
已尝试的解决方案(均无效)
- 指定
setJars(new String[]{"/app/wordcount.jar"})共享应用JAR包(路径为Driver镜像中JAR的绝对路径) - 使用maven-shade-plugin修改依赖分发方式,导致Driver Pod出现
ClassNotFoundException: SparkSession异常 - 重构代码,将Lambda表达式替换为实现
FlatMapFunction的静态内部类,无效果:public static class SplitLine implements FlatMapFunction<String, String> { @Override public Iterator<String> call(String line) throws Exception { List<String> words = new ArrayList<>(); for (String word : line.split(" ")) { words.add(word); } return words.iterator(); } ... Dataset<String> wordsDataset = text.flatMap(new SplitLine(), Encoders.STRING());
恳请各位提供配置建议或代码修改指导,以解决当前环境下的序列化异常问题。
内容的提问来源于stack exchange,提问作者itaydafna
相关产品推荐
相关产品推荐

