You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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);

补充注意事项

  1. sinkTo 方法是Flink 1.12及以上版本才提供的API,如果你的项目依赖Flink版本低于1.12,要么升级Flink版本到1.12+,要么改用旧版的org.apache.flink.streaming.api.functions.sink.filesystem.StreamingFileSink。
  2. 你当前的测试代码中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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.04 15:24:00