【问题标题】:Spark Dataframe write to kafka topic in avro format?Spark Dataframe以avro格式写入kafka主题?
【发布时间】:2018-06-05 16:58:35
【问题描述】:

我在 Spark 中有一个看起来像

的数据框

事件DF

   Sno|UserID|TypeExp
    1|JAS123|MOVIE
    2|ASP123|GAMES
    3|JAS123|CLOTHING
    4|DPS123|MOVIE
    5|DPS123|CLOTHING
    6|ASP123|MEDICAL
    7|JAS123|OTH
    8|POQ133|MEDICAL
    .......
    10000|DPS123|OTH

我需要以 Avro 格式将其写入 Kafka 主题 目前我可以使用以下代码在 Kafka 中以 JSON 格式编写

val kafkaUserDF: DataFrame = eventDF.select(to_json(struct(eventDF.columns.map(column):_*)).alias("value"))
  kafkaUserDF.selectExpr("CAST(value AS STRING)").write.format("kafka")
    .option("kafka.bootstrap.servers", "Host:port")
    .option("topic", "eventdf")
    .save()

现在我想以 Avro 格式将其写入 Kafka 主题

【问题讨论】:

    标签: scala apache-spark dataframe apache-kafka avro


    【解决方案1】:

    火花 >= 2.4

    您可以使用spark-avro 库中的to_avro 函数。

    import org.apache.spark.sql.avro._
    
    eventDF.select(
      to_avro(struct(eventDF.columns.map(column):_*)).alias("value")
    )
    

    火花

    你必须这样做:

    • 创建一个函数,将序列化的 Avro 记录写入ByteArrayOutputStream 并返回结果。一个简单的实现(仅支持平面对象)可能类似于(由Kafka Avro Scala Example 采用Sushil Kumar Singh

      import org.apache.spark.sql.Row
      
      def encode(schema: org.apache.avro.Schema)(row: Row): Array[Byte] = {
        val gr: GenericRecord = new GenericData.Record(schema)
        row.schema.fieldNames.foreach(name => gr.put(name, row.getAs(name)))
      
        val writer = new SpecificDatumWriter[GenericRecord](schema)
        val out = new ByteArrayOutputStream()
        val encoder: BinaryEncoder = EncoderFactory.get().binaryEncoder(out, null)
        writer.write(gr, encoder)
        encoder.flush()
        out.close()
      
        out.toByteArray()
      }
      
    • 将其转换为udf:

      import org.apache.spark.sql.functions.udf
      
      val schema: org.apache.avro.Schema
      val encodeUDF = udf(encode(schema) _)
      
    • 用它代替to_json

      eventDF.select(
        encodeUDF(struct(eventDF.columns.map(column):_*)).alias("value")
      )
      

    【讨论】:

    • 错字:import org.apache.spark.sql.funcitons.udf 请将funcitons 改为functions
    • 不幸的是,Avro 模式在生产中很少是平坦的,这不支持 Kafka 的 Confluent Avro 格式。除此之外,您只需要一个 GenericDatumWriter
    • 从 kafka 消费后如何恢复或反序列化? @OneCricketeer
    • @supernatural 我建议创建一个新帖子
    猜你喜欢
    • 2020-09-05
    • 1970-01-01
    • 2020-10-04
    • 2019-04-26
    • 2019-07-12
    • 2018-03-19
    • 2022-01-14
    • 1970-01-01
    • 2019-09-10
    相关资源
    最近更新 更多