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