【发布时间】: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