【问题标题】:Flink cassandraOutputFormat tuple needs frozen valueFlink cassandraOutputFormat 元组需要冻结值
【发布时间】: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&lt;frozen&lt;dataobj&gt; ??缺少&gt;
  • 是的,它仍然可以正常运行(老实说,这很奇怪,缺少它也很好)。我添加了缺少的“>”以确保也是如此。
  • 添加您的表并输入架构
  • 尝试将@Frozen("list&lt;frozen&lt;dataobj&gt;")改为@Frozen
  • 改为@Frozen,同样的错误。

标签: java cassandra apache-flink


【解决方案1】:

CustomDataObj类上添加UDT注解

@UDT(name = "dataobj")
public class CustomDataObj { 
    //...... 
}

已编辑

jobid 注释更改为@Column(name = "jobid")dataobjs Frozen Annotation 更改为@Frozen

@Table(keyspace="mykeyspace", name="mytablename")
public class AlphaGroupingObject implements Serializable {

    @Column(name = "jobid")
    private String jobId;

    @Column(name = "datalist")
    @Frozen
    private List<CustomDataObj> dataobjs;
    @Column(name = "userid")
    private String userid;

    //Getters and Setters
}

【讨论】:

  • 进行了您在第一次编辑中概述的更改,仍然抛出相同的错误。
  • 我看到它是如何为注释为 UDT 的 pojo 建议自定义编解码器的,但我的问题是如何将元组的字段设置为冻结。您在 pojo 中执行此操作的方式只是添加 @Frozen,但对于无法完成的元组。
【解决方案2】:

我相信我找到了一种比向 cassandraOutputFormat 提供元组更好的方法,但它在技术上仍然无法回答这个问题,所以我不会将此标记为答案。我最终使用了 cassandra 的对象映射器,所以我可以将 pojo 发送到桌面。仍然需要验证数据是否已成功发送到表中,并且一切都按照它的实现方式正常工作,但我认为这会对面临类似问题的任何人有所帮助。

这是概述解决方案的文档:http://docs.datastax.com/en/developer/java-driver/2.1/manual/object_mapper/using/

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-04-27
    • 1970-01-01
    • 2013-02-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-02-09
    相关资源
    最近更新 更多