【发布时间】:2017-05-17 14:45:01
【问题描述】:
我有一个 flink 项目,它将作为批处理作业在 cassandra 表中插入数据。我已经有一个 flink 流项目,它正在将 pojo 写入同一个 cassandra 表,但是 cassandraOutputFormat 需要将数据作为元组(希望在某些时候改变为接受像 CassandraSink 那样的 pojo)。所以这是我拥有的pojo:
@Table(keyspace="mykeyspace", name="mytablename")
public class AlphaGroupingObject implements Serializable {
@Column(name = "jobId")
private String jobId;
@Column(name = "datalist")
@Frozen("list<frozen<dataobj>")
private List<CustomDataObj> dataobjs;
@Column(name = "userid")
private String userid;
//Getters and Setters
}
以及我从这个 pojo 中制作的元组数据集:
DataSet<Tuple3<String, List<CustomDataObj>, String>> outputDataSet = listOfAlphaGroupingObject.map(new AlphaGroupingObjectToTuple3Mapper());
这也是触发输出的那一行:
outputDataSet.output(new CassandraOutputFormat<>("INSERT INTO mykeyspace.mytablename (jobid, datalist, userid) VALUES (?,?,?);", clusterThatWasBuilt));
现在我遇到的问题是,当我尝试运行它时,当它尝试将其输出到 cassandra 表时出现此错误:
Caused by: com.datastax.driver.core.exceptions.CodecNotFoundException:
Codec not found for requested operation: [frozen<mykeyspace.dataobj> <-> flink.custom.data.CustomDataObj]
所以我知道它什么时候是 pojo,我只需要在字段中添加 @Frozen 注释,但我不知道如何为元组执行此操作。解决此问题的最佳/正确方法是什么?还是我做了一些不必要的事情,因为实际上有一种方法可以通过我还没有找到的 cassandraOutputFormat 发送 pojo?
感谢您提前提供的所有帮助!
编辑:
这也是 CustomDataObj 类的代码:
@UDT(name="dataobj", keyspace = "mykeyspace")
public class CustomDataObj implements Serializable {
@Field(name = "userid")
private String userId;
@Field(name = "groupid")
private String groupId;
@Field(name = "valuetext")
private String valueText;
@Field(name = "comments")
private String comments;
//Getters and setters
}
编辑 2
包括与 CustomDataObj 绑定的 cassandra 中的表架构和 mytablename 架构。
CREATE TYPE mykeyspace.dataobj (
userid text,
groupid text,
valuetext text,
comments text
);
CREATE TABLE mykeyspace.mytablename (
jobid text,
datalist list<frozen<dataobj>>,
userid text,
PRIMARY KEY (jobid, userid)
);
【问题讨论】:
-
对吗
list<frozen<dataobj>??缺少> -
是的,它仍然可以正常运行(老实说,这很奇怪,缺少它也很好)。我添加了缺少的“>”以确保也是如此。
-
添加您的表并输入架构
-
尝试将
@Frozen("list<frozen<dataobj>")改为@Frozen -
改为@Frozen,同样的错误。
标签: java cassandra apache-flink