Writing to DB2 database gets stuck inside foreachpartition() from my Java code on Spark Cluster with 20 worker nodes configured

Viewed 235

I am trying to read data from a file uploaded on a server path, do some manipulation and then save this data to DB2 database. We have around 300K records which may increase further in future so we are trying to do all the manipulation as well write to DB2 inside foreachpartition. Below are the steps followed to do so.

  1. Create spark context as global and static.
static SparkContext sparkContext = new SparkContext();
static JavaSparkContext jc = sparkContext.getJavaSparkContext();
static SparkSession sc = sparkContext.getSparkSession();
  1. Create a dataset of file present on server
Dataset<Row> dataframe = sparkContext.getSparkSession().read().option("delimiter","|").option("header","false").option("inferSchema","false").schema(SchemaClass.generateSchema()).csv(filePath).withColumn("ID",monotonically_increasing_id())).withColumn("Partition_Id",spark_partition_id());
  1. Calling foreachpartition
dataframe.foreachpartition(new ForeachPartitionFunction<Row>()){
               @Override
               public void call(Iterator<Row> _row) throws Exception{
                List<SparkPositionDto> li = new ArrayList<>();
               while(_row.hasNext()){
               PositionDto positionDto = AnotherClass.method(row,1);
               SparkPositionDto spd = copyToSparkDto(positionDto);
               if(spd != null){
                li.add(spd);
                }
               }
             System.out.println("Writing via Spark : List size : "+li.size());
             JavaRDD<SparkPositionDto> finalRdd = jc.parallelize(li);
             Dataset<Row> dfToWrite = sc.createDataFrame(finalRDD, SparkPositionDto.class);
             System.out.println("Writing Data");   
             if(dfToWrite != null){
              dfToWrite.write()
                       .format("jdbc")
                       .option("url","jdbc:msdb2://"+"Database_Name"+";useKerberos=true")
                       .option("driver","DRIVER_NAME")
                       .option("dbtable","TABLE_NAME")
                       .mode(SaveMode.Append)
                       .save();
              }
            }
          }

The weird observation is that when i run this code outside foreachpartition for a small set of data, it works fine and in my spark cluster just 1 driver and 1 application runs, but when the same code is running inside foreachpartition, I could see 1 driver and 2 applications running with 1 app in running state and other in waiting. If I add numberOfPartitions as 5 in my schema then 5 applications can be seen running. It is running continuously, nothing in logs, seems it got stuck somewhere.

0 Answers
Related