【问题标题】:Intentionally drop state when using suppress for rate limiting updates to KTable使用抑制对 KTable 进行速率限制更新时故意丢弃状态
【发布时间】:2019-12-04 14:50:52
【问题描述】:

我正在使用 Kafka Streams 2.3.1 suppress() 运算符来限制发送到底层 KTable 的更新数量。

这里的用例是,在我的处理逻辑中,我想进行 HTTP 调用,但是为了限制调用次数,我将流窗口化并聚合落入同一时间窗口的源主题消息以进行单个 API 调用。

代码大致如下

KTable<Windowed<String>, List<Event>> windowedEventKTable = inputKStream
    .groupByKey()
    .windowedBy(TimeWindows.of(Duration.ofSeconds(30)).grace(Duration.ofSeconds(5))
    .aggregate(Aggregator::new, ((key, value, aggregate) -> aggregate.aggregate(value)), stateStore)
    .suppress(Suppressed.untilTimeLimit(Duration.ofSeconds(5), maxRecords(500).emitEarlyWhenFull())
    .mapValues((windowedKey, groupedTriggerAggregator) -> {//code here returning a list})
    .toStream((k,v) -> k.key())
    .flatMapValues((readOnlyKey, value) -> value);

我遇到的问题是,当发出超过记录限制的窗口时,状态被保留。在某些时候,单个时间窗口的状态会增长到多个 MB,导致抑制存储更改日志消息超过主题的 max.message.bytes 限制。对于我们的用例,一旦窗口被发射,我们实际上并不关心剩余状态,丢弃它是安全的。

由于我们在多个团队之间共享 Kafka 集群,因此运行集群的团队不愿将集群级别 max.message.bytes 属性提高到超过我们要求的 10 MB。

除了使用transformValues 实现我的逻辑之外,我还有其他选择吗?如果没有,未来是否有任何 Kafka Streams 增强功能能够开箱即用地处理这个问题?

【问题讨论】:

    标签: apache-kafka-streams


    【解决方案1】:

    对于我们的用例,一旦窗口被发射,我们实际上并不关心剩余状态,将其丢弃是安全的。

    对于这种情况,您可以通过aggregation() 参数Materialized.withRetentiontTime(...) 将存储保留时间(默认为1 天)设置为与指定宽限期相同的值。

    我遇到的问题是,当发出超过记录限制的窗口时,状态被保留。在某些时候,单个时间窗口的状态会增长到多个 MB,导致抑制存储更改日志消息超过主题的 max.message.bytes 限制。

    这实际上是一个有趣的陈述,看看你的代码,我只是想澄清一些事情:当你限制时间并允许根据缓存大小提前发出时,你似乎有很多记录在外面即使在发出中间结果之后,也可以进一步更新状态。如果您如上所述通过保留时间清除状态,则需要考虑以下事项:

    • 清除状态不会影响基于缓存大小触发的任何发射,因为只有在保留时间过去后才会清除状态。 0 此外,清除状态意味着清除后出现的所有乱序记录都将被处理,而是将被删除(因为保留时间隐含地将时间戳较小的输入记录标记为“延迟”) .

    但是,总体而言,您似乎并不真正关心乱序数据和事件时间窗口,因为您可以“任意”将记录放入窗口中,因为唯一的目标是减少外部的数量API 调用。因此,您实际上通过使用WallclockTimetampExtractor(而不是默认提取器)切换到处理时间语义似乎是合适的。为了确保每条记录只发出一次,您应该将suppress() 配置更改为只发出“最终”结果。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-09-24
      • 2015-02-20
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多