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

Spark 2.2.0流作业每日触发java.io.OptionalDataException故障求助

Spark Streaming作业运行1.5天后崩溃:OptionalDataException分析与解决

看起来你遇到的是Spark Streaming作业运行一段时间后触发的Java序列化异常,结合你的作业配置、代码逻辑和异常栈,我来帮你拆解问题根源和可行的解决方案:

问题根源分析

从异常栈可以看到,java.io.OptionalDataException发生在HashMap的反序列化过程中,异常出现在executor向driver传输序列化对象的阶段(也就是你调用collect()拉取InsertRequestWrapper列表的时候)。

为什么会运行1.5天才崩溃?因为这个异常不是普遍触发的,大概率是以下两种情况:

  1. 特定的Kafka消息导致你的自定义对象(InsertRequestWrapper/SingleEventBaseDocument)序列化后的数据不完整
  2. 大量数据累积传输时,Java序列化的稳定性不足,出现数据流截断或解析错误

另外,你当前的架构把所有处理后的数据拉到driver再写入ES,不仅让driver成为性能瓶颈,也放大了序列化传输的风险——一旦某个executor传输的序列化数据有问题,整个批次就会失败。

解决方案

1. 修复自定义对象的序列化实现

首先检查你的自定义类是否符合序列化规范:

  • 确保InsertRequestWrapper、SingleEventBaseDocument以及它们嵌套的所有对象(包括HashMap等集合的泛型类型)都正确实现了java.io.Serializable接口
  • 如果你的类自定义了readObject()/writeObject()方法,仔细检查逻辑是否正确——HashMap的默认readObject需要完整读取所有键值对,手动实现时很容易出现数据读取越界的情况
  • 绝对不要在序列化对象中包含非序列化资源(比如数据库连接、IO流、未序列化的第三方对象),这些会导致序列化失败或数据损坏

2. 切换到Kryo序列化替代Java序列化

Spark默认的Java序列化效率低且稳定性不足,尤其是对于复杂对象。换成Kryo序列化能大幅降低序列化异常的概率:
在你的Spark配置中添加:

SparkConf sparkConf = new SparkConf();
sparkConf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer");
// 注册你的自定义类,提升序列化效率
sparkConf.registerKryoClasses(new Class[]{
    InsertRequestWrapper.class,
    SingleEventBaseDocument.class
});

3. 重构ES写入逻辑,避免Driver单点瓶颈

你之前尝试过在executor直接写ES但速度慢,大概率是ES批量配置不合理导致的。调整配置后,在executor端写入ES不仅能解决序列化问题,还能提升作业性能:

先修正ES批量配置(注意单位!)

你的当前配置es.bulk.action.bytes=20没有指定单位,默认是字节,这会导致频繁的小批量提交,速度极慢。改成合理的配置:

es.bulk.action.count=5000
es.bulk.action.bytes=20mb  # 单批最大字节数
es.bulk.action.flush.interval=30s
es.bulk.backoff.policy.interval=1s
es.bulk.number.of.retries=3  # 增加重试次数,避免写入失败

用Spark-ES库在Executor端写入

替换你当前的collect()+driver写入逻辑,直接用官方库在executor端写入:

kafkaStream.foreachRDD(rdd -> {
    String batchIdentifier = Long.toHexString(Double.doubleToLongBits(Math.random()));
    LOGGER.info("@@ [" + batchIdentifier + "] Starting batch ...");
    Instant batchStart = Instant.now();

    OffsetRange[] offsetRanges = ((HasOffsetRanges) rdd.rdd()).offsetRanges();

    // 在Executor端处理并写入ES
    long processedCount = rdd.mapPartitionsWithIndex((index, eventsIterator) -> {
        LOGGER.info("@@ Consuming " + offsetRanges[index].count() + " events partition: " + index);
        if (!eventsIterator.hasNext()) {
            return Collections.emptyIterator();
        }
        List<SingleEventBaseDocument> eventList = getSingleEventBaseDocuments(eventsIterator);
        LOGGER.info("@@ Processed " + offsetRanges[index].count() + " events partition: " + index + " list size: " + eventList.size());
        return eventList.iterator();
    }).map(doc -> doc) // 根据你的ES文档结构调整映射
      .saveToEs("your_es_index/your_doc_type"); // 替换为你的ES索引和类型

    LOGGER.info("@@ [" + batchIdentifier + "] Finished batch of " + processedCount + " messages in " + (Instant.now().toEpochMilli() - batchStart.toEpochMilli()) + "ms");

    // 提交偏移量到Kafka
    ((CanCommitOffsets) kafkaStream.inputDStream()).commitAsync(offsetRanges);
});

4. 排查触发异常的特定数据

因为异常是偶发的,建议在数据处理逻辑中添加异常捕获,记录可能触发问题的原始消息:

private List<SingleEventBaseDocument> getSingleEventBaseDocuments(final Iterator<ConsumerRecord<String, byte[]>> eventsIterator) {
    Iterable<ConsumerRecord<String, byte[]>> iterable = () -> eventsIterator;
    return StreamSupport.stream(iterable.spliterator(), true)
        .map(record -> {
            try {
                return toEnrichedEvent(record);
            } catch (Exception e) {
                LOGGER.error("Failed to enrich event, key: {}, value: {}", record.key(), new String(record.value()), e);
                return null;
            }
        })
        .filter(Objects::nonNull)
        .filter(this::isValidEvent)
        .map(this::toEventDocument)
        .collect(Collectors.toList());
}

这样可以定位到具体的坏消息,针对性修复数据解析逻辑。


内容的提问来源于stack exchange,提问作者alina

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:16:26