【发布时间】:2020-12-23 14:11:51
【问题描述】:
我已经创建了一个增量表,现在我正在尝试使用 foreachBatch() 将数据插入到该表中。我关注了这个example。唯一的区别是我使用的是 Java 而不是笔记本,但我想这应该没有什么区别?
我的代码如下:
spark.sql("CREATE TABLE IF NOT EXISTS example_src_table(id int, load_date timestamp) USING DELTA LOCATION '/mnt/delta/events/example_src_table'");
Dataset<Row> exampleDF = spark.sql("SELECT e.id as id, e.load_date as load_date FROM example e");
try {
exampleDF
.writeStream()
.format("delta")
.foreachBatch((dataset, batchId) -> {
dataset.persist();
// Set the dataframe to view name
dataset.createOrReplaceTempView("updates");
// Use the view name to apply MERGE
// NOTE: You have to use the SparkSession that has been used to define the `updates` dataframe
dataset.sparkSession().sql("MERGE INTO example_src_table e" +
" USING updates u" +
" ON e.id = u.id" +
" WHEN NOT MATCHED THEN INSERT (e.id, e.load_date) VALUES (u.id, u.load_date)");
})
.outputMode("update")
.option("checkpointLocation", "/mnt/delta/events/_checkpoints/example_src_table")
.start();
} catch (TimeoutException e) {
e.printStackTrace();
}
这段代码运行没有任何问题,但没有数据写入 url '/mnt/delta/events/example_src_table' 的 delta 表。有谁知道我做错了什么?
我使用的是 Spark 3.0 和 Java 8。
编辑
在使用 Scala 的 Databricks Notebook 上进行了测试,然后运行良好。
【问题讨论】:
-
您是否发现任何错误? Apache Spark 不支持从表中读取作为流式查询。如果保存
start()返回的查询并调用query.processAllAvailable(),会抛出后台发生的异常。 -
不,我没有收到任何错误。而且我不想从这里的增量表中读取。 “exampleDF”的选择查询位于我从流中创建的临时表上。根据我的经验,这在 Apache Spark 中是可能的。
标签: java apache-spark spark-structured-streaming delta-lake