【问题标题】:Structured Streaming extract most recent values for each id结构化流提取每个 id 的最新值
【发布时间】:2019-02-10 21:44:03
【问题描述】:

我有包含 ID、类型和值的数据流:对于具有给定 ID 的一组用户,我从不同的传感器(输入)。 传入数据示例:

ID type value
1  A    70
2  B    16
1  A    71
2  A    72

我需要创建 Spark Structured Streaming 应用程序,该应用程序将对获取的数据执行自定义集群。但是,我一开始就卡住了>我不知道如何创建一组数据,其中将包含每种类型的每个用户的最后测量值。我需要为系统中出现过的每个用户设置此设置。

因此,基本上,对于上述数据流,我需要一个结构化流式应用程序,它可以为我提供一组针对每种类型的每个用户的最后测量值>

  ID type value
  1  A    71
  2  B    16
  2  A    72

用户可能有一段时间不活跃,我仍然需要保留他们的记录。如果输出是数据框,这将很有用。

任何关于如何做到这一点的想法都将受到欢迎。

PS 我对 Spark Structured Streaming 还很陌生,如果这是一个微不足道的问题,我很抱歉。

【问题讨论】:

  • 请确定所需的输出模式。
  • 我认为,鉴于问题,我需要一个完整的输出,
  • 如何才能真正看到最后的测量值是多少?您肯定需要某种时间戳吗?
  • 我可以为每条记录添加时间戳...
  • 无论如何它都行不通。我尝试了所有技巧,看看那些链接,我只检查了 2 个。实际上不现实。

标签: apache-spark dataframe spark-structured-streaming


【解决方案1】:

简短的回答是:这不可能使用 Spark 结构化流式处理(目前)。

许多关于此的帖子都没有提出实际可行的解决方案。

仔细想想,实际上这是一项艰巨的任务。

我尝试了各种方法 - 即使我知道这是不可能的 - 并且总是从 Spark 中得到某种错误。这些都记录在 Stack Overflow 上。例如:

Structured streaming custom deduplication

Retain last row for given key in spark structured streaming

【讨论】:

  • 你会建议我放弃结构化流并尝试使用 RDD 解决这个问题吗?或者有更聪明的解决方案:)
  • 您可以写入数据存储并读回数据帧和进程,因此很简单。但不是流式用例。
  • 我不应该通过流式用例。
猜你喜欢
  • 2012-05-10
  • 2022-08-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-10-02
  • 1970-01-01
  • 2021-10-21
相关资源
最近更新 更多