Flink SQL MiniBatch仅触发于Checkpoint问题求助
问题:Flink MiniBatch未按预期触发,仅在Checkpoint时输出
作业配置
'table.exec.sink.upsert-materialize': 'NONE', 'table.exec.mini-batch.enabled': true, 'table.exec.mini-batch.allow-latency': '1 s', 'table.exec.mini-batch.size': '20000', 'table.optimizer.agg-phase-strategy': 'ONE_PHASE', 'execution.checkpointing.interval': env_var('ENV', 'local') == 'local' and '10 s' or '900 s',
已启用Checkpoint,作业逻辑:从两个Kafka主题消费数据,按Key关联后执行聚合操作,最终通过Upsert Kafka连接器输出到Kafka。
官方文档说明
MiniBatch可利用最大延迟时间缓冲输入记录,是一种通过缓冲输入记录减少状态访问的优化手段。当达到允许的延迟间隔或缓冲记录数上限时,MiniBatch会触发。注意:若
table.exec.mini-batch.enabled设为true,其值必须大于0。
问题现象
根据文档描述,MiniBatch应在满足1s延迟或20000条缓冲记录任一条件时触发,但实际运行中,作业运行20多分钟仍未立即输出所有预期记录,仅在Checkpoint触发时才会全量输出(本地环境Checkpoint间隔设为10s,输出速度更快)。
提问
请问仅依赖以下配置,能否让MiniBatch正常工作?
'table.exec.mini-batch.allow-latency': '1 s', 'table.exec.mini-batch.size': '20000',
注:在Flink 1.19.1和2.1.0版本中均遇到此问题。
内容的提问来源于stack exchange,提问作者hitesh
相关产品推荐
相关产品推荐

