如何使用from_json与Kafka connect 0.10和Spark Structured Streaming?

car*_*ues 10 scala apache-kafka apache-spark apache-kafka-connect spark-structured-streaming

我试图重现[Databricks] [1]中的示例并将其应用于Kafka的新连接器并激发结构化流媒体,但我无法使用Spark中的开箱即用方法正确解析JSON ...

注意:该主题以JSON格式写入Kafka.

val ds1 = spark
          .readStream
          .format("kafka")
          .option("kafka.bootstrap.servers", IP + ":9092")
          .option("zookeeper.connect", IP + ":2181")
          .option("subscribe", TOPIC)
          .option("startingOffsets", "earliest")
          .option("max.poll.records", 10)
          .option("failOnDataLoss", false)
          .load()
Run Code Online (Sandbox Code Playgroud)

以下代码不起作用,我相信这是因为列json是一个字符串而且与from_json签名方法不匹配...

    val df = ds1.select($"value" cast "string" as "json")
                .select(from_json("json") as "data")
                .select("data.*")
Run Code Online (Sandbox Code Playgroud)

有小费吗?

[更新]工作示例:https: //github.com/katsou55/kafka-spark-structured-streaming-example/blob/master/src/main/scala-2.11/Main.scala

aba*_*hel 22

首先,您需要为JSON消息定义架构.例如

val schema = new StructType()
  .add($"id".string)
  .add($"name".string)
Run Code Online (Sandbox Code Playgroud)

现在,您可以在from_json下面的方法中使用此架构.

val df = ds1.select($"value" cast "string" as "json")
            .select(from_json($"json", schema) as "data")
            .select("data.*")
Run Code Online (Sandbox Code Playgroud)

  • 如果你有编译器警告"价值$不是会员..."请不要忘记导入spark.implicits._我需要额外的5-10分钟才能搞清楚 (2认同)