我正在使用协议缓冲区将数据流发送到 Apache Flink。我有两节课。一种是生产者,一种是消费者。Producer 是一个 java 线程类,它从 socket 中读取数据,Protobuf 将其反序列化,然后将其存储在我的 BlockingQueue 中。 Consumer 是一个在 Flink 中实现 SourceFunction 的类。我使用以下方法测试了这个程序:
DataStream<Event.MyEvent> stream = env.fromCollection(queue);
Run Code Online (Sandbox Code Playgroud)
而不是自定义源,它工作正常。但是当我尝试使用 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)
...
Run Code Online (Sandbox Code Playgroud)
在另一次尝试中,我将两个分类为一个(实现 SourceFunction 的类)。我从套接字获取数据并使用 protobuf 将其反序列化并将其存储在 BlockingQueue 中,然后我立即从 BlockingQueue 中读取。我的代码也适用于这种方法。
但我想使用两个单独的类(多线程),但它会抛出该异常。我试图在过去 2 天内解决它,也做了很多搜索,但没有运气。任何帮助都会受到重视。
制作人:
public class Producer implements Runnable {
Boolean running = true;
Socket socket = null, bufferSocket = null;
PrintStream ps = …Run Code Online (Sandbox Code Playgroud)