【问题标题】:Delete rows from cassandra table using pyspark or cql query使用 pyspark 或 cql 查询从 cassandra 表中删除行
【发布时间】:2020-10-05 02:48:57
【问题描述】:

例如,我有一张包含很多列的表格。 test_event 并且我在同一个键空间中有另一个表 test,其中包含我必须从 test_event 中删除的行的 id。

我试过 deleteFromCassandra,但它不起作用,因为 spark 看不到 SparkContext。 我发现一些解决方案使用了 DELETE FROM,但它是用 scala 编写的。

经过大约一百次尝试,我终于感到困惑并寻求您的帮助。有人可以一步一步来吗?

【问题讨论】:

  • 你能显示代码吗 - deleteFromCassandra 可以正常工作,也许你错过了一些导入
  • @AlexOtt test.deleteFromCassandra(keyspace, test_event)
  • 需要指定keyColumns参数
  • 用ID和主表的主键显示表结构
  • 我昨天错过了它是用于 pyspark - 查看答案...

标签: apache-spark pyspark cassandra spark-cassandra-connector


【解决方案1】:

看看这段代码:

from pyspark.sql import SQLContext

def main_function():

  sql = SQLContext(sc)
  tests = sql.read.format("org.apache.spark.sql.cassandra").\
               load(keyspace="your keyspace", table="test").where(...)
  for test in tests:
    delete_sql = "delete from test_event where id = " + test.select('id')
    sql.execute(delete_sql)

请注意,一次删除一行并不是 spark 的最佳做法,但上面的代码只是一个示例,可帮助您弄清楚您的实现。

【讨论】:

  • 'TypeError: 'Column' 对象不可调用'。我应该将其重写为 UDF 吗?
  • 我制作了带有参数“id”的函数,现在我不明白如何应用它。我的意思是用 id 遍历测试表。
  • 检查数据框的名称,列 'id' 可能不存在。如果您想通过 id 删除 id 那么是的,您将需要进行迭代,请记住这不是 spark 的好习惯,建议您过滤然后删除所有匹配的内容。例如,对于 Cassandra,您可以通过分区键删除。
【解决方案2】:

Spark Cassandra 连接器 (SCC) 本身仅提供适用于 Python 的 Dataframe API。但是有一个pyspark-cassandra package在SCC之上提供了RDD API,所以可以如下进行删除。

使用(我已尝试使用 Spark 2.4.3)启动 pyspark shell:

bin/pyspark --conf spark.cassandra.connection.host=IPs\
    --packages anguenot:pyspark-cassandra:2.4.0

并从一个表中读取数据,并进行删除。您需要有源数据才能拥有与主键对应的列。它可以是完整主键、部分主键或仅分区键 - 根据它,Cassandra 将使用相应的墓碑类型(行/范围/分区墓碑)。

在我的示例中,表的主键由一列组成 - 这就是我在数组中仅指定一个元素的原因:

rdd = sc.cassandraTable("test", "m1")
rdd.deleteFromCassandra("test","m1", keyColumns = ["id"])

【讨论】:

    猜你喜欢
    • 2014-10-18
    • 2020-11-12
    • 2014-11-15
    • 2019-05-18
    • 1970-01-01
    • 2015-01-22
    • 2013-06-19
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多