Skip to main content
Question

Unable to write spark DF in Vertica using API

  • December 4, 2019
  • 38 replies
  • 61 views

Show first post

38 replies

Bryan_H
Forum|alt.badge.img+2
  • Participating Frequently
  • December 5, 2019

I am confused now: ACCT_final schema says col2: string (nullable = false) so I would not expect NULL to work. However, the exception should be "[Vertica]JDBC Parameter 2 is not nullable." Is the table ACCT defined with NOT NULL on col2?
Also, what version of Vertica JDBC driver do you have? I am running my test with vertica-jdbc-9.3.0-1.jar


Prakhar84
Forum|alt.badge.img
  • Author
  • Participating Frequently
  • December 5, 2019

i guess col2 value is " ,if you see below
[Row(col1=u'ABC', col2=u'', col3=u'123456789'), Row(col1=u'ABC', col2=u'', col3=u'12345678')
target vertica table DDL.all columns are nullable =true in vertica DB
CREATE TABLE test_vertica
(
col1 varchar(1024),
col2 varchar(1024),
col3 varchar(1024)
)
driver is :vertica-jdbc-9.2.0-0.jar


Bryan_H
Forum|alt.badge.img+2
  • Participating Frequently
  • December 5, 2019

As LenoyJ suggested, we'll try a custom dialect. Vertica is almost verbatim Postgres, but Spark doesn't know that. So try adding the attached VerticaDialect.scala class to your project (remove .txt at end!), and register the dialect as follows:
JdbcDialects.registerDialect(VerticaDialect)
val df = spark.read.jdbc(urlWithUserAndPass, "TESTDB", new
Properties()
// your code here
......
JdbcDialects.unregisterDialect(VerticaDialect)


Prakhar84
Forum|alt.badge.img
  • Author
  • Participating Frequently
  • December 6, 2019

I am just starting with pyspark so couple of questions
a) Where should I save this file
b) Will this work in pyspark as well ,I have to use pyspark for coding this
c) How will below code change if pyspark (also i see after new below there is no closing brackets etc)
val df = spark.read.jdbc(urlWithUserAndPass, "TESTDB", new
Properties()
Please suggest my complete code is below :
pyspark2 --jars /home/x/vertica-9.0.1_spark2.1_scala2.11.jar,/home/x/vertica-jdbc-9.2.0-0.jar
from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField
from pyspark.sql.types import DoubleType, IntegerType, StringType,DateType
from pyspark.sql import functions
from pyspark.sql import HiveContext
scSpark = SparkSession.builder.appName("A").config("spark.sql.warehouse.dir","hdfs:///user/hive/warehouse").enableHiveSupport().getOrCreate()
out_file="/x/y/z/bus_date=20190430/source=abc/"
ACCT_df = spark.read.format("csv").option("header", "true").load(out_file)
ACCT_df.createOrReplaceTempView("ACCT")
query="select col1,col2,col3 from ACCT"
ACCT_final=scSpark.sql(query)
ACCT_final.printSchema()
url1 = "jdbc:vertica://vertica-dw-dev.net/db"
ACCT_final.write.format("jdbc").option("driver", "com.vertica.jdbc.Driver").option("url", url1).option("dbtable", "test_vertica").option("savemode","append").option("user", "xxxxx").option("password", "xxxxxxx").mode("append").save()


LenoyJ
Forum|alt.badge.img+1
  • Participating Frequently
  • December 6, 2019

@Bryan_H's comment started me on a path to modify Postgres's Spark Dialect for Vertica. A few customers requested it before so now's a good time as ever. For one, Vertica does not have a TEXT data type - it's all VARCHAR. And there are no ARRAY types in Vertica (yet). I haven't added Vertica specific data types (like UUID) - probably scope for the future.
@Prakhar84 , I did the below with Scala. You should be able to extrapolate to pyspark (google how to import Spark JDBC dialects)

  1. I have a CSV file as follows with nulls:

    $ cat out.csv
    col1,col2,col3
    ABC,,123
    DEF,,456
    
  2. Start Spark:

    ./bin/spark-shell --jars /opt/vertica/packages/SparkConnector/lib/vertica-spark2.1_scala2.11.jar,/opt/vertica/java/vertica-jdbc.jar,/home/dbadmin/VerticaDialect.scala --master spark://172.31.4.157:7077
    
  3. Import the following Vertica Dialect and Register it. Modified Postgres's dialect from here.

    //#Begin Scala code
    import org.apache.spark.sql.jdbc.{JdbcDialects, JdbcDialect, JdbcType}
    import java.sql.{Connection, Types}
    import java.util.Locale
    
    import org.apache.spark.sql.execution.datasources.jdbc.{JDBCOptions, JdbcUtils}
    import org.apache.spark.sql.types._
    
    
    val VerticaDialect = new JdbcDialect { 
    
      val MAX_PRECISION = 38
      val MAX_SCALE = 38
    
      import scala.math.min
    
      def bounded(precision: Int, scale: Int): DecimalType = {
          DecimalType(min(precision, MAX_PRECISION), min(scale, MAX_SCALE))
      }
    
      override def canHandle(url: String): Boolean =
        url.toLowerCase(Locale.ROOT).startsWith("jdbc:vertica")
    
      override def getCatalystType(
          sqlType: Int, typeName: String, size: Int, md: MetadataBuilder): Option[DataType] = {
        if (sqlType == Types.REAL) {
          Some(FloatType)
        } else if (sqlType == Types.SMALLINT) {
          Some(ShortType)
        } else if (sqlType == Types.BIT && typeName.equals("bit") && size != 1) {
          Some(BinaryType)
        } else if (sqlType == Types.OTHER) {
          Some(StringType)
        } 
        else None
      }
    
      private def toCatalystType(
          typeName: String,
          precision: Int,
          scale: Int): Option[DataType] = typeName match {
        case "bool" => Some(BooleanType)
        case "bit" => Some(BinaryType)
        case "int2" => Some(ShortType)
        case "int4" => Some(IntegerType)
        case "int8" | "oid" => Some(LongType)
        case "float4" => Some(FloatType)
        case "money" | "float8" => Some(DoubleType)
        case "text" | "varchar" | "char" | "cidr" | "inet" | "json" | "jsonb" | "uuid" =>
          Some(StringType)
        case "bytea" => Some(BinaryType)
        case "timestamp" | "timestamptz" | "time" | "timetz" | "datetime" | "smalldatetime" => Some(TimestampType)
        case "date" => Some(DateType)
        case "numeric" | "decimal" if precision > 0 => Some(bounded(precision, scale))
        case "numeric" | "decimal" => Some(DecimalType. SYSTEM_DEFAULT)
        case _ => None
      }
    
      override def getJDBCType(dt: DataType): Option[JdbcType] = dt match {
        case StringType => Some(JdbcType("VARCHAR", Types.CHAR))
        case BinaryType => Some(JdbcType("BYTEA", Types.BINARY))
        case BooleanType => Some(JdbcType("BOOLEAN", Types.BOOLEAN))
        case FloatType => Some(JdbcType("FLOAT", Types.FLOAT))
        case DoubleType => Some(JdbcType("FLOAT8", Types.DOUBLE))
        case ShortType | ByteType => Some(JdbcType("SMALLINT", Types.SMALLINT))
        case t: DecimalType => Some(
          JdbcType(s"NUMERIC(${t.precision},${t.scale})", java.sql.Types.NUMERIC))
        case _ => None
      }
    
      override def getTableExistsQuery(table: String): String = {
        s"SELECT 1 FROM $table LIMIT 1"
      }
    
      override def isCascadingTruncateTable(): Option[Boolean] = Some(false)
    
      override def getTruncateQuery(
          table: String,
          cascade: Option[Boolean] = isCascadingTruncateTable): String = {
        cascade match {
          case Some(true) => s"TRUNCATE TABLE $table"
          case _ => s"TRUNCATE TABLE $table"
        }
      }
    
      override def beforeFetch(connection: Connection, properties: Map[String, String]): Unit = {
        super.beforeFetch(connection, properties)
    
        if (properties.getOrElse(JDBCOptions.JDBC_BATCH_FETCH_SIZE, "0").toInt > 0) {
          connection.setAutoCommit(false)
        }
      }
    
    }
    
    JdbcDialects.registerDialect(VerticaDialect)
    
  4. Read the CSV File & Show:

    import org.apache.spark.sql.types._
    import org.apache.spark.sql.{DataFrame, Row, SQLContext, SaveMode}
    import com.vertica.spark._
    import org.apache.spark.sql.SQLContext
    
    val sqlContext = new SQLContext(sc)
    
    val out_file="/home/dbadmin/out.csv"
    val df = sqlContext.read.format("csv").option("header", "true").load(out_file)
    
    df.show
    //#+----+----+----+
    //#|col1|col2|col3|
    //#+----+----+----+
    //#| ABC|null| 123|
    //#| DEF|null| 456|
    //#+----+----+----+
    
  5. Create a JDBC connection string and write to Vertica

    val url = "jdbc:vertica://172.31.4.157/dbname?username=dbadmin&password=&ConnectionLoadBalance=1&DirectBatchInsert=1"
    
    //#Using Append Mode. Make sure table exists with the right schema in Vertica...
    val mode = SaveMode.Append
    
    //#Write using JDBC.
    df.write.format("jdbc").option("driver", "com.vertica.jdbc.Driver").option("url", url).option("dbtable", "public.test_vertica").mode(mode).save()
    
  6. Notice that you don't get the "Driver not capable." error anymore as we've used the above Vertica Dialect and Nulls are handled properly.

  7. Switch over to vsql and check the table if Nulls have transferred:

    dbadmin=> select * from test_vertica;
     col1 | col2 | col3
    ------+------+------
     ABC  |      | 123
     DEF  |      | 456
    (2 rows)
    
    dbadmin=> select col1, col2, coalesce(col2,'yes it is null') as Is_Col2_Null, col3 from test_vertica;
     col1 | col2 |  Is_Col2_Null  | col3
    ------+------+----------------+------
     ABC  |      | yes it is null | 123
     DEF  |      | yes it is null | 456
    (2 rows)
    

Prakhar84
Forum|alt.badge.img
  • Author
  • Participating Frequently
  • December 6, 2019

Hi Lenoy,
Thanks for the detailed exaplanation but will need your help to implement this in pyspark code ...searched on google but could not find more
based on below link I tried
https://stackoverflow.com/questions/51731998/how-to-add-custom-jdbc-dialects-in-pyspark
Are you saying
a) save the verticadialect.scala shared earlier in a location and the call pyspark like below?
pyspark2 --jars /home/x/vertica-9.0.1_spark2.1_scala2.11.jar,/home/x/vertica-jdbc-9.2.0-0.jar,/home/x/VerticaDialect.scala
Tried below but getting error
from py4j.java_gateway import java_import
gw = spark.sparkContext._gateway
java_import(gw.jvm, "com.me.VerticaDialect")
gw.jvm.org.apache.spark.sql.jdbc.JdbcDialects.registerDialect(gw.jvm.com.me.VerticaDialect())
but i get an error
Traceback (most recent call last):
File "", line 1, in
TypeError: 'JavaPackage' object is not callable

