【问题标题】:Getting ClassNotFound Exception in Flink SourceFunction在 Flink SourceFunction 中获取 ClassNotFound 异常
【发布时间】:2016-12-26 15:41:50
【问题描述】:

我正在使用协议缓冲区将数据流发送到 Apache Flink。 我有两节课。一个是生产者,一个是消费者。 Producer 是一个 java 线程类,它从 socket 读取数据并 Protobuf 反序列化它,然后我将它存储在我的 BlockingQueue Consumer 是 Fl​​ink 中实现 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


    【解决方案1】:

    提问者已经找到了一种方法来实现这一点。我已经从问题中提取了相关部分。请注意,它发生的原因仍然无法解释。

    我没有使用引号语法,因为它是很多文本,但以下是提问者分享的:

    所以我终于让它工作了。我在 SourceFunction(Consumer)中创建了我的 BlockingQueue 对象,并从 SourceFunction 类(Consumer)内部调用了 Producer 类,而不是在程序的 main 方法中制作 BlockingQueue 和调用 Producer 类。现在它可以工作了!

    这是我在 Flink 中的完整工作代码:

    public class Main {
    
    public static void main(String[] args) throws Exception {
    
        final int port, buffer;
        //final String ip;
        try {
            final ParameterTool params = ParameterTool.fromArgs(args);
            port = params.getInt("p");
            buffer = params.getInt("b");
        } catch (Exception e) {
            System.err.println("No port number and/or buffer size specified.");
            return;
        }
    
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    
    
        DataStream<Event.MyEvent> stream = env.addSource(new Consumer(port, buffer));
        //DataStream<Event.MyEvent> stream = env.fromCollection(queue);
    
    
        Pattern<Event.MyEvent, ?> crashedPattern = Pattern.<Event.MyEvent>begin("start")
                .where(new FilterFunction<Event.MyEvent>() {
                    @Override
                    public boolean filter(Event.MyEvent myEvent) throws Exception {
                        return (myEvent.getItems().getValue() >= 120);
                    }
                })
                .<Event.MyEvent>followedBy("next").where(new FilterFunction<Event.MyEvent>() {
                    @Override
                    public boolean filter(Event.MyEvent myEvent) throws Exception {
                        return (myEvent.getItems().getValue() <= 10);
                    }
                })
                .within(Time.seconds(3));
    
        PatternStream<Event.MyEvent> crashed = CEP.pattern(stream.keyBy(new KeySelector<Event.MyEvent, String>() {
            @Override
            public String getKey(Event.MyEvent myEvent) throws Exception {
                return myEvent.getEventType();
            }
        }), crashedPattern);
    
        DataStream<String> alarm = crashed.select(new PatternSelectFunction<Event.MyEvent, String>() {
            @Override
            public String select(Map<String, Event.MyEvent> pattern) throws Exception {
                Event.MyEvent start = pattern.get("start");
                Event.MyEvent next = pattern.get("next");
                return start.getEventType() + " | Speed from " + start.getItems().getValue() + " to " + next.getItems().getValue() + " in 3 seconds\n";
            }
        });
    
        DataStream<String> rate = alarm.windowAll(TumblingProcessingTimeWindows.of(Time.seconds(1)))
                .apply(new AllWindowFunction<String, String, TimeWindow>() {
                    @Override
                    public void apply(TimeWindow timeWindow, Iterable<String> iterable, Collector<String> collector) throws Exception {
                        int sum = 0;
                        for (String s: iterable) {
                            sum ++;
                        }
                        collector.collect ("CEP Output Rate: " + sum + "\n");
                    }
                });
    
        rate.writeToSocket(InetAddress.getLocalHost().getHostName(), 7070, new SimpleStringSchema());
    
        env.execute("Flink Taxi Crash Streaming");
    }
    
    private static class Producer implements Runnable {
    
        Boolean running = true;
        Socket socket = null, bufferSocket = null;
        PrintStream ps = null;
        BlockingQueue<Event.MyEvent> queue;
        final int port;
    
        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();
            }
        }
    
    }
    
    private static class Consumer implements SourceFunction<Event.MyEvent> {
    
        Boolean running = true;
        final int port;
        BlockingQueue<Event.MyEvent> queue;
    
        Consumer(int port, int buffer){
            queue = new ArrayBlockingQueue<>(buffer);
            this.port = port;
        }
    
        @Override
        public void run(SourceContext<Event.MyEvent> sourceContext) {
            try {
                new Thread(new Producer(queue, port)).start();
                while (running) {
                    sourceContext.collect(queue.take());
                }
            }catch (Exception e){
                e.printStackTrace();
            }
        }
    
        @Override
        public void cancel() {
            running = false;
        }
    }
    

    【讨论】:

      猜你喜欢
      • 2012-06-22
      • 1970-01-01
      • 1970-01-01
      • 2012-06-30
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2012-10-13
      • 2016-08-11
      相关资源
      最近更新 更多