【问题标题】:Spark Scala Cassandra connector delete all all rows is failing with IllegalArgumentException requirement failed ExceptionSpark Scala Cassandra 连接器删除所有行失败并出现 IllegalArgumentException 要求失败异常
【发布时间】:2021-10-14 17:49:06
【问题描述】:

创建表-

CREATE TABLE test.word_groups (group text, word text, count int,PRIMARY KEY (group,word));

插入数据 -

INSERT INTO test.word_groups (group , word , count ) VALUES ( 'A-group', 'raj', 0) ;
INSERT INTO test.word_groups (group , word , count ) VALUES ( 'b-group', 'jaj', 0) ;
INSERT INTO test.word_groups (group , word , count ) VALUES ( 'A-group', 'raff', 3) ;

 SELECT * FROM word_groups ;

 group   | word | count
---------+------+-------
 b-group |  jaj |     0
 A-group | raff |     3
 A-group |  raj |     0

脚本 -

val cassandraUrl = "org.apache.spark.sql.cassandra"
val wordGroup: Map[String, String] = Map("table" ->"word_groups", 
  "keyspace" -> "test", "cluster" -> "test-cluster")
val groupData = {spark.read.format(cassandraUrl).options(wordGroup).load()
  .where(col("group") === "b-group")}
groupData.rdd.deleteFromCassandra("sunbird_courses", "word_groups")

例外-

java.lang.IllegalArgumentException: requirement failed: Invalid row size: 3 instead of 2.
    at scala.Predef$.require(Predef.scala:224)
    at com.datastax.spark.connector.writer.SqlRowWriter.readColumnValues(SqlRowWriter.scala:23)
    at com.datastax.spark.connector.writer.SqlRowWriter.readColumnValues(SqlRowWriter.scala:12)
    at com.datastax.spark.connector.writer.BoundStatementBuilder.bind(BoundStatementBuilder.scala:102)
    at com.datastax.spark.connector.writer.GroupingBatchBuilder.next(GroupingBatchBuilder.scala:105)
    at com.datastax.spark.connector.writer.GroupingBatchBuilder.next(GroupingBatchBuilder.scala:30)
    at scala.collection.Iterator$class.foreach(Iterator.scala:891)
    at com.datastax.spark.connector.writer.GroupingBatchBuilder.foreach(GroupingBatchBuilder.scala:30)
    at com.datastax.spark.connector.writer.TableWriter$$anonfun$writeInternal$1.apply(TableWriter.scala:229)
    at com.datastax.spark.connector.writer.TableWriter$$anonfun$writeInternal$1.apply(TableWriter.scala:198)
    at com.datastax.spark.connector.cql.CassandraConnector$$anonfun$withSessionDo$1.apply(CassandraConnector.scala:112)
    at com.datastax.spark.connector.cql.CassandraConnector$$anonfun$withSessionDo$1.apply(CassandraConnector.scala:111)
    at com.datastax.spark.connector.cql.CassandraConnector.closeResourceAfterUse(CassandraConnector.scala:129)
    at com.datastax.spark.connector.cql.CassandraConnector.withSessionDo(CassandraConnector.scala:111)
    at com.datastax.spark.connector.writer.TableWriter.writeInternal(TableWriter.scala:198)
    at com.datastax.spark.connector.writer.TableWriter.delete(TableWriter.scala:194)
    at com.datastax.spark.connector.RDDFunctions$$anonfun$deleteFromCassandra$1.apply(RDDFunctions.scala:119)
    at com.datastax.spark.connector.RDDFunctions$$anonfun$deleteFromCassandra$1.apply(RDDFunctions.scala:119)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
    at org.apache.spark.scheduler.Task.run(Task.scala:123)
    at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
    at java.lang.Thread.run(Thread.java:745)
21/08/11 09:01:24 WARN TaskSetManager: Lost task 0.0 in stage 11.0 (TID 2953, localhost, executor driver): java.lang.IllegalArgumentException: requirement failed: Invalid row size: 3 instead of 2.

Spark 版本 - 2.4.4 和 Spark Cassandra 连接器版本 - 2.5.0

Spark Cassandra 连接器文档链接 - https://github.com/datastax/spark-cassandra-connector/blob/master/doc/5_saving.md#deleting-rows-and-columns

我正在尝试删除这些列的所有记录,包括主键。

有什么解决办法吗?

仅供参考 - 我需要从 word_groups 表中删除组“A-group”的所有记录,包括主键/分区键

【问题讨论】:

    标签: dataframe apache-spark cassandra rdd spark-cassandra-connector


    【解决方案1】:

    这是 2.5.x 中有趣的变化,我不知道 - 即使指定了 keyColumns,您现在也需要有正确的行大小,以前没有它也可以工作 - 对我来说似乎是一个错误。

    删除整行时只需要保留主键 - 将删除更改为:

    groupData.select("group", "word").rdd.deleteFromCassandra("test", "word_groups")
    

    但在您的情况下,最好根据分区键列删除 - 在这种情况下,您将只有一个墓碑(您仍然需要只选择必要的列):

    import com.datastax.spark.connector._
    {groupData.select("group").rdd
      .deleteFromCassandra("test", "word_groups", keyColumns = SomeColumns("group"))}
    

    您甚至不需要从 Cassandra 读取输入数据 - 如果您知道分区键的值,那么您可以创建 RDD 并删除数据(类似于 doc 中所示):

    case class Key (group:String)
    { sc.parallelize(Seq(Key("b-group")))
       .deleteFromCassandra("test", "word_groups", keyColumns = SomeColumns("group"))}
    

    【讨论】:

    • 不提 SomeColumns 就不可能吗?有没有其他不提主键的删除方式?
    • java.io.IOException: Failed to prepare statement DELETE "group", "word" FROM "test"."word_groups" WHERE "group" = :"group" AND "word" = :"word": Invalid identifier group for deletion (should not be a PRIMARY KEY part) at com.datastax.spark.connector.writer.TableWriter.com$datastax$spark$connector$writer$TableWriter$$prepareStatement(TableWriter.scala:142) 如果执行上述查询(SomeColumns("group", "word"))会得到这个异常
    • groupData.rdd.deleteFromCassandra("sunbird_courses", "word_groups",keyColumns = SomeColumns("word", "group")) 我得到相同的异常错误java.lang.IllegalArgumentException: requirement failed: Invalid row size: 3 instead of 2. @Alex Ott
    • 抱歉,2.5.x 中的行为似乎发生了变化,我已经更新了答案
    • 特别是我们必须选择主键并调用删除任务权限
    猜你喜欢
    • 1970-01-01
    • 2023-03-14
    • 2016-02-04
    • 1970-01-01
    • 2021-12-21
    • 2021-03-04
    • 1970-01-01
    • 2021-11-11
    • 1970-01-01
    相关资源
    最近更新 更多