【问题标题】:How to pass a function name As input for Flink Map function如何传递函数名作为 Flink Map 函数的输入
【发布时间】:2018-12-16 06:29:08
【问题描述】:

我有一个基于 Flink Java API 的类:

public class SP implements Serializable {

    private transient StreamExecutionEnvironment env;
    private DataStream<byte[]> data ;
}

然后我尝试为类 SP 编写一个方法,该方法获取函数名称并将该函数应用于 data 字段行。

public DataStream<Object> myMap(Function<Object, Object> func) {
        return data.map(x -> func.apply(x));
    }

所以在 main 方法中,我创建了一个简单的函数并将其传递给myMap 函数。

public static void main(String[] args) throws Exception {
        SP temp = new SP();
        DataStream<Object> datastream = temp.getDataFromKakfa("7798", 1).myMap(Test::print) ;
        datastream.print() ;

        temp.execute();
    }

public static Object print(Object o) {
        try {
            StringBuilder res = new StringBuilder();
            for (byte b : serializeObject(o)) {
                res.append(String.format("%02X ", b));
                res.append(" "); // delimiter
            }
            return res.toString();
        } catch (NullPointerException e){
            return 0 ;
        } catch (IOException e) {
            return 0;
        }
    }
public static byte[] serializeObject(Object obj) throws IOException
    {
        ByteArrayOutputStream bytesOut = new ByteArrayOutputStream();
        ObjectOutputStream oos = new ObjectOutputStream(bytesOut);
        oos.writeObject(obj);
        oos.flush();
        byte[] bytes = bytesOut.toByteArray();
        bytesOut.close();
        oos.close();
        return bytes;
    }

但我得到了错误:

Exception in thread "main" org.apache.flink.api.common.InvalidProgramException: The implementation of the MapFunction is not serializable. The object probably contains or references non serializable fields.

它指的是 myMap 函数。我该如何解决这个问题?是否有更直接的方法来处理这种情况?

【问题讨论】:

  • 能把完整的代码分享给我们吗?
  • 您能否将“DataStream 数据”声明为静态或将其设为瞬态。我在这里发现了类似的问题:stackoverflow.com/questions/9411292/…
  • @hequn8128 我认为声明一个字段statictransient 使其在集群上运行期间不会序列化。

标签: lambda apache-flink


【解决方案1】:

不看这个没有太多细节,看来你的Function&lt;Object, Object&gt; func需要实现Serializable

您可以通过创建标记界面来完成:

@FunctionalInterface
interface SerializableFuncton<I, O> extends Function<I, O>, Serializable { }

然后将DataStream&lt;Object&gt; myMap(Function&lt;Object, Object&gt; func)更改为DataStream&lt;Object&gt; myMap(SerializableFuncton&lt;Object, Object&gt; func)

【讨论】:

    猜你喜欢
    • 2023-02-05
    • 2021-04-19
    • 1970-01-01
    • 2011-09-27
    • 1970-01-01
    • 1970-01-01
    • 2020-12-23
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多