【问题标题】:Merge into delta table not working with java foreachbatch合并到增量表中不适用于 java foreachbatch
【发布时间】: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


【解决方案1】:

如果您想用新数据更新数据,请尝试遵循如下语法

WHEN NOT MATCHED THEN 
    UPDATE SET e.load_date = u.load_date AND  e.id = u.id
    

如果你只想添加它占据的数据是这样的

WHEN NOT MATCHED THEN INSERT *

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-11-28
    • 2018-06-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多