【问题标题】:Is there a way to subtract column 'A' present in Stream1 from column 'B' present in Stream2?有没有办法从 Stream2 中的“B”列中减去 Stream1 中的“A”列?
【发布时间】:2019-04-11 20:43:55
【问题描述】:

我正在使用 Spark 结构化流 (pyspark) 从 Kafka 读取 2 个流(stream1 和 stream2)。我必须计算stream1和stream2的偏移量之间的差异。

我正在尝试这样的事情:

<class 'pyspark.sql.dataframe.DataFrame'>
root
|--timestamp: timestamp (nullable = true)
|-- value: string (nullable = true)
|-- offset: double (nullable = true)
|-- string_val: string (nullable = true)
|-- ping: double (nullable = true)
|-- date: string (nullable = true)
|-- time: string (nullable = true)
|-- offset_v1: double (nullable = true)
|-- date_time: string (nullable = true)
|-- date_format: timestamp (nullable = true)

<class 'pyspark.sql.dataframe.DataFrame'>
|-- Mean: double (nullable = true)
|-- pingTime: timestamp (nullable = true)
|-- Std_Deviation: double (nullable = true)
|-- devTime: timestamp (nullable = true)
|-- offset_v2: double (nullable = true)
|-- upperBound: double (nullable = true)
|-- lowerBound: double (nullable = true)

stream2 = stream2.withColumn('difference',stream2.offset_v2-stream1.offset_v1)

它会抛出一个错误:

pyspark.sql.utils.AnalysisException: u'Resolved attribute(s) offset_v1#95 缺失 upperBound#182,Std_Deviation#149,lowerBound#189,Mean#133,pingTime#129-T30000ms,devTime#144-T30000ms,offset_v2#155 in operator !Project [Mean#133, pingTime#129-T30000ms, Std_Deviation#149,devTime#144-T30000ms,offset_v2#155, 上界#182, 下界#189, (offset_v2#155 - offset_v1#95) AS 差异#233]

【问题讨论】:

  • 您需要进行连接以获得Dataframe,其中包含您计算差异所需的所有相关列
  • 我没有任何列可以用来执行连接操作。不加入,有什么办法吗?
  • 如何匹配两边的行并计算偏移量的差异?那只需要加入。

标签: pyspark apache-kafka pyspark-sql spark-structured-streaming


【解决方案1】:

就像Venki说的,你需要先join才能比较相关的行。你有什么专栏可以做到这一点吗? dateid 可以解决问题。假设您在两个数据框中都有一个名为 join_col 的人:

from pyspark.sql.functions import *

stream_final = stream1.join(stream2, 'join_col', 'inner')

# Now compute difference by adding a new column 'offset_diff':

stream_final = stream_final.withColumn('offset_diff', stream_final.offset_v1 - stream_final.offset_v2)

如果您找不到合适的连接,这对于您比较不同长度的列的情况是个问题,我相信这就是您正在处理的问题。

【讨论】:

  • 我没有任何列使用/可以执行连接操作。不加入有什么办法吗?
猜你喜欢
  • 2021-06-27
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多