Please help us to implement in pyspark as well so that it will be helful for clients using pyspark not scala.


Bryan_H
Forum|alt.badge.img+2
  • Participating Frequently
  • December 6, 2019

Hi, I believe you would need to compile the scala file into a JAR and add to the classpath. We will run some tests to determine the best approach but it may be fairly complex to set up.


Prakhar84
Forum|alt.badge.img
  • Author
  • Participating Frequently
  • December 6, 2019

Hi Bryan
It will be very beneficial for clients who use pyspark for pushing data in vertica to have this setup properly,if you can guide us step by step to work on my code then others can also leverage the same.
Is there a way that dialect can be directly written for pyspark? (just asking)
Also what will be the best way to integrate this with spark rather then registering dialect everytime
Prakhar


LenoyJ
Forum|alt.badge.img+1
  • Participating Frequently
  • December 6, 2019

@Prakhar84, quick tangent - I'm curious, what became of the checking of the connectivity between Vertica & HDFS? Most of my customers just get the Spark Connector working which usually solves most issues. I believe you raised another discussion here. Let's continue that discussion there.


Prakhar84
Forum|alt.badge.img
  • Author
  • Participating Frequently
  • December 6, 2019

@Bryan_H

Please let us know of any solution for this dialect in pyspark .
Thanks for your guidance so far.

