如何在Apache Flink中删除Cassandra数据?自定义Sink序列化异常解决
问题描述
在Apache Flink里用CassandraSink插入数据很顺手,但我找不到删除数据的方法。尝试写自定义Sink的时候碰到了NotSerializableException,求问怎么写出能实现删除操作的代码?
我写的自定义Sink代码
public class MyCassandraSink implements SinkFunction<String> { private Cluster cluster = Cluster.builder() .addContactPoint("127.0.0.1") .build(); private Session cassandra = cluster.connect("mykeyspace"); @Override public void invoke(String value, Context context) throws Exception { cassandra.execute("SOME DELETE QUERY"); } }
抛出的异常
Exception in thread "main" org.apache.flink.api.common.InvalidProgramException: [com.datastax.driver.core.SessionManager@3b0fe47a] is not serializable. The object probably contains or references non serializable fields. at org.apache.flink.api.java.ClosureCleaner.clean(ClosureCleaner.java:151) at org.apache.flink.api.java.ClosureCleaner.clean(ClosureCleaner.java:126) at org.apache.flink.api.java.ClosureCleaner.clean(ClosureCleaner.java:126) at org.apache.flink.api.java.ClosureCleaner.clean(ClosureCleaner.java:126) at org.apache.flink.api.java.ClosureCleaner.clean(ClosureCleaner.java:126) at org.apache.flink.api.java.ClosureCleaner.clean(ClosureCleaner.java:71) at org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.clean(StreamExecutionEnvironment.java:1574) at org.apache.flink.streaming.api.datastream.DataStream.clean(DataStream.java:185) at org.apache.flink.streaming.api.datastream.DataStream.addSink(DataStream.java:1227) at com.meshkan.streaming.entry.EventListener.main(EventListener.java:42) Caused by: java.io.NotSerializableException: com.datastax.driver.core.SessionManager at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1184) at java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:348) at java.util.concurrent.CopyOnWriteArrayList.writeObject(CopyOnWriteArrayList.java:973) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at java.io.ObjectStreamClass.invokeWriteObject(ObjectStreamClass.java:1140) at java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1496) at java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1432) at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1178) at java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1548) at java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1509) at java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1432) at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1178) at java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:348) at org.apache.flink.util.InstantiationUtil.serializeObject(InstantiationUtil.java:586) at org.apache.flink.api.java.ClosureCleaner.clean(ClosureCleaner.java:133) ... 9 more
解决方案
这个问题我之前也踩过坑,核心原因是Flink会把Sink的实例序列化后分发到各个TaskManager节点,但你代码里直接初始化的Cluster和Session都是不可序列化的对象,所以才会触发这个异常。
正确的做法是使用RichSinkFunction(它继承了SinkFunction,还提供了生命周期管理方法),把Cassandra连接的初始化放在open()方法里,资源释放放在close()方法里,这样连接对象不会被序列化,而是在每个Task节点上单独初始化。
修改后的代码如下:
public class MyCassandraDeleteSink extends RichSinkFunction<String> { // 用transient修饰,告诉序列化机制忽略这些字段 private transient Cluster cluster; private transient Session session; @Override public void open(Configuration parameters) throws Exception { super.open(parameters); // 每个Task启动时初始化Cassandra连接 cluster = Cluster.builder() .addContactPoint("127.0.0.1") .build(); session = cluster.connect("mykeyspace"); } @Override public void invoke(String value, Context context) throws Exception { // 根据输入的value构造具体的删除语句,这里只是示例 session.execute(String.format("DELETE FROM your_table WHERE id = '%s'", value)); } @Override public void close() throws Exception { super.close(); // Task结束时释放连接资源,避免泄漏 if (session != null) { session.close(); } if (cluster != null) { cluster.close(); } } }
关键细节说明
- 用RichSinkFunction替代SinkFunction:它的
open()和close()方法分别对应Task的启动和销毁阶段,完美适配外部资源的初始化与释放。 - transient关键字:标记不需要序列化的成员变量,避免Flink尝试序列化Cassandra的连接对象。
- 连接本地化初始化:每个Task节点单独创建自己的Cassandra连接,既解决了序列化问题,也能更好地控制连接资源的使用。
另外,如果你想更贴合Flink Cassandra连接器的生态,也可以研究下CassandraSink的自定义语句扩展,但用RichSinkFunction实现删除是最直接高效的方式。
内容的提问来源于stack exchange,提问作者user3640028
相关产品推荐
相关产品推荐

