【问题标题】:SQL over Spark Streaming基于 Spark 流的 SQL
【发布时间】:2014-10-18 13:01:31
【问题描述】:

这是通过 Spark Streaming 运行简单 SQL 查询的代码。

import org.apache.spark.streaming.{Seconds, StreamingContext}
import org.apache.spark.streaming.StreamingContext._
import org.apache.spark.sql.SQLContext
import org.apache.spark.streaming.Duration

object StreamingSQL {

  case class Persons(name: String, age: Int)

  def main(args: Array[String]) {

    val sparkConf = new SparkConf().setMaster("local").setAppName("HdfsWordCount")
    val sc = new SparkContext(sparkConf)
    // Create the context
    val ssc = new StreamingContext(sc, Seconds(2))

    val lines = ssc.textFileStream("C:/Users/pravesh.jain/Desktop/people/")
    lines.foreachRDD(rdd=>rdd.foreach(println))

    val sqc = new SQLContext(sc);
    import sqc.createSchemaRDD

    // Create the FileInputDStream on the directory and use the
    // stream to count words in new files created

    lines.foreachRDD(rdd=>{
      rdd.map(_.split(",")).map(p => Persons(p(0), p(1).trim.toInt)).registerAsTable("data")
      val teenagers = sqc.sql("SELECT name FROM data WHERE age >= 13 AND age <= 19")
      teenagers.foreach(println)
    })

    ssc.start()
    ssc.awaitTermination()
  }
}

如您所见,要在流上运行 SQL,必须在 foreachRDD 方法中进行查询。 我想对从两个不同流接收的数据运行 SQL 连接。有什么办法可以吗?

【问题讨论】:

    标签: apache-spark spark-streaming


    【解决方案1】:

    好吧,我想总结一下我们在 Spiro 的回答中讨论后得出的解决方法。他建议首先创建一个空表,然后将 RDD 插入其中。唯一的问题是 Spark 还不允许插入到表格中。以下是可以做的:

    首先,创建一个与您期望从流中获得相同架构的 RDD:

    import sqlContext.createSchemaRDD
    val d1=sc.parallelize(Array(("a",10),("b",3))).map(e=>Rec(e._1,e._2))
    

    然后将其保存为 Parquet 文件

    d1.saveAsParquetFile("/home/p1.parquet")
    

    现在,加载 parquet 文件并使用 registerAsTable() 方法将其注册为表格。

    val parquetFile = sqlContext.parquetFile("/home/p1.parquet")
    parquetFile.registerAsTable("data")
    

    现在,当您收到您的流时,只需在您的流上应用 foreachRDD() 并使用 insertInto() 继续在上面创建的表中插入各个 RDD方法

    dStream.foreachRDD(rdd=>{
    rdd.insertInto("data")
    })
    

    此 insertInto() 工作正常,允许将数据收集到表中。现在您可以对任意数量的流执行相同的操作,然后运行查询。

    【讨论】:

      【解决方案2】:

      按照您编写代码的方式,每次运行 SQL 查询时,最终都会生成一系列小 SchemaRDD。诀窍是将这些中的每一个保存到累积 RDD 或累积表中。

      首先,表格方法,使用insertInto

      对于您的每个流,首先创建一个您注册为表的 emty RDD,获取一个空表。对于您的示例,假设您将其称为“allTeenagers”。

      然后,对于您的每个查询,使用 SchemaRDD 的 insertInto 方法将结果添加到该表中:

      teenagers.insertInto("allTeenagers")
      

      如果您对两个流都执行此操作,创建两个单独的累积表,然后您可以使用普通的旧 SQL 查询将它们连接起来。

      (注意:我实际上并没有让他工作,稍微搜索一下让我怀疑其他人有,但我很确定我已经理解insertInto的设计意图,所以我认为这个解决方案值得记录。)

      第二种unionAll 方法(还有一个 union 方法,但这使得获取正确的类型变得更加棘手):

      这涉及到创建一个初始 RDD——我们再次称它为allTeenagers

      // create initial SchemaRDD even if it's empty, so the types work out right
      var allTeenagers = sqc.sql("SELECT ...")
      

      然后,每次:

      val teenagers = sqc.sql("SELECT ...")
      allTeenagers = allTeenagers.unionAll(teenagers)
      

      也许不用说您需要列匹配。

      【讨论】:

      • 感谢您的回复。我尝试了类似var p1 = Person("Hari",22); val rdd1 = sc.parallelize(Array(p1)); rdd1.registerAsTable("data"); var p2 = Person("sagar", 22); var rdd2 = sc.parallelize(Array(p2)); rdd2.insertInto("data"); 并收到错误“java.lang.AssertionError: assertion failed: No plan for InsertIntoTable Map(), false” 似乎我使用 insertInto() 的方式错误?
      • @Pravesh:我也有同样的问题。我很确定它应该可以工作,但是一些搜索让我想知道是否有人使用它。我很好奇您对您在 Spark 列表上发布的问题的答复。我用第二个解决方案更新了我的答案,我很确定它会基于unionAll 工作,我很惊讶没有人提出建议。后者的一个简单示例对我来说很好。
      • 感谢您的宝贵建议。如果你发现新的东西,请让我更新。也会这样做。
      • @Pravesh:出于某种原因,您是否已经排除了我的答案中的第二个解决方案(unionAll)?
      • @Pravesh:我不是建议您将来自不同流的数据收集到单个 RDD 中,而是将 foreachRDD 提供给您的 RDD 片段从每个流收集到一个累加器表中/ RDD,生成两个表或两个 RDD,然后您可以将它们连接起来(每个都包含到目前为止相应流中的所有数据)
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-10-13
      • 2019-06-01
      • 1970-01-01
      • 2022-08-18
      相关资源
      最近更新 更多