【问题标题】:Flink issues with getting data from cassandra to use as datasetFlink 从 cassandra 获取数据以用作数据集的问题
【发布时间】:2017-05-12 23:53:48
【问题描述】:

我正在尝试从 cassandra 表中提取数据以用作数据集,但遇到了两个问题。

第一个是 cassandraInputFormat 只返回一个元组,我宁愿没有一个 tuple12 并且只使用一个 pojo 来定义它期望返回的内容。所以我不知道这是否只是我必须接受的事情,是否有办法使用 pojo 代替 cassandraConnector (https://ci.apache.org/projects/flink/flink-docs-release-1.2/dev/connectors/cassandra.html),或者使用 cassandraInputFormat 不是最好的方法获取数据。

另一个问题是即使我从 cassandraInputFormat 中提取的数据(无论是否为元组)我都不知道如何将其设置为数据源。对于文件、csv 和 HDFS,有很多方法 (https://ci.apache.org/projects/flink/flink-docs-master/api/java/org/apache/flink/api/java/ExecutionEnvironment.html#ExecutionEnvironment--),但没有一个明确用于 cassandra。所以我的猜测是,我需要使用 cassandraInputFormat 提取数据并使用 .fromElements() 或 .fromCollecton() 之类的东西,以及正确的方法是什么。

提前感谢您的帮助!

更新:

这个“有效”(感谢 Chesnay Schepler 的帮助):

DataSet<Tuple2<String, String>> testSet = 
exEnv.createInput(cassandraInputFormat, TypeInformation.of(newTypeHint<Tuple2<String, String>>(){})); 

但是这个错误现在正在发生...

Exception in thread "main" org.apache.flink.optimizer.CompilerException: 
Error translating node 'Data Source "at execute(CodeBatchProcessorImpl.java:85) 
(org.apache.flink.batch.connectors.cassandra.CassandraInputFormat)" : NONE 
[[ GlobalProperties [partitioning=RANDOM_PARTITIONED] ]] [[ LocalProperties [ordering=null, grouped=null, unique=null] ]]':
Could not write the user code wrapper class org.apache.flink.api.common.operators.util.UserCodeObjectWrapper :
java.io.NotSerializableException: flink.streaming.code.CodeBatchProcessorImpl

进一步包括:

Caused by: java.io.NotSerializableException: org.apache.flink.api.java.LocalEnvironment

更新 2:

必须将环境设置为瞬态。现已修复!

【问题讨论】:

    标签: java cassandra apache-flink


    【解决方案1】:

    您可以通过调用 ExecutionEnvironment#createInput(InputFormat) 来使用 CassandraInputFormat 和所有 InputFormats。

    目前没有将元素直接读取为 POJO 的选项。最简单的解决方法是在将元组转换为所需 POJO 的接收器之后添加一个 MapFunction。

    【讨论】:

    • 我试过“DataSet> testSet = exEnv.createInput(cassandraInputFormat);”但它返回错误:“org.apache.flink.api.common.InvalidProgramException:输入格式返回的类型无法自动确定。请使用'createInput(InputFormat)明确指定生成类型的TypeInformation , TypeInformation)' 方法代替"。关于如何解决的任何想法?我尝试添加 Tuple2 和一堆其他方法,但没有任何东西是正确的。
    • 弄清楚我必须做些什么来解决这个错误。这有效: DataSet> testSet = exEnv.createInput(cassandraInputFormat, TypeInformation.of(new TypeHint>(){}));
    • 啊,是的,您必须明确定义输出类型,原因是 InputFormat 本身在创建作业时没有任何关于它将发出什么的信息。
    猜你喜欢
    • 2017-08-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-11-16
    • 1970-01-01
    相关资源
    最近更新 更多