Apache Flink使用Avro SpecificRecord的FileSink编译类型不匹配问题
问题原因
org.apache.flink.connector.file.sink.FileSink 是Flink推出的新一代文件写入Sink,它实现的是 org.apache.flink.api.connector.sink.Sink 接口,而 DataStream.addSink() 方法仅接收实现了旧版 org.apache.flink.streaming.api.functions.sink.SinkFunction 接口的实例,二者类型不匹配,因此编译报错。
修复方案
只需要将代码中添加Sink的行修改为适配新Sink接口的调用方式即可:
将
source.addSink( sink);
替换为
source.sinkTo(sink);
补充注意事项
sinkTo方法是Flink 1.12及以上版本才提供的API,如果你的项目依赖Flink版本低于1.12,要么升级Flink版本到1.12+,要么改用旧版的org.apache.flink.streaming.api.functions.sink.filesystem.StreamingFileSink。- 你当前的测试代码中
getUser方法仅给id字段赋值,没有给name字段赋值,而你定义的Avro Schema中name字段的类型是必填string,没有允许为null,运行时可能会触发序列化错误,建议补充name字段的赋值:
public static User getUser() { User u = new User(); u.setId(1L); u.setName("test"); return u; }
内容的提问来源于stack exchange,提问作者user51
相关产品推荐
相关产品推荐

