【问题标题】:Producing Avro type message in spark sql 2.4.4 data frame to Kafka将 spark sql 2.4.4 数据帧中的 Avro 类型消息生成到 Kafka
【发布时间】:2020-05-20 00:35:02
【问题描述】:

我正在尝试使用 spark SQL 将 Avro 消息写入 Kafka。有人可以建议我如何在java中实现它吗?我找到了一个 scala 参考代码,但没有找到 Java。

我试过了,但抛出错误,我在哪里可以配置模式注册表。

aggr.selectExpr("CAST(order_id AS String) AS key", "to_avro(struct(*)) AS value").write().format("kafka").option("kafka.bootstrap.servers", "localhost:9092").option("topic", "aggr_topic").save();

或者请将scala代码复制到java。

val df = spark
  .readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", kafkaURL)
  .option("subscribe", "t")
  .load()
  .select(
    from_avro($"key", "t-key", schemaRegistryURL).as("key"),
    from_avro($"value", "t-value", schemaRegistryURL).as("value"))

提前致谢。

【问题讨论】:

    标签: java apache-spark apache-kafka spark-structured-streaming confluent-schema-registry


    【解决方案1】:

    除了val df之外,该代码在Java中完全相同

    from_avro 只存在于databricks 环境中,顺便说一句,你想要writeStream 和to_avro,无论如何。

    另一种方法是用foreachPartition将dataframe转成RDD,然后手动新建一个KafkaProducer来发送事件

    您可能还对https://github.com/AbsaOSS/ABRiS感兴趣

    【讨论】:

      猜你喜欢
      • 2018-09-01
      • 2022-11-25
      • 1970-01-01
      • 1970-01-01
      • 2018-03-09
      • 2017-12-12
      • 1970-01-01
      • 1970-01-01
      • 2018-04-18
      相关资源
      最近更新 更多