Flink 1.4.1升级至1.13.2作业提交报detached及无sink错误如何解决
作业代码调整方案
- 替换调试类eager操作
原代码中的print()属于触发本地客户端拉取结果的eager执行方法,不适合在detached集群模式下使用,做调试输出时请替换为Flink 1.13提供的正式PrintSink:
// 原写法(仅支持本地测试,detached模式下报错) dataStream.map(/* 你的图像处理逻辑 */).print(); // 替换后写法 dataStream.map(/* 你的图像处理逻辑 */) .sinkTo(PrintSink.<图像处理结果类>builder() .setParallelism(1) .build()) .name("调试输出Sink");
如果不需要输出调试信息,只是需要调用节点本地处理器完成图像处理落地,请将处理器调用逻辑封装到自定义Sink中,既符合Flink的编程规范,也满足Sink校验要求:
// 自定义图像处理Sink public class ImageProcessSink extends RichSinkFunction<图像处理中间结果类> { private transient 本地处理器实例 processor; @Override public void open(Configuration parameters) throws Exception { // 初始化本地处理器 processor = new 本地处理器实例(); } @Override public void invoke(图像处理中间结果类 value, Context context) throws Exception { // 调用节点本地处理器执行逻辑 processor.process(value); } } // 作业中添加Sink dataStream.map(/* 预处理逻辑 */) .addSink(new ImageProcessSink()) .name("图像处理落地Sink");
- 临时兜底方案
如果暂时不想调整业务逻辑位置,仅需要绕过Sink校验,可以添加DiscardingSink作为空Sink兜底:
dataStream.map(/* 原有逻辑 */).addSink(new DiscardingSink<>());
注意该方案仅用作临时验证,生产环境建议将落地逻辑迁移到Sink中,方便后续监控、重试策略配置
版本迁移核心注意事项
- 执行模式校验加严:Flink 1.12版本正式引入流批一体架构后,对detached集群提交模式做了明确的eager操作限制,
collect()、print()、executeAndCollect()这类需要客户端等待返回结果的操作,在detached模式下会直接被拦截报错,该限制是为了避免提交逻辑与执行模式不匹配导致的无意义运行,老版本1.4无此校验所以可以提交,但逻辑本身并不合理。 - Sink强校验规则上线:1.13版本新增了作业DAG合法性校验,要求所有作业必须至少包含一个Sink节点才允许提交,避免用户忘记配置数据落地逻辑导致作业空跑消耗资源,老版本1.4无此校验,所以之前将所有处理逻辑放在
map、process等算子中、未显式声明Sink的作业可以正常提交,新版本会被拦截。 - API兼容性说明:Flink 1.13对旧版
SinkFunction做了向下兼容,原有自定义Sink不需要立刻重构为新的Sink接口即可运行,后续迭代可以逐步迁移。 - 状态兼容限制:Flink 1.4生成的Checkpoint/Savepoint无法直接在1.13版本中恢复,如有状态迁移需求,需要先通过Flink 1.8、Flink 1.10等中间版本做格式转换后再恢复到1.13版本。
- 提交参数变更:1.13版本的
flink run命令部分参数与1.4不兼容,尤其是资源配置、队列绑定相关的参数,提交作业前需要核对参数规则调整提交脚本。
内容的提问来源于stack exchange,提问作者Burcu
相关产品推荐
相关产品推荐

