【发布时间】:2018-04-19 12:08:04
【问题描述】:
我们正在对从 MySQL 收集的 kafka 数据进行流式处理。现在,一旦完成所有分析,我想将我的数据直接保存到 Hbase。我通过了 spark 结构化流文档,但找不到任何带有 Hbase 的接收器。我用来从 Kafka 读取数据的代码如下。
val records = spark.readStream.format("kafka").option("subscribe", "kaapociot").option("kafka.bootstrap.servers", "XX.XX.XX.XX:6667").option("startingOffsets", "earliest").load
val jsonschema = StructType(Seq(StructField("header", StringType, true),StructField("event", StringType, true)))
val uschema = StructType(Seq(
StructField("MeterNumber", StringType, true),
StructField("Utility", StringType, true),
StructField("VendorServiceNumber", StringType, true),
StructField("VendorName", StringType, true),
StructField("SiteNumber", StringType, true),
StructField("SiteName", StringType, true),
StructField("Location", StringType, true),
StructField("timestamp", LongType, true),
StructField("power", DoubleType, true)
))
val DF_Hbase = records.selectExpr("cast (value as string) as Json").select(from_json($"json",schema=jsonschema).as("data")).select("data.event").select(from_json($"event", uschema).as("mykafkadata")).select("mykafkadata.*")
现在最后,我想将 DF_Hbase 数据帧保存在 hbase 中。
【问题讨论】:
标签: scala apache-spark apache-kafka hbase spark-streaming