Spark job getting stuck at the final stage of importing from Oracle DB - Data is not skewed

Viewed 1501

I am trying to pull data from Oracle DB and putting it to AWS S3 using Apache Spark 2.3.1. The job is running fine till the last stage and getting stuck there. I don't think the data is skewed because each stage has an equal number of records. Below is the query I'm using in spark.

url = "jdbc:oracle:thin:@IP:PORT/SID"
user = "user"
password = "password"
driver = "oracle.jdbc.driver.OracleDriver"
table = "table"
fetchSize = 1000
partitionColumn = "num_rows"

date1 = (datetime.today() - td(days=42)).date().strftime('%d-%b-%Y')
date2 = (datetime.today() - td(days=2)).date().strftime('%d-%b-%Y')

query = "(select min(rownum) as min, max(rownum) as max from "+table+" where date>='"+str(date1)+"' and date<='"+str(date2)+"') tmp1"
print(query)

DF = spark.read.format("jdbc").option("url", url) \
                              .option("dbtable", query) \
                              .option("user", user) \
                              .option("password", password) \
                              .option("driver", driver) \
                              .load()

lower_bound, upper_bound = DF.first()
lower_bound = int(lower_bound)
upper_bound = int(upper_bound)
numPartitions = int(upper_bound/fetchSize)+1
print(lower_bound,upper_bound)
print(numPartitions)

query = "(select t1.*, ROWNUM as num_rows from (select * from " + table + " where date>='"+str(date1)+"' and date<='"+str(date2)+"') t1) tmp2"
print(query)

DF = spark.read.format("jdbc").option("url", url) \
                              .option("dbtable", query) \
                              .option("user", user) \
                              .option("password", password) \
                              .option("fetchSize",fetchSize) \
                              .option("numPartitions", numPartitions) \
                              .option("partitionColumn", partitionColumn) \
                              .option("lowerBound", lower_bound) \
                              .option("upperBound", upper_bound) \
                              .option("driver", driver) \
                              .load()

path = "s3://my_path"
DF.write.mode("overwrite").parquet(path)

The code basically pulls last 42 days data and put it into S3 bucket. Below is the output till the write statement. Code was run on '10-Sep-2018'

(select min(rownum) as min, max(rownum) as max from table where date>='30-Jul-2018' and date<='08-Sep-2018') tmp1
(1, 2195427)
2196
(select t1.*, ROWNUM as num_rows from (select * from table where date>='30-Jul-2018' and date<='08-Sep-2018') t1) tmp2

As you can see,

  • the total number of records = 2195427
  • records per partition = 1000
  • number of partitions = 2196

So the job has 2196 stages and each stage pulls 1000 records. The job is getting stuck at 2191/2196 with 5 more stages to go.

Hardware Specs:

I am using r4.xlarge machines. My cluster is 1 Master, 2 Slaves of r4.xlarge. Below are my driver and executor specs.

spark.driver.cores  8
spark.driver.memory 24g
spark.driver.memoryOverhead 3072M
spark.executor.cores    1
spark.executor.memory   3g
spark.executor.memoryOverhead   512M
spark.yarn.am.cores 1
spark.yarn.am.memory    3g
spark.yarn.am.memoryOverhead    512M

Spark Executor UI

The stages 1 to 2191 got completed in 1.3 hrs but remaining 5 stages are stuck for more than three hours.

Please find the log here : https://github.com/rinazbelhaj/stackoverflow/blob/master/Spark_Log_10_Sept_2018

I am not able to figure out the root cause of this issue.

1 Answers

I guess you have two possible scenarios:

1 - Problem related to Oracle DB

I'm dealing with a very similar problem. But instead of stuck the task was being interrupted by SQL Server. The interruption was caused by connection reset and it happens randomly.

To avoid this I set some parameters on JDBC connection string. The error was stopped but the task never ends.

  • original:

jdbc:sqlserver://host:port;database=db_name;

  • modified:
db_url=jdbc:sqlserver://host:port;
database=db_name;
applicationIntent=readonly;
applicationName=app-name;
columnEncryptionSetting=Disabled;
disableStatementPooling=true;
encrypt=false;
integratedSecurity=false;
lastUpdateCount=true;
lockTimeout=-1;
loginTimeout=15;
multiSubnetFailover=false;
packetSize=8000;
queryTimeout=-1;
responseBuffering=adaptive;
selectMethod=direct;
sendStringParametersAsUnicode=true;
serverNameAsACE=false;
TransparentNetworkIPResolution=true;
trustServerCertificate=false;
trustStoreType=JKS;
sendTimeAsDatetime=true;
xopenStates=false;
authenticationScheme=nativeAuthentication;
authentication=NotSpecified;
socketTimeout=0;
fips=false;
enablePrepareOnFirstPreparedStatementCall=false;
serverPreparedStatementDiscardThreshold=10;
statementPoolingCacheSize=0;
jaasConfigurationName=SQLJDBCDriver;
sslProtocol=TLS;
cancelQueryTimeout=-1;
useBulkCopyForBatchInsert=false;

Thus, I decided to remove the added parameters on JDBC connection string and started to pass spark configuration on cluster creation. I changed the max number of retries from 4 (default) to 50.

  • spark.task.maxFailures=50

So, the connection problem persists but at least the task successfully ends.

I would suggest you to set any connection timeout because maybe it's unlimited - usually set as 0 or -1. Check Oracle's driver documentation and try to change the default behavior.

2 - Problem related to S3

We also faced a problem related to write operation on S3. I couldn't find the exact error message but it was something similar to An error occurred while calling o70.parquet.

When we solved the mentioned problem the writing speed was taking too long.

One guy from our team suggested to use HDFS to write data from database and then copy operation from HDFS to S3. The performance increased substantially.

  • set HDFS as destination (maybe will need to increase disk size from master node)
destination = 'hdfs:///path-to-hdfs'

DF.write                \
  .mode("overwrite")    \
  .parquet(destination)
  • perform copy from HDFS to S3

Reference: https://docs.aws.amazon.com/emr/latest/ReleaseGuide/UsingEMR_s3distcp.html

I hope this helps you! I'll report any improvement from my task ;)

Related