【发布时间】:2021-04-26 00:21:57
【问题描述】:
我有这样的管道:
env.addSource(kafkaConsumer)
.keyBy { value -> value.f0 }
.window(EventTimeSessionWindows.withGap(Time.minutes(2)))
.reduce(::reduceRecord)
.addSink(kafkaProducer)
我想用 TTL 使键控数据过期。
一些博客文章指出我需要一个ValueStateDescriptor。
我做了一个这样的:
val desc = ValueStateDescriptor("val state", MyKey::class.java)
desc.enableTimeToLive(ttlConfig)
但是我如何将这个描述符实际应用到我的管道中,以便它实际执行 TTL 到期?
【问题讨论】:
标签: kotlin apache-flink flink-streaming stream-processing