【问题标题】:How to expire keyed state with TTL in Apache Flink?如何在 Apache Flink 中使用 TTL 使键控状态过期?
【发布时间】: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


    【解决方案1】:

    您描述的管道不使用任何可以从设置状态 TTL 中受益的键控状态。管道中唯一的键控状态是会话窗口的内容,并且随着会话关闭,该状态将被尽快清除。 (此外,由于您使用的是 reduce 函数,因此该状态仅包含每个键的一个值。)

    在大多数情况下,过期状态仅与您明确创建的状态相关,在这种情况下,您可以随时访问状态描述符并可以将其配置为使用状态 TTL。 Flink SQL 确实会代表您创建可能不会自动过期的状态,在这种情况下,您需要使用 Idle State Retention Time 来配置它。 CEP 库还代表您创建状态,在这种情况下,您应该确保您的模式最终匹配​​或超时。

    【讨论】:

    • 密钥本身呢。窗口关闭后密钥是否被清除?还是永远存在?如果是这样,如果不断向流中添加唯一键怎么办?
    • 密钥不是单独存储的。当键状态被清除时,键/值对被清除,没有任何东西留下。
    • 我明白了。我是否正确:在上面的管道代码中,假设我只得到 1 个事件,键为 x,值为 y。一旦密钥 x 的会话在 2 分钟后关闭,“x”和“y”都会从状态中删除,状态为空。
    • 听起来你明白,但要迂腐:一旦水印到达,表明至少有两分钟的时间间隔,在此期间x 没有事件,xy 从状态中删除。但是,该州将保留任何用作关闭 x 会话的水印证据的事件。
    • 所以我对此进行了测试,似乎这个答案不准确:stackoverflow.com/questions/67429035/…
    猜你喜欢
    • 2023-04-10
    • 1970-01-01
    • 2021-10-18
    • 1970-01-01
    • 1970-01-01
    • 2021-05-07
    • 1970-01-01
    • 1970-01-01
    • 2023-04-05
    相关资源
    最近更新 更多