Spark 2.2.0流作业每日触发java.io.OptionalDataException故障求助
看起来你遇到的是Spark Streaming作业运行一段时间后触发的Java序列化异常,结合你的作业配置、代码逻辑和异常栈,我来帮你拆解问题根源和可行的解决方案:
问题根源分析
从异常栈可以看到,java.io.OptionalDataException发生在HashMap的反序列化过程中,异常出现在executor向driver传输序列化对象的阶段(也就是你调用collect()拉取InsertRequestWrapper列表的时候)。
为什么会运行1.5天才崩溃?因为这个异常不是普遍触发的,大概率是以下两种情况:
- 特定的Kafka消息导致你的自定义对象(
InsertRequestWrapper/SingleEventBaseDocument)序列化后的数据不完整 - 大量数据累积传输时,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

