NoClassDefFoundError при запуске flink с соединителем KafkaJAVA

Программисты JAVA общаются здесь
Ответить
Anonymous
 NoClassDefFoundError при запуске flink с соединителем Kafka

Сообщение Anonymous »

Я пытаюсь передать данные из Kafka в потоковом режиме с помощью Flink. Мой код компилируется без ошибок, но при запуске я получаю следующую ошибку:

Код: Выделить всё

Error: A JNI error has occurred, please check your installation and try again
Exception in thread "main" java.lang.NoClassDefFoundError:
org/apache/flink/streaming/util/serialization/DeserializationSchema
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.streaming.util.serialization.DeserializationSchema
at java.net.URLClassLoader.findClass(URLClassLoader.java:381)
at java.lang.ClassLoader.loadClass(ClassLoader.java:424)
at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:335)
at java.lang.ClassLoader.loadClass(ClassLoader.java:357)
... 7 more
Мой список зависимостей POM выглядит следующим образом:

Код: Выделить всё

    

org.apache.flink
flink-java
1.3.2


org.apache.flink
flink-streaming-core
0.9.1


org.apache.flink
flink-clients
0.10.2


org.apache.flink
flink-connector-kafka-0.9_2.11
1.3.2


com.googlecode.json-simple
json-simple
1.1


Код Java, который я пытаюсь запустить, просто подписывается на тему Kafka под названием «streamer»:

Код: Выделить всё

import java.util.Properties;
import java.util.Arrays;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer09;
import org.apache.flink.streaming.util.serialization.SimpleStringSchema;
import org.apache.flink.streaming.util.serialization.DeserializationSchema;

public class StreamConsumer {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
Properties properties = new Properties();
properties.setProperty("bootstrap.servers", "localhost:9092");
properties.setProperty("group.id", "samplegroup");

DataStream messageStream = env.addSource(new FlinkKafkaConsumer09("streamer", new SimpleStringSchema(), properties));

messageStream.rebalance().map(new MapFunction() {
private static final long serialVersionUID = -6867736771747690202L;
@Override
public String map(String value) throws Exception {
return "Streamed data: " + value;
}
}).print();
env.execute();
}
}
Информация о системе:

1. Версия Кафки: 0.9.0.1

2. Версия Flink: 1.3.2

3. Версия OpenJDK: 1.8

Хотя я использую maven, я не думаю, что это проблема с maven, поскольку я получаю ту же ошибку, даже когда пытаюсь использовать maven. Я вручную загрузил все необходимые файлы .jar в папку и указал путь к этой папке с опцией -cp при компиляции с помощью javac. Во время выполнения я получаю ту же ошибку, что и выше, но во время компиляции ошибок нет.

Подробнее здесь: https://stackoverflow.com/questions/466 ... -connector
Ответить

Быстрый ответ

Изменение регистра текста: 
Смайлики
:) :( :oops: :roll: :wink: :muza: :clever: :sorry: :angel: :read: *x)
Ещё смайлики…
   
К этому ответу прикреплено по крайней мере одно вложение.

Если вы не хотите добавлять вложения, оставьте поля пустыми.

Максимально разрешённый размер вложения: 15 МБ.

Вернуться в «JAVA»