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 Job运行参数
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
相关产品推荐
相关产品推荐

