【问题标题】:How to actually discard late records?如何实际丢弃迟到的记录?
【发布时间】: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


【解决方案1】:

认为 Hartimer 的自我回答实际上是不正确的。让我试着解释一下会发生什么,至少据我自己所知。 :-)

  • 迟到的数据是根据通过配置的时间戳提取器为您的应用程序配置的时间语义处理的。在@Hartimer 的情况下,这是事件时间(此处使用自定义时间戳提取器)。
  • FWIW,在 处理时间 的情况下,根据定义,没有迟到的记录:每个记录都“及时”到达。 “迟到”的记录(同样,在此上下文中没有此类记录)包含在当前窗口中,但从不回填到较早的窗口中。
  • 设置窗口保留时间的调用TimeWindows#until() 是保留时间的下限。 Kafka 可能会在窗口周围保留“更长一点”(我在这里故意模糊,见下文),而不是配置的保留时间。因此,@Hartimer 进行的严格测试可能不会产生直观预期的结果。

关于窗口保留时间是一个下限的幕后实际发生的事情有点棘手(并且可能超出了这个问题的范围),所以我暂缓解释除非我有特殊要求。

更新:此外,问题的 sn-p 中的这段代码甚至不应该工作,因为它应该抛出一个 IllegalArgumentException

TimeWindows.of(TimeUnit.HOURS.toMillis(1L))
           .until(TimeUnit.SECONDS.toMillis(1L))

要求是,对于它们各自的输入参数,until() &gt;= of()。您不能定义大小为 1 小时但保留期仅为 1 秒的窗口(此处的保留时间必须 >= 1 小时)。

更新 2: 幕后发生的事情是 TimeWindows#until() 的设置用于创建/管理本地窗口存储的段文件。只要窗口段存在,该窗口的迟到记录就会被接受。我将跳过有关如何删除/过期的部分,因为我真的需要深入研究代码(我不知道)。

【讨论】:

  • 感谢您的回复! Gist 共享中的单元测试不会抛出IllegalArgumentException,我相信如果指定了advanceByadvanceBy() &lt; of(),就会抛出它。我假设until() 部分是窗口大小加上until() 作为下限。所述测试是可运行的,如果这里有人可以检查出来,那就太好了。我可能将“处理时间”与“摄取时间”混淆了。无论如何,随着我在回答中提到的更改,测试按预期工作,我担心它可能会出于其他原因这样做......进一步解释会很好
  • 也许 IAE 只在最新的 Kafka 0.10.2 中抛出——也许您使用的是早期版本?
  • 我正在使用0.10.1.1
  • 关于你的最后一次更新,我对“下限”的概念没有任何问题,我想了解更多关于窗口管理的细节,如果我应该为此创建一个单独的问题,请告诉我.我仍然需要一种方法来可靠地测试“迟到的记录”及其处理......
  • 好吧,您可以创建一个单独的问题,但请记住,我们在这里讨论的是 Kafka 的 Streams API 的内部行为,您不应在代码/测试中依赖它。合同是“下限”,这是您应该使用的(恕我直言)。
【解决方案2】:

我相信我发现了自己的问题。它归结为TimestampExtractor 以及我用来评估“延迟记录”的值。

在 Kafka Stream 术语中,有三个“时间”(see here):

  • 事件时间:记录数据的时间
  • 处理时间:流处理器接收数据的时间
  • 摄取时间:(与问题无关)

在我的示例中,我实际上是使用 Event-time 来确定是否有延迟,但这并不代表延迟记录。收集数据的人都会将此值设置为他们当地的时间感知(至少在我的用例中)。

重要的日期是处理时间。无论何时生成,我们需要多长时间才能接收到该事件。我的聚合已经按“事件时间”处理分组。

我创建了一个新的 Gist,其中包含现在通过的测试的更新版本。添加了一个额外的字段receivedAt,模拟“处理时间”。

https://gist.github.com/Hartimer/c79569ad517ab95d08dbe8e84bfa6789

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-01-15
    • 1970-01-01
    • 2019-06-05
    • 1970-01-01
    • 1970-01-01
    • 2010-11-26
    • 2019-09-10
    • 1970-01-01
    相关资源
    最近更新 更多