Storm 1.2.1拓扑运行时抛出NotSerializableException异常求助
首先,这个异常的核心原因很明确:Storm在集群模式下会把你的拓扑组件(Spout、Bolt)以及相关配置对象序列化后,分发到集群的各个Worker节点上运行。任何不可序列化的对象被包含在这些需要分发的内容里,就会触发NotSerializableException。你开启TOPOLOGY_TESTING_ALWAYS_TRY_SERIALIZE配置让本地模式也做序列化检查,这一步非常棒——能提前在本地暴露问题,不用等到集群部署才踩坑。
下面是具体的排查和解决步骤:
1. 先定位具体的不可序列化类
你给出的异常堆栈被截断了,完整的堆栈末尾一定会明确指出哪个类不可序列化(比如java.io.NotSerializableException: com.yourcompany.YourProblemClass)。先找到这个类,这是解决问题的关键。
2. 针对不同场景的解决方案
场景A:自定义类不可序列化
如果是你自己写的类导致的问题,只需要让它实现java.io.Serializable接口就行。在Scala里这个很简单:
// 非case类的情况,手动继承Serializable class CustomData(id: String, value: Int) extends Serializable // Scala的Case类默认自动实现Serializable,所以如果是case类出现这个问题,大概率是它引用了其他不可序列化的成员 case class Order(id: String, user: User) // 这里如果User不可序列化,Order也会出问题
场景B:Spout/Bolt引用了不可序列化的成员变量
很多人会犯这个错:把数据库连接、第三方服务客户端这类不可序列化的对象,作为Spout/Bolt的构造函数参数或者成员变量。比如这种错误写法:
// 错误:把不可序列化的RedisClient直接作为构造参数传递 class RedisWriterBolt(redisClient: RedisClient) extends BaseBasicBolt { override def execute(tuple: Tuple, collector: BasicOutputCollector): Unit = { redisClient.write(tuple.getStringByField("data")) } }
这类对象(连接、客户端)只能在Worker节点本地初始化,绝对不能被序列化分发。正确的写法是在prepare(Bolt)或open(Spout)方法里初始化:
// 正确:在prepare方法里初始化客户端,避免序列化 class RedisWriterBolt extends BaseBasicBolt { private var redisClient: RedisClient = _ override def prepare(config: Map[String, AnyRef], context: TopologyContext, collector: BasicOutputCollector): Unit = { // 这里的初始化逻辑是在Worker节点本地执行的,不需要序列化 redisClient = new RedisClient(config.get("redis_host").asInstanceOf[String], 6379) } override def execute(tuple: Tuple, collector: BasicOutputCollector): Unit = { redisClient.write(tuple.getStringByField("data")) } }
场景C:闭包/匿名函数引用了不可序列化对象
Scala的闭包很方便,但如果在Spout/Bolt里用闭包引用了外部的不可序列化对象,也会触发异常。比如:
// 错误:闭包引用了外部的非序列化对象 val nonSerializableObj = new NonSerializableService() val bolt = new BaseBasicBolt { override def execute(tuple: Tuple, collector: BasicOutputCollector): Unit = { nonSerializableObj.process(tuple) // 闭包捕获了nonSerializableObj,导致整个匿名Bolt不可序列化 } }
解决方法是把闭包改成显式的类,或者把不可序列化对象的初始化移到Bolt的prepare方法里。
场景D:Storm配置里放了不可序列化对象
检查你往Config里put的所有值,确保它们都是可序列化的。比如不要把非序列化的对象直接放到配置里,而是传递字符串、数字这类基本类型,或者可序列化的自定义配置类。
3. 额外的Scala专属注意事项
- 避免用
lazy val引用不可序列化的对象:lazy val的初始化逻辑会被序列化机制触发,导致不可序列化的对象被尝试序列化。 - 检查伴生对象:如果你的类的伴生对象包含不可序列化的成员,也可能影响类的序列化(不过这种情况比较少见)。
总结一下:Storm的序列化要求是集群运行的核心前提——所有需要跨节点分发的对象必须是可序列化的。你用TOPOLOGY_TESTING_ALWAYS_TRY_SERIALIZE提前发现问题的思路非常正确,接下来只要定位到具体的不可序列化类,按照上面的场景对应修复就可以了。
内容的提问来源于stack exchange,提问作者Saurabh

