spring boot with spark-submit cluster mode

Viewed 186

We are trying to achieve parallelism through the spark executors below are the steps which we are following -

  • Read from the hive
  • Data transformation (custom spring library).
  • Ship it via rest endpoint in batches (1000 records per batch).

Problem - We want to do all these steps in parallel and want to use a spring-boot based library.

Understanding - If we are using the custom code for transformation (the code that we want to run parallelly most probably inside rdd.map() method) then our classes and composite dependencies need to be serialized or those classes need to implement serialization.

We know we can achieve this by performing these tasks in sequence over the driver in such case we need to collect the data over the driver again + again and then pass it to the next step. In this case, we are not leveraging the power of executors.

Needs your assistance here -

If we ship this spring boot dependency to executors then is there any way that the executor understands the spring boot code and resolves the annotations over there ?

sample code -

Code from spring boot library -

public class Process{
   String convert(Row row) {
    return row.mkString();
  }
}

@Component
@ConditionalOnProperty(name = "process.dummy.serialize", havingValue = "true")
class ProcessNotSerialized extends Process {
    @Autowired
    private RecordService recorService; //not-serialized

    public void setName(String name) {
       this.name = name;
    }

    private String name;

    @Override
    public String toString() {
      return "ProcessNotSerialized";
    }
  }

Code from my spark spring boot application -

  Dataset<Row> sqlDf = sparkSession.sql(sqlQuery); // millions of data
  ProcessNotSerialized process = new ProcessNotSerialized();

    System.out.println("object name=>" + process.toString());
    List<String> listColumns = sqlDf.select(column)
            .javaRDD()
            .map(row -> {
                return process.convert(row);
            })
            .collect();

here is the code inside map() will execute in parallel.

Please let me know if you have any better way other than serialization or running over a driver.

0 Answers
Related