【发布时间】:2016-12-26 15:41:50
【问题描述】:
我正在使用协议缓冲区将数据流发送到 Apache Flink。 我有两节课。一个是生产者,一个是消费者。 Producer 是一个 java 线程类,它从 socket 读取数据并 Protobuf 反序列化它,然后我将它存储在我的 BlockingQueue Consumer 是 Flink 中实现 SourceFunction 的类。 我使用以下方法测试了这个程序:
DataStream<Event.MyEvent> stream = env.fromCollection(queue);
而不是自定义源,它工作正常。 但是当我尝试使用 SourceFunction 类时,它会抛出这个异常:
Caused by: java.lang.RuntimeException: Unable to find proto buffer class
at com.google.protobuf.GeneratedMessageLite$SerializedForm.readResolve(GeneratedMessageLite.java:775)
...
Caused by: java.lang.ClassNotFoundException: event.Event$MyEvent
at java.net.URLClassLoader.findClass(URLClassLoader.java:381)
at java.lang.ClassLoader.loadClass(ClassLoader.java:424)
at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:331)
...
在另一次尝试中,我将两者混合为一个(实现 SourceFunction 的类)。我从套接字获取数据并使用 protobuf 对其进行反序列化并将其存储在 BlockingQueue 中,然后我立即从 BlockingQueue 中读取数据。我的代码也适用于这种方法。
但我想使用两个单独的类(多线程),但它会引发异常。 我试图在过去 2 天解决它,也做了很多搜索,但没有运气。 任何帮助将不胜感激。
制作人:
public class Producer implements Runnable {
Boolean running = true;
Socket socket = null, bufferSocket = null;
PrintStream ps = null;
BlockingQueue<Event.MyEvent> queue;
final int port;
public Producer(BlockingQueue<Event.MyEvent> queue, int port){
this.port = port;
this.queue = queue;
}
@Override
public void run() {
try {
socket = new Socket("127.0.0.1", port);
bufferSocket = new Socket(InetAddress.getLocalHost(), 6060);
ps = new PrintStream(bufferSocket.getOutputStream());
while (running) {
queue.put(Event.MyEvent.parseDelimitedFrom(socket.getInputStream()));
ps.println("Items in Queue: " + queue.size());
}
}catch (Exception e){
e.printStackTrace();
}
}
}
消费者:
public class Consumer implements SourceFunction<Event.MyEvent> {
Boolean running = true;
BlockingQueue<Event.MyEvent> queue;
Event.MyEvent event;
public Consumer(BlockingQueue<Event.MyEvent> queue){
this.queue = queue;
}
@Override
public void run(SourceContext<Event.MyEvent> sourceContext) {
try {
while (running) {
event = queue.take();
sourceContext.collect(event);
}
}catch (Exception e){
e.printStackTrace();
}
}
@Override
public void cancel() {
running = false;
}
}
Event.MyEvent 是我的 protobuf 类。我使用的是 2.6.1 版本,并使用 v2.6.1 编译了类。我仔细检查了版本以确保它不是问题。 Producer 类工作正常。 我使用 Flink v1.1.3 和 v1.1.4 对此进行了测试。 我在本地模式下运行它。
编辑:答案包含在问题中,单独发布并在此处删除。
2016 年 12 月 28 日更新
... 但我还是很好奇。是什么导致了这个错误?是 Flink 的 bug 还是我做错了什么?
...
【问题讨论】:
标签: java protocol-buffers apache-flink flink-streaming