【问题标题】:How to do LAG like implementation in KSQLDB?如何在 KSQLDB 中进行 LAG 类实现?
【发布时间】:2020-09-02 18:25:47
【问题描述】:

我最近开始使用 ksql 并想检查是否有人可以帮助我进行查询设计。问题陈述是我有一个视频会议应用程序,广播公司可以在其中多次启动和暂停流。我想获得该流的总播放时间和总暂停时间。我有一个由开始和暂停时间戳组成的点击流数据。我应该怎么做才能生成优化的视图。

非常感谢任何帮助:)

谢谢

【问题讨论】:

  • 您好 Abhimanyu,如果您提供更多信息,通常会从社区成员那里获得更多时间,尤其是表明您已经花费了一些时间进行研究。

标签: sql apache-kafka confluent-platform ksqldb


【解决方案1】:

分组事件

您需要解决的第一个问题是如何将开始/停止事件组合在一起?

您可能希望通过某种USER_ID 或其他唯一标识正在启动/停止流的广播公司的属性对它们进行分组。

您可能还希望按某种STREAM_ID 或其他唯一标识正在播放的流的属性进行分组。

这可能就足够了,如果您只需要每个广播公司、每个视频的总播放时间。但是,您可能还需要考虑时间。例如,如果我今天看了一个视频,然后明天再看,那是两个观看会话,有两个独立的观看总时间,还是你不在乎?

按时间对事件进行分组的一种方法是使用会话窗口。在对数据进行会话之前,您需要定义定义会话的参数。这是good example of using session windows in ksqlDB

按时间分组事件的另一种方法是使用翻转窗口。这是good example of using tumbling windows

计算播放时间

将活动分组后,您可能需要计算播放时间。例如,如果我在时间 5 开始播放,在时间 8 停止播放,那么我观看视频的时间为5 - 8 = 3

这需要捕获播放事件并等待停止事件,然后输出时间差。并以容错的方式做一些事情。

在撰写本文时,这需要自定义 UDAF(自定义用户定义的聚合函数)。

自定义 UDAF 可以捕获开始事件,将其存储以供将来参考,并输出“0”作为播放时间,然后当它看到相应的停止事件时,它可以将开始事件从其状态中移除,计算播放时间并返回。

这是一个good example of writing a custom UDF in ksqlDB,尽管您需要一个自定义的 UDAF,它已被 here 覆盖。

目前有一个PR open with an enhancement to the LATEST_BY_OFFSET method 可以很好地满足您的目的。这增强了该方法以允许它捕获最后一个 N 值,而不仅仅是最后一个 1 值。很可能,这将在 ksqlDB v0.13 中发布,如果您有任何开发经验,您可以随时拉取代码并在本地编译。如果它不能满足您的目的,那么您可以将其作为开发自己的起点。

当然,这些解决方案要求您的源事件流正确排序,以便停止事件永远不会出现在它们相关的播放事件之前

聚合

计算完一对开始/停止事件之间的播放时间后,您需要将它们聚合起来。这是good example of how to aggregate in ksqlDB

【讨论】:

  • 您好安德鲁感谢您的回复。过去两个月我一直在研究 ksql,我没有太多的开发背景,尤其是 java,我来自统计背景。因此,我想重申一下我一直坚持的地方-:假设广播公司在下午 12 点开始直播,在 12:05 暂停,然后在 12:08 重新开始并在 12:10 结束。这里总的播放时间是 7 分钟,3 分钟是暂停时间。
  • 为了获得这个数字,我一直在尝试创建两个表,即暂停和播放使用 1 秒的窗口翻滚捕获此类事件的所有时间戳,然后使用两者中都存在的计数器减去它们.但我想检查它是否可以大规模扩展,我预计每天大约有 20-30k 次实时会话。
  • 我不确定我是否遵循。考虑更新您的问题以包含您一直在尝试的 SQL、示例输入数据和预期/所需的输出。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2013-07-15
  • 1970-01-01
  • 1970-01-01
  • 2021-12-25
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多