From 9964c996aa4f65dd80f03d20cde24847a5104884 Mon Sep 17 00:00:00 2001 From: Cerdore Date: Wed, 29 Sep 2021 15:32:02 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E6=94=B9=E9=83=A8=E5=88=86=E4=BB=A3?= =?UTF-8?q?=E7=A0=81=EF=BC=8C=E6=B7=BB=E5=8A=A0=E9=83=A8=E5=88=86=E6=B3=A8?= =?UTF-8?q?=E9=87=8A=E4=BB=A5=E5=8F=8A=E5=88=A0=E9=99=A4=E6=97=A0=E7=94=A8?= =?UTF-8?q?=E6=B3=A8=E9=87=8A=E4=B8=8E=E4=BB=A3=E7=A0=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Cerdore --- .../src/main/scala/SQLDataSourceExample.scala | 3 +-- .../sources/opengauss/OpenGaussDataSource.scala | 15 +++++++-------- .../org/opengauss/spark/OpenGaussExample.scala | 17 +++-------------- 3 files changed, 11 insertions(+), 24 deletions(-) diff --git a/SparkOpOpenGauss/src/main/scala/SQLDataSourceExample.scala b/SparkOpOpenGauss/src/main/scala/SQLDataSourceExample.scala index f2126d7f..e7766e01 100644 --- a/SparkOpOpenGauss/src/main/scala/SQLDataSourceExample.scala +++ b/SparkOpOpenGauss/src/main/scala/SQLDataSourceExample.scala @@ -26,7 +26,7 @@ object SQLDataSourceExample { // $example on:jdbc_dataset$ // Note: JDBC loading and saving can be achieved via either the load/save or jdbc methods // Loading data from a JDBC source - val dburl = "jdbc:postgresql://x.x.x.x:port/school" + val dburl = "jdbc:postgresql://x.x.x.x:port/school" //注意修改此处的ip与端口地址 val jdbcDF = spark.read .format("jdbc") .option("url", dburl) @@ -65,6 +65,5 @@ object SQLDataSourceExample { jdbcDF3.write .option("createTableColumnTypes", "cla_id INT, cla_name VARCHAR(20)") .jdbc(dburl, "customtable3", connectionProperties) - // $example off:jdbc_dataset$ } } \ No newline at end of file diff --git a/SparkOpOpenGauss/src/main/scala/org/opengauss/spark/sources/opengauss/OpenGaussDataSource.scala b/SparkOpOpenGauss/src/main/scala/org/opengauss/spark/sources/opengauss/OpenGaussDataSource.scala index 89625577..25dc0246 100644 --- a/SparkOpOpenGauss/src/main/scala/org/opengauss/spark/sources/opengauss/OpenGaussDataSource.scala +++ b/SparkOpOpenGauss/src/main/scala/org/opengauss/spark/sources/opengauss/OpenGaussDataSource.scala @@ -127,18 +127,17 @@ class OpenGaussWriter(connectionProperties: ConnectionProperties) extends DataWr connectionProperties.password ) - // TODO:待修改 - val statement = "insert into ${connectionProperties.tableName} (cor_name, cor_type, credit) values (?,?,?)" + val statement = "insert into ${connectionProperties.tableName} (cla_id, cla_name, cla_teacher) values (?,?,?)" val preparedStatement = connection.prepareStatement(statement) override def write(record: InternalRow): Unit = { - val cor_name = record.getString(0) - val cor_type = record.getString(1) - val credit = record.getDouble(2) + val cla_id = record.getInt(0) + val cla_name = record.getString(1) + val cla_teacher = record.getInt(2) - preparedStatement.setString(0, cor_name) - preparedStatement.setString(1, cor_type) - preparedStatement.setDouble(2, credit) + preparedStatement.setInt(0, cla_id) + preparedStatement.setString(1, cla_name) + preparedStatement.setInt(2, cla_teacher) preparedStatement.executeUpdate() } diff --git a/SparkOpOpenGauss/src/test/scala/org/opengauss/spark/OpenGaussExample.scala b/SparkOpOpenGauss/src/test/scala/org/opengauss/spark/OpenGaussExample.scala index acc414c1..0f7270db 100644 --- a/SparkOpOpenGauss/src/test/scala/org/opengauss/spark/OpenGaussExample.scala +++ b/SparkOpOpenGauss/src/test/scala/org/opengauss/spark/OpenGaussExample.scala @@ -10,7 +10,7 @@ import org.scalatest.Matchers.convertToAnyShouldWrapper class OpenGaussExample extends FlatSpec { val testTableName = "course" - val dburl = "jdbc:postgresql://x.x.x.x:port/school" + val dburl = "jdbc:postgresql://x.x.x.x:port/school" //注意将此处修改成你的机器对应的的ip与端口 "Simple data source" should "read" in{ val sparkSession = SparkSession.builder @@ -43,7 +43,6 @@ class OpenGaussExample extends FlatSpec { .option("password", "pwdofsparkuser") .option("tableName", testTableName) .option("partitionSize", 10) -// .option("partitionColumn", "name") .load() .show() @@ -59,7 +58,7 @@ class OpenGaussExample extends FlatSpec { import spark.implicits._ - val df = (60 to 70).map(_.toLong).toDF("product_no") + val df = (60 to 70).map(_.toLong).toDF("cla_id") df .write @@ -69,22 +68,12 @@ class OpenGaussExample extends FlatSpec { .option("password", "pwdofsparkuser") .option("tableName", testTableName) .option("partitionSize", 10) - .option("partitionColumn", "product_no") + .option("partitionColumn", "cla_id") .mode(SaveMode.Append) .save() spark.stop() } - - - object Queries { - lazy val createTableQuery = s"CREATE TABLE $testTableName (user_id BIGINT PRIMARY KEY);" - - lazy val testValues: String = (1 to 50).map(i => s"($i)").mkString(", ") - - lazy val insertDataQuery = s"INSERT INTO $testTableName VALUES $testValues;" - } - }