【发布时间】: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 我认为声明一个字段
static或transient使其在集群上运行期间不会序列化。
标签: lambda apache-flink