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

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方法末尾添加阻塞逻辑,比如:
    executor.awaitTermination(Long.MAX_VALUE, TimeUnit.NANOSECONDS);
    
    或者使用CountDownLatch等工具阻塞主线程,确保Executor能持续运行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 18:25:54