代码之家  ›  专栏  ›  技术社区  ›  Akash Sethi

ApacheSpark2.4结构流从Kafka读取经过多次转换后始终是二进制结构

  •  -1
  • Akash Sethi  · 技术社区  · 7 年前

    我正在使用apache spark 2.4,在对流式查询应用多个转换之后,我正在从kafka读取json数据最终输出仍然是二进制的。

    val streamingDF = sparkSession.readStream
          .format("kafka")
          .option("subscribe", "test")
          .option("startingOffsets", "latest")
          .option("failOnDataLoss", value = false)
          .option("maxOffsetsPerTrigger", 50000L)
          .option("kafka.bootstrap.servers", "kafka_server")
          .option("enable.auto.commit" , "false")
          .load()
    
    val dataSet = streamingDF.selectExpr("CAST(value AS STRING)").as[String] 
    val stream = dataSet.map{value => convertJSONToCaseClass(value)}
    .map{data => futherconvertions(data)}.writeStream.format("console")
    .outputMode(OutputMode.Update()).start()
    

    在这之后,我在控制台上得到这样的输出。

    Batch: 8
    -------------------------------------------
    +--------------------+
    |               value|
    +--------------------+
    |[01 00 63 6F 6D 2...|
    |[01 00 63 6F 6D 2...|
    |[01 00 63 6F 6D 2...|
    

    预期输出假定为具有多列的数据帧

    我做错什么了吗? 任何帮助都将不胜感激。

    谢谢

    2 回复  |  直到 7 年前
        1
  •  0
  •   jay singh Tanwer    7 年前

    不建议按照文档设置“enable.auto.commit”,请参考卡夫卡的具体配置 https://spark.apache.org/docs/2.4.0/structured-streaming-kafka-integration.html 您也可以按以下方式尝试:

    val streamingDF = sparkSession.readStream
      .format("kafka")
      .option("subscribe", "test")
      .option("startingOffsets", "latest")
      .option("failOnDataLoss", value = false)
      .option("maxOffsetsPerTrigger", 50000L)
      .option("kafka.bootstrap.servers", "kafka_server")
      .load()
    val df = streamingDF.selectExpr("CAST(value as STRING)")         
    
     val mySchema = StructType(Array(
      StructField("X", StringType, true),
      StructField("Y", StringType, true),
      StructField("Z", StringType, true))                            
    
    val Resultdf = df.select(from_json($"value", mySchema).as("data")).select("data.*")
    
        2
  •  0
  •   Kaushal    7 年前

    Spark 2.4不支持多聚合链。

    https://spark.apache.org/docs/2.4.0/structured-streaming-programming- guide.html#unsupported-operations

    多个流聚合(即 流式数据集尚不支持流式数据集。

    推荐文章