カウンター列を持つテーブルの UPDATE は、spark-cassandra-connector を介して実行できます。モード「追加」(または必要に応じてSaveMode .Append)で保存するDataFramesおよびDataFrameWriterメソッドを使用する必要があります。コードDataFrameWriter.scalaを確認してください。
たとえば、次の表があるとします。
cqlsh:test> SELECT * FROM name_counter ;
name | surname | count
---------+---------+-------
John | Smith | 100
Zhang | Wei | 1000
Angelos | Papas | 10
コードは次のようになります。
val updateRdd = sc.parallelize(Seq(Row("John", "Smith", 1L),
Row("Zhang", "Wei", 2L),
Row("Angelos", "Papas", 3L)))
val tblStruct = new StructType(
Array(StructField("name", StringType, nullable = false),
StructField("surname", StringType, nullable = false),
StructField("count", LongType, nullable = false)))
val updateDf = sqlContext.createDataFrame(updateRdd, tblStruct)
updateDf.write.format("org.apache.spark.sql.cassandra")
.options(Map("keyspace" -> "test", "table" -> "name_counter"))
.mode("append")
.save()
更新後:
name | surname | count
---------+---------+-------
John | Smith | 101
Zhang | Wei | 1002
Angelos | Papas | 13
DataFrame 変換は、RDD を DataFrame に暗黙的に変換しimport sqlContext.implicits._
、.toDF()
.
このおもちゃのアプリケーションの完全なコードを確認してください:
https://github.com/kyrsideris/SparkUpdateCassandra/tree/master
ここではバージョンが非常に重要であるため、上記は Scala 2.11.7、Spark 1.5.1、spark-cassandra-connector 1.5.0-RC1-s_2.11、Cassandra 3.0.5 に適用されます。DataFrameWriter は since と指定され@Experimental
てい@since 1.4.0
ます。