Spark+Kafka服务无报错静默退出问题排查求助
问题:Spark+Kafka服务无报错静默退出排查
我正在构建一个微服务架构的后端应用,其中一个服务是Apache Spark任务运行器——通过Kafka接收带日期范围的消息,提取对应日期范围内CSV文件的数据。为简化测试,我把逻辑调整为计算CSV文件中所有记录的平均价格(忽略未使用的request变量),但进程出现无报错的静默退出,日志无错误提示,需要排查原因。
服务代码
package com.samurai.lab.sparkjobservice; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; import java.io.IOException; import java.nio.file.Files; import java.nio.file.Paths; import java.time.Duration; import java.util.Collections; import java.util.Properties; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; public class App2 { private static final ObjectMapper mapper = new ObjectMapper(); private static final String KAFKA_BROKERS = "localhost:9092"; private static final String KAFKA_GROUP = "sparkGroup"; private static final String KAFKA_TOPIC = "marketDataRequests"; private static final String MARKETDATA_CSV_PATH = "C:/Users/USER/Desktop/TickData/data-12-01-until-12-03-year-2020/fake.csv"; public static void main(String[] args) { ExecutorService executor = Executors.newSingleThreadExecutor(); executor.submit(App2::consume); } private static void consume() { try (KafkaConsumer<String, String> consumer = createKafkaConsumer()) { consumer.subscribe(Collections.singletonList(KAFKA_TOPIC)); while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); for (ConsumerRecord<String, String> record : records) { processRecord(record); } } } catch (Exception e) { e.printStackTrace(); } } private static KafkaConsumer<String, String> createKafkaConsumer() { Properties props = new Properties(); props.put("bootstrap.servers", KAFKA_BROKERS); props.put("group.id", KAFKA_GROUP); props.put("key.deserializer", StringDeserializer.class.getName()); props.put("value.deserializer", StringDeserializer.class.getName()); return new KafkaConsumer<>(props); } private static void processRecord(ConsumerRecord<String, String> record) { try { String jsonMessage = record.value(); JsonNode rootNode = mapper.readTree(jsonMessage); processMarketDataRequest(rootNode); } catch (IOException e) { System.out.println("Error while processing record: " + e.getMessage()); e.printStackTrace(); } } private static void processMarketDataRequest(JsonNode request) { // Create a Spark session SparkSession spark = SparkSession.builder() .appName("ClientModeExample") .master("local[*]") .getOrCreate(); // Check if file exists if (!Files.exists(Paths.get(MARKETDATA_CSV_PATH))) { System.err.println("File does not exist: " + MARKETDATA_CSV_PATH); return; } // Read CSV into a DataFrame Dataset<Row> df = spark.read().option("header", "true").csv(MARKETDATA_CSV_PATH); // Convert the Price column to double Dataset<Row> dfDouble = df.withColumn("Price", df.col("Price").cast("double")); // Calculate the average price Dataset<Row> avgPrice = dfDouble.agg(org.apache.spark.sql.functions.avg("Price").alias("average_price")); // Print the average price avgPrice.show(); // Stop the Spark session spark.stop(); } }
控制台日志
SLF4J: Class path contains multiple SLF4J providers. SLF4J: Found provider [org.apache.logging.slf4j.SLF4JServiceProvider@620b29e6] SLF4J: Found provider [org.slf4j.simple.SimpleServiceProvider@121ff335] SLF4J: See https://www.slf4j.org/codes.html#multiple_bindings for an explanation. SLF4J: Actual provider is of type [org.apache.logging.slf4j.SLF4JServiceProvider@620b29e6] Using Spark's default log4j profile: org/apache/spark/log4j2-defaults.properties 23/07/21 13:00:36 INFO SparkContext: Running Spark version 3.4.1 23/07/21 13:00:37 INFO ResourceUtils: ============================================================== 23/07/21 13:00:37 INFO ResourceUtils: No custom resources configured for spark.driver. 23/07/21 13:00:37 INFO ResourceUtils: ============================================================== 23/07/21 13:00:37 INFO SparkContext: Submitted application: ClientModeExample 23/07/21 13:00:37 INFO ResourceProfile: Default ResourceProfile created, executor resources: Map(cores -> name: cores, amount: 1, script: , vendor: , memory -> name: memory, amount: 1024, script: , vendor: , offHeap -> name: offHeap, amount: 0, script: , vendor: ), task resources: Map(cpus -> name: cpus, amount: 1.0) 23/07/21 13:00:37 INFO ResourceProfile: Limiting resource is cpu 23/07/21 13:00:37 INFO ResourceProfileManager: Added ResourceProfile id: 0 23/07/21 13:00:37 INFO SecurityManager: Changing view acls to: USER 23/07/21 13:00:37 INFO SecurityManager: Changing modify acls to: USER 23/07/21 13:00:37 INFO SecurityManager: Changing view acls groups to: 23/07/21 13:00:37 INFO SecurityManager: Changing modify acls groups to: 23/07/21 13:00:37 INFO SecurityManager: SecurityManager: authentication disabled; ui acls disabled; users with view permissions: USER; groups with view permissions: EMPTY; users with modify permissions: USER; groups with modify permissions: EMPTY 23/07/21 13:00:37 INFO Utils: Successfully started service 'sparkDriver' on port 5396. 23/07/21 13:00:37 INFO SparkEnv: Registering MapOutputTracker 23/07/21 13:00:37 INFO SparkEnv: Registering BlockManagerMaster 23/07/21 13:00:37 INFO BlockManagerMasterEndpoint: Using org.apache.spark.storage.DefaultTopologyMapper for getting topology information 23/07/21 13:00:37 INFO BlockManagerMasterEndpoint: BlockManagerMasterEndpoint up 23/07/21 13:00:37 INFO ConsumerCoordinator: [Consumer clientId=consumer-sparkGroup-1, groupId=sparkGroup] Revoke previously assigned partitions marketDataRequests-0 23/07/21 13:00:37 INFO ConsumerCoordinator: [Consumer clientId=consumer-sparkGroup-1, groupId=sparkGroup] Member consumer-sparkGroup-1-89cbdb45-88a1-4d93-9817-26920020efe5 sending LeaveGroup request to coordinator DESKTOP-SCGERG2.lan:9092 (id: 2147483647 rack: null) due to the consumer is being closed 23/07/21 13:00:37 INFO ConsumerCoordinator: [Consumer clientId=consumer-sparkGroup-1, groupId=sparkGroup] Resetting generation and member id due to: consumer pro-actively leaving the group 23/07/21 13:00:37 INFO ConsumerCoordinator: [Consumer clientId=consumer-sparkGroup-1, groupId=sparkGroup] Request joining group due to: consumer pro-actively leaving the group 23/07/21 13:00:37 INFO Metrics: Metrics scheduler closed 23/07/21 13:00:37 INFO Metrics: Closing reporter org.apache.kafka.common.metrics.JmxReporter 23/07/21 13:00:37 INFO Metrics: Metrics reporters closed 23/07/21 13:00:37 INFO AppInfoParser: App info kafka.consumer for consumer-sparkGroup-1 unregistered Process finished with exit code 0
参考pom.xml
<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 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <parent> <groupId>com.samurai.lab</groupId> <artifactId>samurai-lab</artifactId> <version>1.0-SNAPSHOT</version> </parent> <artifactId>spark-job-service</artifactId> <packaging>jar</packaging> <name>spark-job-service</name> <url>https://maven.apache.org</url> <properties> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <apache-spark.version>3.4.1</apache-spark.version> <apache-kafka.version>3.5.0</apache-kafka.version> <extraJavaTestArgs> -XX:+IgnoreUnrecognizedVMOptions --add-opens=java.base/java.lang=ALL-UNNAMED --add-opens=java.base/java.lang.invoke=ALL-UNNAMED --add-opens=java.base/java.lang.reflect=ALL-UNNAMED --add-opens=java.base/java.io=ALL-UNNAMED --add-opens=java.base/java.net=ALL-UNNAMED --add-opens=java.base/java.nio=ALL-UNNAMED --add-opens=java.base/java.util=ALL-UNNAMED --add-opens=java.base/java.util.concurrent=ALL-UNNAMED --add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED --add-opens=java.base/sun.nio.ch=ALL-UNNAMED --add-opens=java.base/sun.nio.cs=ALL-UNNAMED --add-opens=java.base/sun.security.action=ALL-UNNAMED --add-opens=java.base/sun.util.calendar=ALL-UNNAMED -Djdk.reflect.useDirectMethodHandle=false </extraJavaTestArgs> </properties> <dependencies> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.13</artifactId> <version>${apache-spark.version}</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.13</artifactId> <version>${apache-spark.version}</version> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>${apache-kafka.version}</version> </dependency> <dependency> <groupId>org.slf4j</groupId> <artifactId>slf4j-simple</artifactId> <version>2.0.7</version> </dependency> <dependency> <groupId>org.mongodb.spark</groupId> <artifactId>mongo-spark-connector_2.12</artifactId> <version>3.0.1</version> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-compiler-plugin</artifactId> <version>3.11.0</version> <configuration> <source>17</source> <target>17</target> </configuration> </plugin> </plugins> </build> </project>
排查与修复思路
1. Spark Session关闭触发JVM退出
Spark的spark.stop()默认会触发JVM的shutdown hook,直接终止整个进程。你的Kafka消费逻辑运行在单线程Executor中,一旦Spark任务执行完毕调用stop,整个JVM就会被关闭,导致Kafka消费者也随之退出。
- 修复方案:
- 复用Spark Session:提前初始化一个全局Spark Session,所有消息处理都使用这个Session,不再每次创建后关闭。
- 禁用shutdown hook:如果必须关闭Session,使用
spark.sparkContext().stop(false)(注意Spark版本兼容性),避免触发JVM退出。
2. 主线程未阻塞导致进程退出
main方法中创建单线程Executor后没有阻塞主线程,主线程执行完毕后直接退出,导致Executor中的Kafka消费线程也被终止。
- 修复方案:在main方法末尾添加阻塞逻辑,比如:
或者使用CountDownLatch等工具阻塞主线程,确保Executor能持续运行。executor.awaitTermination(Long.MAX_VALUE, TimeUnit.NANOSECONDS);
3. Scala版本不兼容(Mongo Connector)
你的Spark依赖使用Scala 2.13版本(spark-core_2.13),但Mongo Spark Connector使用的是Scala 2.12版本(mongo-spark-connector_2.12),这种版本不兼容会导致类加载异常,可能引发静默进程退出。
- 修复方案:更换Mongo Connector为对应Scala 2.13的版本,比如
mongo-spark-connector_2.13:10.2.0(需注意与Spark版本的兼容性)。
4. SLF4J依赖冲突掩盖日志
日志显示存在多个SLF4J provider,虽然不是直接导致退出的原因,但可能掩盖真实的错误日志,导致无法排查问题。
- 修复方案:排除Spark自带的SLF4J依赖,只保留
slf4j-simple,或者统一使用log4j2作为日志实现。例如在Spark依赖中添加排除:<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.13</artifactId> <version>${apache-spark.version}</version> <exclusions> <exclusion> <groupId>org.slf4j</groupId> <artifactId>slf4j-log4j12</artifactId> </exclusion> <exclusion> <groupId>org.apache.logging.log4j</groupId> <artifactId>log4j-slf4j-impl</artifactId> </exclusion> </exclusions> </dependency>
5. Spark本地模式线程资源抢占
Spark在local[*]模式下会占用所有可用CPU核心的线程,可能导致Kafka消费线程被抢占或资源耗尽,进而引发进程退出。
- 修复方案:指定Spark本地模式的线程数,比如
master("local[2]"),避免占用全部资源;或者将Spark部署为集群模式,与Kafka服务分离。
内容的提问来源于stack exchange,提问作者Niko Konovalov
相关产品推荐
相关产品推荐

