【问题标题】:How to convert Spark Streaming data into Spark DataFrame如何将 Spark Streaming 数据转换为 Spark DataFrame
【发布时间】:2016-05-16 16:46:53
【问题描述】:

目前Spark还没有为流数据创建DataFrame,但是我在做异常检测的时候,使用DataFrame进行数据分析更加方便快捷。我已经完成了这部分,但是当我尝试使用流数据进行实时异常检测时,问题就出现了。试了好几种方法,仍然无法将DStream转为DataFrame,也无法将DStream内部的RDD转为DataFrame。

这是我最新版本代码的一部分:

import sys
import re

from pyspark import SparkContext
from pyspark.sql.context import SQLContext
from pyspark.sql import Row
from pyspark.streaming import StreamingContext
from pyspark.mllib.clustering import KMeans, KMeansModel, StreamingKMeans
from pyspark.sql.functions import *
from pyspark.sql.types import *
from pyspark.sql.functions import udf
import operator


sc = SparkContext(appName="test")
ssc = StreamingContext(sc, 5)
sqlContext = SQLContext(sc)

model_inputs = sys.argv[1]

def streamrdd_to_df(srdd):
    sdf = sqlContext.createDataFrame(srdd)
    sdf.show(n=2, truncate=False)
    return sdf

def main():
    indata = ssc.socketTextStream(sys.argv[2], int(sys.argv[3]))
    inrdd = indata.map(lambda r: get_tuple(r))
    Features = Row('rawFeatures')
    features_rdd = inrdd.map(lambda r: Features(r))
    features_rdd.pprint(num=3)
    streaming_df = features_rdd.flatMap(streamrdd_to_df)

    ssc.start()
    ssc.awaitTermination()

if __name__ == "__main__":
    main()

正如您在 main() 函数中看到的,当我使用 ssc.socketTextStream() 方法读取输入流数据时,它会生成 DStream,然后我尝试将 DStream 中的每个个体转换为 Row,希望我可以转换稍后将数据放入 DataFrame。

如果我在这里使用 ppprint() 打印 features_rdd,它可以工作,这让我想, features_rdd 中的每个个体都是一批 RDD,而整个 features_rdd 是一个 DStream。

然后我创建了streamrdd_to_df()方法并希望将每批RDD转换为数据帧,它给了我错误,显示:

ERROR StreamingContext:启动上下文时出错,将其标记为已停止 java.lang.IllegalArgumentException: 要求失败:没有注册输出操作,所以没有要执行的操作

有没有想过如何对 Spark 流数据进行 DataFrame 操作?

【问题讨论】:

    标签: python pyspark spark-streaming


    【解决方案1】:

    仔细阅读错误。它说没有注册输出操作。 Spark 是懒惰的,只有当它有东西要产生时才执行作业/鳕鱼。在您的程序中没有“输出操作”,Spark 也抱怨同样的问题。

    在 DataFrame 上定义一个 foreach() 或原始 SQL 查询,然后打印结果。它会正常工作的。

    【讨论】:

    • 谢谢@Sumit,我确实尝试输出结果。在 streamrdd_to_df() 函数中,我使用sdf.show(n=2, truncate=False) 打印出结果,但它不能....
    • 我执行了代码,它对我有用。我所做的唯一更改是将get_tuple 替换为tuple 并删除pyspark.mllib... 导入。虽然输出没有意义,但没有错误。错误中还有其他内容吗?你能粘贴完整的堆栈跟踪吗?
    【解决方案2】:

    Spark 为我们提供了结构化流媒体,可以解决此类问题。它可以生成流数据帧,即连续附加的数据帧。请查看以下链接

    http://spark.apache.org/docs/latest/structured-streaming-programming-guide.html

    【讨论】:

    • 我刚读到这个,他们在 Spark 2.0 中添加了它,它可能是这个问题的解决方案。我使用的是 Spark 1.5
    • 现在相比它的推出有了很大的改进。
    【解决方案3】:

    你为什么不使用这样的东西:

    def socket_streamer(sc): # retruns a streamed dataframe
        streamer = session.readStream\
            .format("socket") \
            .option("host", "localhost") \
            .option("port", 9999) \
            .load()
        return streamer
    

    上面这个函数的输出本身(或一般的readStream)是一个DataFrame。在那里你不需要担心 df,它已经由 spark 自动创建。 见Spark Structured Streaming Programming Guide

    【讨论】:

      【解决方案4】:

      1 年后,我开始探索 Spark 2.0 流方法,终于解决了我的异常检测问题。 Here's my code in IPython,也可以找how does my raw data input look like

      【讨论】:

        【解决方案5】:

        借助 Spark 2.3 / Python 3 / Scala 2.11(使用数据块),我能够在 scala 中使用临时表和代码 sn-p(在笔记本中使用魔法):

        Python部分:

        ddf.createOrReplaceTempView("TempItems")
        

        然后在一个新的单元格上:

        %scala
        import java.sql.DriverManager
        import org.apache.spark.sql.ForeachWriter
        
        // Create the query to be persisted...
        val tempItemsDF = spark.sql("SELECT field1, field2, field3 FROM TempItems")
        
        val itemsQuery = tempItemsDF.writeStream.foreach(new ForeachWriter[Row] 
        {      
          def open(partitionId: Long, version: Long):Boolean = {
            // Initializing DB connection / etc...
          }
        
          def process(value: Row): Unit = {
            val field1 = value(0)
            val field2 = value(1)
            val field3 = value(2)
        
            // Processing values ...
          }
        
          def close(errorOrNull: Throwable): Unit = {
            // Closing connections etc...
          }
        })
        
        val streamingQuery = itemsQuery.start()
        

        【讨论】:

          【解决方案6】:

          无需将 DStream 转换为 RDD。根据定义,DStream 是 RDD 的集合。只需使用 DStream 的方法 foreach() 循环遍历每个 RDD 并采取行动。

          val conf = new SparkConf()
            .setAppName("Sample")
          val spark = SparkSession.builder.config(conf).getOrCreate()
          sampleStream.foreachRDD(rdd => {
              val sampleDataFrame = spark.read.json(rdd)
          }
          

          【讨论】:

            【解决方案7】:

            spark documentation 介绍了如何使用 DStream。基本上,您必须在流对象上使用foreachRDD 才能与之交互。

            这是一个示例(确保您创建了一个 Spark 会话对象):

            def process_stream(record, spark):
                if not record.isEmpty():
                    df = spark.createDataFrame(record) 
                    df.show()
            
            
            def main():
                sc = SparkContext(appName="PysparkStreaming")
                spark = SparkSession(sc)
                ssc = StreamingContext(sc, 5)
                dstream = ssc.textFileStream(folder_path)
                transformed_dstream = # transformations
            
                transformed_dstream.foreachRDD(lambda rdd: process_stream(rdd, spark))
                #                   ^^^^^^^^^^
                ssc.start()
                ssc.awaitTermination()
            

            【讨论】:

              猜你喜欢
              • 1970-01-01
              • 2019-05-14
              • 1970-01-01
              • 1970-01-01
              • 1970-01-01
              • 2021-03-26
              • 2017-12-31
              • 2017-03-17
              • 2016-05-12
              相关资源
              最近更新 更多