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