【发布时间】: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 增强功能能够开箱即用地处理这个问题?
【问题讨论】: