【问题标题】:Spark Cassandra write Dataframe, how to find which keys already exist in database during insertionSpark Cassandra编写Dataframe,插入过程中如何查找数据库中已存在哪些键
【发布时间】:2023-04-06 13:42:01
【问题描述】:

我编写了以下 JAVA 方法,通过 Apache Spark 将多个 POJO 的数据持久保存到 Apache Cassandra 数据库。

这似乎工作正常,但是 Spark 没有提供有关记录是已插入(cassandra 中不存在密钥)还是已更新(数据库中已存在密钥)的任何信息。

有没有一种成本最低的方法(我想避免在数据框中加载表的内容并检查重复键),以便在插入时找出哪些记录已经存在(有重复键)在数据库中?

具体代码如下:

@Service
public class WriteDB {

    @Autowired
    private SparkSession sparkSession;

    Logger LOG = LoggerFactory.getLogger(WriteDB.class);

    public <T> void uploadData(List<T> objects, Class<T> clazz, String keyspaceName, String tableName) {

        LOG.info("Number of records to be committed to database: " + objects.size());

        //Create dataset from entity object
        Dataset<Row> df = sparkSession.createDataFrame(objects, clazz);


        //Write data from spark dataframe to cassandra schema
        df.write().mode(SaveMode.Append).format("org.apache.spark.sql.cassandra").options(new HashMap<String, String>() {{
            put("keyspace", keyspaceName);
            put("table", tableName);
        }}).save();

        LOG.info("Records Commited");
    }
}

【问题讨论】:

  • 顺便说一句,你用的是 Cassandra 还是 DSE
  • 我们目前使用的是 Cassandra 的 Apache 发行版。 Spark 和 Cassandra 之间的通信是使用 Dastax Cassandra 连接器实现的。

标签: apache-spark cassandra insert duplicates spark-cassandra-connector


【解决方案1】:

在 Cassandra 中,一切都是 upsert - 插入和更新之间没有区别。 Cassandra 在插入或更新时不会检查数据是否存在(LWT 除外)——它只是添加数据,而之前的副本在压缩期间会被删除。

完成任务的唯一方法是从表中加载数据 - 使用 Dataframe API,它将通过将整个表读入 Dataframe 然后加入,或在 RDD API 中使用 joinWithCassandra 或 @ 在 Spark 级别上完成987654323@(见doc)。

【讨论】:

  • 谢谢,我希望避免这种情况,因为桌子很大,但似乎这是唯一的方法。
  • 但是你真的需要知道什么更新什么不更新吗?
  • 主要原因是我们收到了大量的测量数据,我们希望在数据库中插入重复数据时能够记录(或者可能返回警告)。
  • 你可以做一些优化,比如只加入主键列,而不是整个数据集等。但是 joinWithCassandra 可能仍然比 Spark 级别的加入更快
猜你喜欢
  • 1970-01-01
  • 2017-05-06
  • 2021-12-28
  • 2021-03-16
  • 1970-01-01
  • 2021-09-14
  • 2018-10-05
  • 1970-01-01
  • 2017-05-01
相关资源
最近更新 更多