【发布时间】:2019-07-21 21:40:24
【问题描述】:
我正在尝试了解 databricks delta 并考虑使用 Kafka 进行 POC。基本上计划是使用来自 Kafka 的数据并将其插入到 databricks 增量表中。
这些是我执行的步骤:
- 在数据块上创建增量表。
%sql
CREATE TABLE hazriq_delta_trial2 (
value STRING
)
USING delta
LOCATION '/delta/hazriq_delta_trial2'
- 使用来自 Kafka 的数据。
import org.apache.spark.sql.types._
val kafkaBrokers = "broker1:port,broker2:port,broker3:port"
val kafkaTopic = "kafkapoc"
val kafka2 = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", kafkaBrokers)
.option("subscribe", kafkaTopic)
.option("startingOffsets", "earliest")
.option("maxOffsetsPerTrigger", 100)
.load()
.select($"value")
.withColumn("Value", $"value".cast(StringType))
.writeStream
.option("checkpointLocation", "/delta/hazriq_delta_trial2/_checkpoints/test")
.table("hazriq_delta_trial2")
但是,当我查询表格时,它是空的。
我可以确认数据即将到来。当我向 Kafka 主题生成消息时,我通过查看图表中的峰值来验证它。
我错过了什么吗?
我需要有关如何将我从 Kafka 获得的数据插入到表中的帮助。
【问题讨论】:
-
只需在将数据发送到kafka之前运行您的代码,一旦代码运行,您需要将数据发送到kafka。
-
可以尝试直接插入Delta表目录而不是插入表吗?例如。
kafka .writeStream .format("delta") .outputMode("append") .option("checkpointLocation", "/delta/events/_checkpoints/etl-from-json") .start("/delta/events") // as a path
标签: scala apache-spark apache-kafka spark-structured-streaming delta-lake