使用Spark Runner运行含KafkaIO的Apache Beam代码报StackOverflowError
报错根因
该栈溢出错误出现在Scala集合的序列化阶段,由版本不兼容与序列化配置错误共同导致:
- Apache Beam 2.33.0 默认的
beam-runners-spark依赖仅适配Spark 2.x版本,与你使用的Spark 3.1.2 依赖的Scala版本、序列化逻辑不兼容,会触发序列化死循环导致栈溢出。 - Beam 2.33.0 官方适配的最高Kafka版本为2.8.x,你使用的Kafka 3.0.0 跨大版本存在API不兼容问题,会加重序列化异常。
- 未配置Spark使用Kryo序列化器,默认的Java序列化器处理Beam复合对象时更容易出现栈溢出。
解决方案
按照以下步骤修改即可解决:
1. 替换适配Spark 3.x的Runner依赖
将原有beam-runners-spark依赖替换为Spark 3专属版本,同时统一所有Scala依赖为2.12.x版本(与Spark 3.1.2 默认的Scala版本对齐):
<dependency> <groupId>org.apache.beam</groupId> <artifactId>beam-runners-spark-3</artifactId> <version>2.33.0</version> </dependency>
2. 调整Kafka版本到兼容范围
将Kafka客户端版本降级到Beam 2.33.0 官方验证的兼容版本:
<kafka.version>2.8.1</kafka.version>
如果业务必须使用Kafka 3.0.0,需要手动排除Kafka客户端传递的冲突依赖,兼容性不做保证。
3. 补充Spark序列化配置
初始化PipelineOptions时添加以下序列化相关配置:
// 使用Kryo序列化器替代默认Java序列化 options.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer"); // 注册Beam专用的Kryo注册器 options.set("spark.kryo.registrator", "org.apache.beam.runners.spark.structuredstreaming.KryoRegistrator"); // 调大序列化缓存与线程栈深度避免溢出 options.set("spark.kryoserializer.buffer.max", "64m"); options.set("spark.driver.extraJavaOptions", "-Xss1024k"); options.set("spark.executor.extraJavaOptions", "-Xss1024k");
4. 功能验证
先运行不含KafkaIO的简单Spark Runner任务,确认Runner本身可正常运行后,再加入KafkaIO逻辑验证全链路可用性。
内容的提问来源于stack exchange,提问作者Rahul Dawn
相关产品推荐
相关产品推荐

