【发布时间】: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