【问题标题】:Produce Avro topic to Kafka using Apache Spark使用 Apache Spark 为 Kafka 生成 Avro 主题
【发布时间】:2019-09-10 18:10:21
【问题描述】:

我已经在本地安装了 kafka(目前没有集群/模式注册表)并尝试生成 Avro 主题,下面是与该主题关联的模式。

{
  "type" : "record",
  "name" : "Customer",
  "namespace" : "com.example.Customer",
  "doc" : "Class: Customer",
  "fields" : [ {
    "name" : "name",
    "type" : "string",
    "doc" : "Variable: Customer Name"
  }, {
    "name" : "salary",
    "type" : "double",
    "doc" : "Variable: Customer Salary"
  } ]
}

我想创建一个简单的SparkProducerApi,根据上面的架构创建一些数据,然后发布到kafka。 考虑创建示例数据转换为dataframe,然后将其更改为avro,然后发布。

val df = spark.createDataFrame(<<data>>)

然后,如下所示:

df.write
  .format("kafka")
  .option("kafka.bootstrap.servers","localhost:9092")
  .option("topic","customer_avro_topic")
  .save()
}

现在可以通过manually 将架构附加到这个 avro 主题。

这可以通过使用Apache Spark APIs 而不是使用Java/Kafka Apis 来完成吗?这是用于批处理而不是streaming

【问题讨论】:

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


    【解决方案1】:

    我不认为这是直接可能的,因为 Spark 中的 Kafka 生产者需要两列 key 和 value,它们都必须是字节数组。

    如果您从磁盘读取现有的 Avro 文件,您的 Avro 数据帧阅读器可能会为姓名和薪水创建两列。因此,您需要一个操作来从包含整个 Avro 记录的其他列中构造一个 value 列,然后删除这些其他列,然后您必须使用诸如 Bijection 之类的库将其序列化为字节数组,例如,因为您不使用架构注册表。

    如果您想生成数据并且没有文件,那么您需要为 Kafka 消息键和字节数组值构建一个 Tuple2 对象列表,然后您可以 parallelize 这些RDD,然后将它们转换为Dataframe。但到那时,只使用常规的 Kafka Producer API 就简单多了。

    另外,如果你已经知道你的架构,试试Ways to generate test data in Kafka中提到的项目

    【讨论】:

    • 谢谢@@cricket_007。好的,那么假设我已经在 kafka 上拥有了这个 customer-avro-topic 和关联的架构 customer.avsc,那么使用 Spark Consumer Apis 至少可以实现反向(消费)。还看到了自版本2.4.0 以来的一些功能from_avroto_avro(不确定这是否足够稳定以支持所有avro 功能)。如果您说即使这也不是直接可能的,那么这意味着如果我没有错的话,我们不能仅仅依靠 Spark APIs for consuming/producing avro topics from kafka。你这么认为吗?
    • 我相信这些功能适用于 Avro 文件。不是 Kafka 事件,但是是的,您需要将字节反序列化为 Avro 对象
    猜你喜欢
    • 2022-08-04
    • 2019-03-14
    • 2018-06-05
    • 1970-01-01
    • 2021-06-14
    • 1970-01-01
    • 2020-07-12
    • 2021-10-11
    • 2021-08-10
    相关资源
    最近更新 更多