【问题标题】:load data from cassandra to flink从 cassandra 加载数据到 flink
【发布时间】:2019-08-22 12:32:29
【问题描述】:

我如何在 java 或 python 中使用 flink 从 cassandra 读取数据。我没有找到关于这个主题的好的文档。我想将此代码连接到我的网站。那是怎么回事。

我有这段代码,但它运行异常

java.lang.NoClassDefFoundError: org/apache/flink/api/common/typeinfo/TypeHint
    at java.lang.Class.getDeclaredMethods0(Native Method)
    at java.lang.Class.privateGetDeclaredMethods(Class.java:2701)
    at java.lang.Class.privateGetMethodRecursive(Class.java:3048)
    at java.lang.Class.getMethod0(Class.java:3018)
    at java.lang.Class.getMethod(Class.java:1784)
    at sun.launcher.LauncherHelper.validateMainClass(LauncherHelper.java:544)
    at sun.launcher.LauncherHelper.checkAndLoadMain(LauncherHelper.java:526)
Caused by: java.lang.ClassNotFoundException: org.apache.flink.api.common.typeinfo.TypeHint
    at java.net.URLClassLoader.findClass(URLClassLoader.java:382)
    at java.lang.ClassLoader.loadClass(ClassLoader.java:424)
    at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:349)
    at java.lang.ClassLoader.loadClass(ClassLoader.java:357)
import com.datastax.driver.core.Cluster;
import org.apache.flink.api.common.typeinfo.TypeHint;
import org.apache.flink.api.java.DataSet;
import org.apache.flink.api.java.ExecutionEnvironment;
import org.apache.flink.api.java.tuple.Tuple5;
import org.apache.flink.api.java.typeutils.TupleTypeInfo;
import org.apache.flink.batch.connectors.cassandra.CassandraInputFormat;
import org.apache.flink.streaming.connectors.cassandra.ClusterBuilder;

public class main {
    public static void main(String[] args) {
        final ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();

        ClusterBuilder clusterBuilder = new ClusterBuilder() {

            @Override
            public Cluster buildCluster(Cluster.Builder builder) {

                return builder.addContactPoint("127.0.0.1").withPort(9042).build();
            }
        };

        DataSet<Tuple5<String,String,String,Integer,String>> inputRecords = env
                .createInput
                        (new CassandraInputFormat<Tuple5<String,String,String,Integer,String>>("Select profilealternative from profiles.profile LIMIT 1 ;",clusterBuilder)
            ,TupleTypeInfo.of(new TypeHint<Tuple5<String,String,String,Integer,String>>() {}));

    }
}

【问题讨论】:

    标签: api cassandra apache-flink


    【解决方案1】:

    您在类路径中缺少的类来自 flink-dist.jar 文件(名称也取决于版本,例如 flink-dist_2.12-1.9.0.jar)。请在运行程序时检查文件是否在您的类路径中。此外,您可能需要其他通常存在于 apache-flink 发行版lib 目录中的 jar。您也可以创建一个包含所有依赖项的 fat jar,然后运行。

    【讨论】:

      【解决方案2】:

      您的问题的主要原因是ClassNotFoundException - 这意味着您的应用程序没有链接所有必需的库来创建 uberjar,或者类路径中缺少相应的库。

      【讨论】:

        猜你喜欢
        • 2018-09-13
        • 2013-12-12
        • 1970-01-01
        • 1970-01-01
        • 2022-01-28
        • 2017-10-27
        • 2015-12-30
        • 2017-08-21
        • 1970-01-01
        相关资源
        最近更新 更多