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

Flink 1.16搭配Flink-Doris-Connector 1.14写入Doris性能慢求助

问题定位与优化方案

一、优先修正Connector版本适配

你当前用Flink 1.16搭配Doris Connector 1.14属于跨版本混用,官方Connector的版本和Flink版本强绑定:1.14版本的Connector仅适配Flink 1.14.x,和1.16的Checkpoint、状态管理逻辑不兼容,会导致批量提交逻辑异常,引发写入延迟。

  • 解决:替换为适配Flink 1.16的Connector版本,比如doris-flink-1.2.0(官方适配Flink 1.15-1.17),Maven依赖调整为:
<dependency>
    <groupId>org.apache.doris</groupId>
    <artifactId>doris-flink-1.2.0</artifactId>
    <version>1.2.0</version>
</dependency>

二、调优批量写入核心参数

Flink 1.16的默认参数和1.13差异较大,可能导致单批次数据量过小、提交触发过慢:

  • 调整以下Connector关键参数:
    • sink.batch.size:设为10000(3万条数据分3批次提交,匹配Doris的处理能力)
    • sink.batch.interval:设为1000(毫秒,超时强制提交,避免数据堆积)
    • sink.max-retries:设为2(减少不必要的重试开销)
  • 示例配置代码:
DorisSink.Builder<String> sinkBuilder = DorisSink.builder();
sinkBuilder.setDorisReadOptions(DorisReadOptions.builder().build())
           .setDorisExecutionOptions(DorisExecutionOptions.builder()
               .setBatchSize(10000)
               .setBatchIntervalMs(1000)
               .setMaxRetries(2)
               .build())
           .setProperties(properties);

Flink 1.16的资源配置、并行度设置可能拖慢写入速度:

  • 调整并行度:设置Job并行度与Doris BE节点数一致,或设为3-5(3万条数据无需过高并行度,避免资源分散)
  • 增加TaskManager内存:设置taskmanager.memory.process.size: 4g,避免GC频繁导致的延迟
  • 测试阶段禁用Checkpoint:如果是离线批量写入,临时关闭Checkpoint排除状态持久化开销:
env.disableCheckpointing();

四、排查Connector写入逻辑与Doris状态

  • 开启异步写入:适配1.16的Connector支持异步写入,添加配置enable.async.write: true,提升并发写入能力
  • 检查Doris节点状态:查看BE节点的be.log,确认是否存在磁盘IO过高、副本同步阻塞的情况,这些会直接拖慢写入速度

五、验证数据写入方式

  • 确认是批量写入:检查代码是否误将单条数据单独提交,必须通过DorisSink的内置批量逻辑处理,不要手动循环调用写入接口
  • 隔离上游数据源:构造3万条数据的List,通过env.fromCollection()生成DataStream直接写入Doris,排除上游数据生成的延迟影响

内容的提问来源于stack exchange,提问作者saiheihua

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 08:23:33