【问题标题】:spark streaming + querying hive table in every streaming batch?火花流+查询每个流批次中的配置单元表?
【发布时间】:2019-03-14 16:38:52
【问题描述】:

我有一个 Spark Streaming 应用程序从 Flume 接收数据,并在经过一些转换后写入 Hbase。

但要进行这些转换,我需要从配置单元表中查询一些数据。然后问题就开始了。

我不能在转换中使用 sqlContext 或 hiveContext(它们不可序列化),当我在转换之外编写代码时,它只运行一次。

如何使此代码在每个流式批处理中运行?

  def TB_PARAMETRIZACAO_TGC(sqlContext: HiveContext): Map[String,(String,String)] = {
  val df_consulta = sqlContext.sql("SELECT TGC,TIPO,DESCRICAO FROM dl_prepago.TB_PARAMETRIZACAO_TGC")
  val resultado = df_consulta.map(x => x(Consulta_TB_PARAMETRIZACAO_TGC.TGC.id).toString
      -> (x(Consulta_TB_PARAMETRIZACAO_TGC.TIPO.id).toString, x(Consulta_TB_PARAMETRIZACAO_TGC.DESCRICAO.id).toString)).collectAsMap()
  resultado
}

【问题讨论】:

  • 答案对 BTW 有帮助吗?

标签: scala apache-spark spark-streaming


【解决方案1】:

尝试下面这个非常简单的方法,注意静态 JOIN 表可以被缓存并且它们不应该太大,否则静态需要是 KV Store LKP,比如 Hbase:

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.streaming.OutputMode

object StreamJoinStatic {

case class Sales(
  transactionId: String,
  customerId:    String,
  itemId:        String,
  amountPaid:    Double)

case class Customer(customerId: String, customerName: String)

def main(args: Array[String]): Unit = {
 val sparkSession = SparkSession.builder
   .master("local") // Not recommended
   .appName("exampleStaticJoinStrStr")
   .getOrCreate()

//create stream from socket
val socketStreamDf = sparkSession.readStream
  .format("socket")
  .option("host", "localhost")
  .option("port", 50050)
  .load()

import sparkSession.implicits._
//take customer data as static df from where ever
val customerDs = sparkSession.read
  .format("csv")
  .option("header", true)
  .load("src/main/resources/customers.csv")
  .as[Customer]

import sparkSession.implicits._
val dataDf = socketStreamDf.as[String].flatMap(value ? value.split(" "))
val salesDs = dataDf
  .as[String]
  .map(value ? {
    val values = value.split(",")
    Sales(values(0), values(1), values(2), values(3).toDouble)
  })

val joinedDs = salesDs.join(customerDs, "customerId")

val query = joinedDs.writeStream.format("console").outputMode(OutputMode.Append())

query.start().awaitTermination()
}
}

然后适应你的具体情况。

【讨论】:

    猜你喜欢
    • 2016-10-02
    • 1970-01-01
    • 1970-01-01
    • 2016-06-25
    • 1970-01-01
    • 2018-08-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多