How to connect to Sqlserver using Spark from a GCP Dataproc cluster in the right manner?

Viewed 304

I am trying to move data from Sqlserver database to Bigquery on GCP. To do that, we created a Dataproc cluster where I can run my spark job which connects to source database on Sqlserver, reads certain tables and ingest them on to Bigquery.

Versions on GCP Dataproc:

Spark: 2.4.7
Scala: 2.12.12

My Spark code:

val dataframe = spark.read.format("jdbc").option("url", s"jdbc:sqlserver://servername:port;DatabaseName=dbname").
option("driver", "com.microsoft.sqlserver.jdbc.SQLServerDriver").
option("user", "username").
option("password", "password").
option("dbtable", s"(select * from tablename where flagCol >= '${YearMonth.of(2019, 1).atDay(1).toString}' and flagCol <= '${YearMonth.of(2019, 12).atEndOfMonth().toString}') as dateDF").
option("partitionColumn", flagCol).
option("lowerBound", s"${YearMonth.of(2019, 1).atDay(1).toString}").
option("numPartitions", s"${YearMonth.of(2019, 12).atEndOfMonth().toString}").
option("numPartitions", 3).
option("fetchsize", 1000).
load()

dataframe.write.format("bigquery").option("table", s"$tablename").mode("append").save()

I create a jar out of it and submit it using the below spark-submit command:

gcloud dataproc jobs submit spark --cluster <CLUSTERNAME> --master yarn --deploy-mode cluster --num-executors 3 --executor-memory 4G --executor-cores 3 --driver-class-path <path_of_driver_jar> --driver-memory 1G --jars <path_of_required_jars> --class com.packageName.className mysparkcode.jar

The problem I am facing is that the connection to the source database from Spark fails randomly with below exception:

Caused by: java.net.SocketException: Connection reset
    at java.net.SocketInputStream.read(SocketInputStream.java:210)
    at java.net.SocketInputStream.read(SocketInputStream.java:141)
    at com.microsoft.sqlserver.jdbc.TDSChannel.read(IOBuffer.java:1981)

At first I thought it is a network/executor timeout and added this configuration:

spark.network.timeout 10000000, spark.executor.heartbeatInterval 10000000

But the issue persists.

So I tried the same code on my local and it worked without any issue on bare minimum resources. I also tried the same code from on of our On-Prem hadoop cluster where spark is available and save the dataframe as a dummy parquet file. To my amusement, without even giving considerable resources, the job ran fine. The job loaded 2.5GB of data into a parquet file in 90seconds.

On-Prem versions:

Spark: 2.3.3
Scala: 2.11.12

But the same code fails on my dataproc cluster.

This is the way I created my cluster.

gcloud dataproc clusters create <CLUSTERNAME> --enable-component-gateway --bucket <GCP_BUCKET_NAME> --region <REGION_NAME> --subnet <SUBNET_ADDRESS> --no-address --zone <ZONE> --master-machine-type n1-standard-8 --master-boot-disk-size 500 --num-workers 2 --worker-machine-type n1-standard-8 --worker-boot-disk-size 500 --metadata 'PIP_PACKAGES=pyspark==2.4.0' --initialization-actions SOME_STARTUP_SCRIPT.sh,//SOME_PATH/pip-install.sh --image-version 1.5-debian10 --project <PROJECT_NAME> --service-account=<SERVICE_ACCOUNT_NAME> --properties <SOME_JARS>,dataproc:dataproc.conscrypt.provider.enable=false --optional-components ANACONDA,JUPYTER

The successful execution on our On-Prem hadoop cluster takes 90seconds to save the dataframe as a parquet file. But on dataproc, even the bare minimum queries like count(*), min/max of columns are resulting in Caused by: java.net.SocketException: Connection reset

I even tried doubling the memory parameters only to face the same exception.

Is this due to the way I wrote my code and executing it the way I am doing it or is it due to a version missmatch issue ? If it is not, is there an issue with the way I created my cluster ? Is there any network setting that I should be using in my cluster creating command or anywhere in my code ? Can anyone let me know what should I do to fix the issue ? Any help is greatly appreciated.

0 Answers
Related