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

Flink去重:无界流中out.collect()与TTL的工作机制及输出验证

无界流场景下的TTL与输出逻辑解析

一、TTL与out.collect()的配合工作机制

在无界流场景中,这段代码的TTL逻辑和输出逻辑是联动的,具体流程如下:

  • 首次处理新Key:当某个(ID, subid)组合的元素第一次到达时,seen.value()返回null,此时会执行out.collect(value)将该元素输出,随后调用seen.update(true)创建状态并启动15秒的TTL计时(因为TTL配置的更新类型是OnCreateAndWrite,只有创建或写入状态时才刷新TTL)。
  • TTL有效期内重复Key:在15秒的TTL有效期内,相同(ID, subid)的元素再次到达时,seen.value()会返回true,不会执行out.collect(),也不会更新状态(没有触发写操作),TTL会继续按原计时倒计时。
  • TTL过期后再次处理同Key:15秒TTL到期后,该Key对应的状态会被后台异步清理,此时相同(ID, subid)的新元素到达时,seen.value()又会返回null,会再次执行out.collect()输出元素,同时重新创建状态并启动新一轮的TTL计时。

注意:TTL的清理是后台异步执行的,不会精确到秒立即删除过期状态,但当状态过期后,Flink会将其视为不存在,后续访问会返回null。

二、flatMap的输出是否会传递到print()算子

是的,flatMap方法会将去重后的流输出至print()算子打印。
代码中明确调用了dataStream.keyBy(...).flatMap(new FilterDuplicate()).print(),flatMap里的out.collect(value)就是将符合条件的元素发送到下游算子,这里的下游就是print(),所以去重后的元素会被正常打印。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 23:25:56