Prakhar


Bryan_H
Forum|alt.badge.img+2
  • Participating Frequently
  • December 8, 2019

Hi, a workaround we have suggested for other customers is to write Parquet files to a temporary folder that Vertica can read, then use vertica-python driver to issue a COPY command to import the Parquet file (see http://github.com/vertica/vertica-python for details)
We are investigating better solutions; however, and in my own opinion, we should implement a complete VerticaDialect and commit to Apache Spark to fix this for all Sprk programs whether Java, Scala, and PySpark. It will take a long time to develop this and then a long time for Apache to accept and publish as part of next Spark though.


LenoyJ
Forum|alt.badge.img+1
  • Participating Frequently
  • December 8, 2019

I'd also recommend checking with the Spark community on how to use JDBC dialects in PySpark.


Bryan_H
Forum|alt.badge.img+2
  • Participating Frequently
  • December 11, 2019

Hi, you can add compiled classes to the classpath. However, you would need to rebuild Spark from a GitHub checkout after applying the following patch:
https://github.com/bryanherger/spark/commit/84d3014e4ead18146147cf299e8996c5c56b377d
This would embed the VerticaDialect into the Spark build. To avoid having to reinstall the custom Spark build, you can extract the compiled VerticaDialect class and apply that at runtime as part of the classpath.