【发布时间】:2017-12-30 04:34:58
【问题描述】:
我是 Apache spark 的新手,并试图在 scala 中使用 Apache kafka 理解结构化流式传输,但到目前为止没有什么对我有利,基本上我想使用 spark 结构化流式传输从 kafka 进程发送 JSON 并发送回 kafka。我尝试了网站上给出的示例,但它不起作用。
这是我的代码:
import org.apache.spark.sql._
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types.StructType
import org.apache.spark.sql.types._
import org.apache.spark.sql.streaming.{OutputMode, Trigger}
object dataset_kafka {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName("kafka-consumer")
.master("local[*]")
.getOrCreate()
import spark.implicits._
spark.sparkContext.setLogLevel("WARN")
val df = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "172.21.0.187:9093")
.option("subscribe", "test")
.load()
df
.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
.writeStream
.format("kafka")
.trigger(Trigger.ProcessingTime("5 seconds"))
.option("kafka.bootstrap.servers", "172.21.0.187:9093")
.option("topic", "test1")
.option("checkpointLocation", "/home/hduser/Desktop/tempo")
.start()
.awaitTermination()
}
}
对我哪里出错有帮助吗?
我正在以这种格式从 kafka 发送 json:
{"schema":"Hiren","payload":"123"}
【问题讨论】:
-
欢迎来到 SO!请参阅此处,了解如何发布一个好问题,一个可能不会被关闭,甚至可能得到回答的问题:stackoverflow.com/help/how-to-ask
-
我的问题无效吗?
-
你应该展示你自己的一些不工作的代码/你自己的一些努力。您要求的称为教程
-
先生,如您所说,我自己尝试过,但它不起作用,请帮助我纠正我的错误
标签: scala apache-kafka spark-streaming