在写flink代码做实时处理的时候,对于延迟的数据,我们添加了乱序时间处理以及允许迟到数据之后,依旧有迟到的数据,可以采用侧输出流进行收集依旧迟到的数据。都写进关系型数据库(如MySQL可以处理,以此来保证数据的不丢)。
但是笔者在写侧输出流的时候,发现执行报错。Caused by: org.apache.flink.api.common.functions.InvalidTypesException: The types of the interface org.apache.flink.util.OutputTag could not be inferred. Support for synthetic interfaces, lambdas, and generic or raw types is limited at this point
定位到是因为侧输出流的问题,就点开源码进行查看。发现侧输出流需要传递的总是一个匿名实现类,笔者的代码里面传递的是一个对象,所以导致报错。源码如下:
所以:
OutputTag> info = new OutputTag >("late-data"){};
在new的时候需要加上大括号,作为匿名实现类进行传递。
欢迎分享,转载请注明来源:内存溢出
评论列表(0条)