【问题标题】:How do I read binary serialized Avro (Confluent Platform) from Kafka using Spark Streaming如何使用 Spark Streaming 从 Kafka 读取二进制序列化的 Avro(Confluent 平台)
【发布时间】:2017-04-26 14:57:40
【问题描述】:

这些是使用 Confluent 平台序列化的 Avros。

我想找到一个这样的工作示例:

https://github.com/seanpquig/confluent-platform-spark-streaming/blob/master/src/main/scala/example/StreamingJob.scala

但对于 Spark 结构化流。

 kafka
   .select("value")
   .map { row => 

     // this gives me test == testRehydrated    
     val test = Foo("bar") 
     val testBytes = AvroWriter[Foo].toBytes(test)
     val testRehydrated = AvroReader[Foo].fromBytes(testBytes)


     // this yields mangled Foo data
     val bytes = row.getAs[Array[Byte]]("value") 
     val rehydrated = AvroReader[Foo].fromBytes(bytes)

【问题讨论】:

  • 您找到可行的解决方案了吗?
  • @aasthetic 见下文

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


【解决方案1】:

我们一直在开发这个库,这可能会有所帮助:ABRiS (Avro Bridge for Spark)

它提供了用于在读取和写入操作(流式传输和批处理)中将 Spark 集成到 Avro 的 API。它还支持 Confluent Kafka 并与 Schema Registry 集成。

免责声明:我为 ABSA 工作,我是这个库背后的主要开发人员。

【讨论】:

  • 在将 ABRiS 依赖项添加到 POM 文件后,当我使用 spark java api 尝试 ABRiS 时,方法 fromAvro 不存在。有什么帮助吗?
  • @Vignesh,该库是用 Scala 编写的,但应该可以导入到 Java 类中。您如何尝试导入它?
  • 在 pom ie(confluent 和 abris) 中添加所有必要的依赖项后,我尝试按照github.com/AbsaOSS/ABRiS 中的说明导入 ABRiS 库,但在Dataset<Row> ds = sparksession.readstream().format("kafka").option("kafka.bootstrap.server","localhost:9092").option("subscribe","topicname"). method fromavro 中没有eclipse intellisense,构建没有问题
  • fromConfluentAvro 即使在所有导入之后也不存在,但在 scala 中它工作正常。
  • fromConfluentAvro 方法在我们使用 java 时不会在 eclipse intelisense 中弹出,但在 scala 中它工作正常。能否请您对使用 java 的 ABRiS 有所了解?
【解决方案2】:

我发现如果你想阅读他们的东西,你必须使用 Confluent 平台解码器。

def decoder: io.confluent.kafka.serializers.KafkaAvroDecoder = {
  val props = new Properties()
  props.put(AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, getSchemaRegistryUrl())
  props.put(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, "true")
  val vProps = new kafka.utils.VerifiableProperties(props)
  new io.confluent.kafka.serializers.KafkaAvroDecoder(vProps)
}

【讨论】:

  • 感谢@zzztimbo,我还为我的密钥注册了模式,它们必须是 Long 类型。我无法在 Spark Streaming 中反序列化它们。有什么想法吗?
  • 获得一个完整的例子来说明如何使用 kafka 主题 + 模式注册表设置结构化流式传输将很有帮助。我不清楚如何处理这个解码器。
猜你喜欢
  • 2019-07-30
  • 2019-11-18
  • 2020-10-15
  • 2019-04-04
  • 1970-01-01
  • 2022-11-24
  • 2017-04-04
  • 1970-01-01
  • 2015-08-01
相关资源
最近更新 更多