【发布时间】:2017-03-28 06:40:19
【问题描述】:
问题来了:
假设有一个数字流,我想在 1 小时的存储桶中收集这些数字中的 MAX,我允许给定存储桶最多延迟 3 小时。
这听起来像是 tumbling windows 的实验室案例。
这是我目前所拥有的:
stream.aggregate(
() -> 0L,
(aggKey, value, aggregate) -> Math.max(value, aggregate),
TimeWindows.of(TimeUnit.HOURS.toMillis(1L)).until(TimeUnit.HOURS.toMillis(3L)),
Serdes.Long(),
"my_store"
)
首先,我无法通过测试验证这是否真的发生。时间戳是通过 TimestampExtractor 提取的,我使用 Thread.sleep 模拟延迟(我将窗口设置为较小的值进行测试),但“迟到的记录”仍会被处理而不是丢弃。
在常规窗口上似乎很少(不是?)示例。有一个关于 SessionWindows 的集成测试,仅此而已。我是否正确理解了这些概念?
编辑 2
示例 JUnit 测试。由于它相当大,我将通过 Gist 分享它。
https://gist.github.com/Hartimer/6018a731753846c1930429716703e5a6
编辑(添加更多代码)
数据点具有时间戳(收集数据的时间)、收集数据的机器的主机名和值。
{
"collectedAt": 12314124134, // timestamp
"hostname": "machine-1",
"reading": 3
}
自定义时间戳提取器用于获取collectedAt。这是我的管道的更完整表示:
source.map(this::fixKey) // Associates record with a key like "<timestamp>:<hostname>"
.groupByKey(Serdes.String(), roundDataSerde)
.aggregate(
() -> RoundData.EMPTY_ROUND,
(aggKey, value, aggregate) -> max(value, aggregate),
TimeWindows.of(TimeUnit.HOURS.toMillis(1L))
.until(TimeUnit.SECONDS.toMillis(1L)), // For testing I allow 1 second delay
roundDataSerde,
"entries_store"
)
.toStream()
.map(this::simpleRoundDataToAggregate) // Associates record with a key like "<timestamp floored to nearest hour>"
.groupByKey(aggregateSerde, aggregateSerde)
.aggregate(
() -> MyAggregate.EMPTY,
(aggKey, value, aggregate) -> aggregate.merge(value), // I know this is not idempotent, that's a WIP
TimeWindows.of(TimeUnit.HOURS.toMillis(1L))
.until(TimeUnit.SECONDS.toMillis(1L)), // For testing I allow 1 second delay
aggregateSerde,
"result_store"
)
.print()
测试的一个sn-p是
Instant roundId = Instant.now().truncatedTo(ChronoUnit.HOURS).minus(9L, ChronoUnit.HOURS);
sendRecord("mytopic", roundId, 3);
sendRecord("mytopic", roundId.plusMillis(15000), 2);
log.info("Waiting a little before sending more usage. (simulating late record)");
Thread.sleep(5000L);
sendRecord("mytopic", roundId.plusMillis(30000), 5);
// Assert stored value is "3".
// It actually is 5 because the last round is accounted for
任何帮助将不胜感激。
【问题讨论】:
-
我怀疑这更多是您的测试设置的影响。你能在这里分享更多的代码吗?例如,您可能会看到正在处理的“迟到的记录”,但在不同的时间窗口中。或者,您编写输入记录的方式(即为记录分配时间戳的位置/时间)可能与您随后的断言/验证不兼容。
-
我编辑了问题以添加更多代码@MichaelG.Noll
-
谢谢哈蒂默。
sendRecord究竟做了什么?它的第二个参数是设置上面显示的 JSON 有效负载中的collectedAt时间戳? (我假设 JSON 有效负载是记录值?) -
@MichaelG.Noll 查看我自己的答案和相关要点
标签: apache-kafka apache-kafka-streams