回答您的问题 - 您需要有一个主键,即使它仅包含分区键 :-)
更详细的答案实际上取决于这些作业的运行频率、总体数据量、集群中有多少节点、使用什么硬件等。通常,我们会尝试将尽可能多的过滤推送到 Cassandra ,所以它只会返回相关数据,而不是所有数据。最有效的过滤发生在第一个聚类列上,例如,如果我只想处理新创建的条目,那么我可以使用具有以下结构的表:
create table test.test (
pk int,
tm timestamp,
c2 int,
v1 int,
v2 int,
primary key(pk, tm, c2));
然后我可以使用以下方法仅获取新创建的条目:
import org.apache.spark.sql.cassandra._
val data = spark.read.cassandraFormat("test", "test").load()
val filtered = data.filter("tm >= cast('2019-03-10T14:41:34.373+0000' as timestamp)")
或者我可以在给定的时间段内获取条目:
val filtered = data.filter("""ts >= cast('2019-03-10T14:41:34.373+0000' as timestamp)
AND ts <= cast('2019-03-10T19:01:56.316+0000' as timestamp)""")
可以通过在数据帧上执行explain 来检查过滤器下推的效果,并检查PushedFilters 部分 - 标有* 的条件将在Cassandra 端执行...
但并非总是可以设计表来匹配所有查询,因此您需要为最常执行的作业设计主键。在您的情况下,update_date_time 可能是一个很好的候选者,但是如果您将它作为聚类列,那么在更新它时需要小心 - 您需要批量执行更改,如下所示:
begin batch
delete from table where pk = ... and update_date_time = old_timestamp;
insert into table (pk, update_date_time, ...) values (..., new_timestamp, ...);
apply batch;
或类似的东西。