· 8 years ago · May 07, 2018, 08:36 PM
1private def getCreateTableStatement(keyspace: String, tableName: String, df: DataFrame, keys: Array[String]) = {
2 val primaryKey = keys.mkString(",")
3 val fields = df.schema.map(field => field.name + " " + field.dataType.simpleString.trim())
4 val fieldsDesc = fields.mkString(", ")
5 "CREATE TABLE if not exists " + keyspace + "." + tableName + " (" + fieldsDesc + ", PRIMARY KEY (" + primaryKey + "))"
6}
7
8private def getInsertStatement(keyspace: String, tableName: String, df: DataFrame) = {
9 val fieldsDesc = df.schema.fields.map(f => f.name).mkString(",")
10 val valuesHolder = df.schema.fields.map(_ => "?").mkString(",")
11 "insert into " + keyspace + "." + tableName + " (" + fieldsDesc + ") values (" + valuesHolder + ")"
12}