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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:30:42