【问题标题】:spark sql to transfer data between Cassandra tablesspark sql 在 Cassandra 表之间传输数据
【发布时间】:2019-01-26 11:52:42
【问题描述】:

请在下面找到 Cassandra 表。

我正在尝试将数据从 1 个 Cassandra 表复制到另一个具有相同结构的 Cassandra 表。

请帮助我。

CREATE TABLE data2 (
        d_no text,
        d_type text,
        sn_perc int,
        tse_dt timestamp,
        f_lvl text,
        ign_f boolean,
        lk_loc text,
        lk_ts timestamp,
        mi_rem text,
        nr_fst text,
        perm_stat text,
        rec_crt_dt timestamp,
        sr_stat text,
        sor_query text,
        tp_dat text,
        tp_ts timestamp,
        tr_rem text,
        tr_type text,
        PRIMARY KEY (device_serial_no, device_type)
    ) WITH CLUSTERING ORDER BY (device_type ASC)

数据插入使用:

Insert into data2(all column names) values('64FCFCFC','HUM',4,'1970-01-02 05:30:00’ ,’NA’,true,'NA','1970-01-02 05:40:00',’NA’,'NA','NA','1970-02-01 05:30:00','NA','NA','NA','1970-02-03 05:30:00','NA','NA');

注意: 当我尝试像这样插入“1970-01-02 05:30:00”时的第 4 列时间戳,并且在 dtaframe 中也正确插入了时间戳,但是当从数据帧插入到 cassandra 并使用 select * from table 时,我看到了它的存在插入像 1970-01-02 00:00:00.000000+0000

类似地,所有时间戳列都会发生这种情况。

pom.xml

<dependencies>
       <dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-core_2.11</artifactId>
    <version>2.3.0</version>
</dependency>
<!-- https://mvnrepository.com/artifact/org.apache.spark/spark-sql -->
<dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-sql_2.11</artifactId>
    <version>2.3.1</version>
</dependency>
<!-- https://mvnrepository.com/artifact/com.datastax.spark/spark-cassandra-connector -->
<dependency>
    <groupId>com.datastax.spark</groupId>
    <artifactId>spark-cassandra-connector_2.11</artifactId>
    <version>2.3.1</version>
</dependency>

我想读取这些值并使用 spark Scala 将其写入另一个 Cassandra 表。见以下代码:

val df2 = spark.read
                       .format("org.apache.spark.sql.cassandra")
                       .option("spark.cassandra.connection.host","hostname")
                       .option("spark.cassandra.connection.port","9042")
                       .option( "spark.cassandra.auth.username","usr")
                       .option("spark.cassandra.auth.password","pas")
                       .option("keyspace","hr")
                       .option("table","data2")
                       .load()
Val df3 =doing some processing on df2.
df3.write
         .format("org.apache.spark.sql.cassandra")
         .mode("append")
         .option("spark.cassandra.connection.host","hostname")
         .option("spark.cassandra.connection.port","9042")
         .option( "spark.cassandra.auth.username","usr")
         .option("spark.cassandra.auth.password","pas")
         .option("spark.cassandra.output.ignoreNulls","true")
         .option("confirm.truncate","true")
         .option("keyspace","hr")
         .option("table","data3")
         .save()

但是当我尝试使用上面的代码插入数据时,我遇到了错误,

java.lang.IllegalArgumentException: requirement failed: Invalid row size: 18 instead of 17.
    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:99)
    at com.datastax.spark.connector.writer.GroupingBatchBuilder.next(GroupingBatchBuilder.scala:106)
    at com.datastax.spark.connector.writer.GroupingBatchBuilder.next(GroupingBatchBuilder.scala:31)
    at scala.collection.Iterator$class.foreach(Iterator.scala:891)
    at com.datastax.spark.connector.writer.GroupingBatchBuilder.foreach(GroupingBatchBuilder.scala:31)
    at com.datastax.spark.connector.writer.TableWriter$$anonfun$writeInternal$1.apply(TableWriter.scala:233)
    at com.datastax.spark.connector.writer.TableWriter$$anonfun$writeInternal$1.apply(TableWriter.scala:210)
    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:145)
    at com.datastax.spark.connector.cql.CassandraConnector.withSessionDo(CassandraConnector.scala:111)
    at com.datastax.spark.connector.writer.TableWriter.writeInternal(TableWriter.scala:210)
    at com.datastax.spark.connector.writer.TableWriter.insert(TableWriter.scala:197)
    at com.datastax.spark.connector.writer.TableWriter.write(TableWriter.scala:183)
    at com.datastax.spark.connector.RDDFunctions$$anonfun$saveToCassandra$1.apply(RDDFunctions.scala:36)
    at com.datastax.spark.connector.RDDFunctions$$anonfun$saveToCassandra$1.apply(RDDFunctions.scala:36)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
    at org.apache.spark.scheduler.Task.run(Task.scala:109)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:345)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)

【问题讨论】:

  • 检查df2.printSchemadf3.printSchema的输出

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


【解决方案1】:

这是一个已知问题 (SPARKC-541) - 您正在将启用 DSE 搜索的表中的数据复制到没有启用 DSE 搜索的表中。您只需将此列作为转换的一部分删除:

val df3 = df2.drop("solr_query").... // your transformations

或者您可以简单地使用较新的驱动程序(如果您使用的是 OSS 驱动程序,则为 2.3.1),或包含此修复程序的相应 DSE 版本。

【讨论】:

  • 嘿,谢谢 Alex,当我删除 solr_query 时,它工作正常。我是 cassandra 的新手,我不明白您建议的第二个选项。请您详细说明一下,添加上面的 pom.xml。
  • 有不同的实现,具体取决于您使用的内容 - 您可以使用 DSE Analytics,然后 Cassandra 连接器就在那里,您可以将其声明为 provided 依赖项(但您的 DSE 版本应该有正确的版本 - 现在不能说它是在哪个修复的,请检查支持),或者您可以使用外部 Spark 集群,然后您可以使用 DSE BYOS jar,或者像上面那样使用连接器。您使用的是哪个 DSE 版本?
  • 你是说下面的版本? com.datastax.sparkspark-cassandra-connector_2.112.3.1
  • 不,我是指 DSE 本身的版本 - 你是在 DSE 上执行 spark 作业,还是在独立的 Spark 上执行?
  • 独立,我在本地设置了 1 个 cassandra 实例,远程另一个集群在那里,我在本地执行..
猜你喜欢
  • 2018-11-11
  • 1970-01-01
  • 1970-01-01
  • 2013-05-14
  • 2017-01-07
  • 1970-01-01
  • 2018-01-24
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多