How to use Prefect's resource manager with a spark cluster

Viewed 918

I have been messing around with Prefect for workflow management, but got stuck with building up and braking down a spark session withing Prefect's resource manager.

I browsed Prefects docs and an example with Dusk is available:

from prefect import resource_manager
from dask.distributed import Client

@resource_manager
class DaskCluster:
    def init(self, n_workers):
        self.n_workers = n_workers

    def setup(self):
        "Create a local dask cluster"
        return Client(n_workers=self.n_workers)

    def cleanup(self, client):
        "Cleanup the local dask cluster"
        client.close()
        
        
with Flow("example") as flow:
    n_workers = Parameter("n_workers")

    with DaskCluster(n_workers=n_workers) as client:
        some_task(client)
        some_other_task(client)        

However I couldn't work out how to do the same with a spark session.

1 Answers

The simplest way to do this is with Spark in local mode:

from prefect import task, Flow, resource_manager

from pyspark import SparkConf
from pyspark.sql import SparkSession

@resource_manager
class SparkCluster:
    def __init__(self, conf: SparkConf = SparkConf()):
        self.conf = conf

    def setup(self) -> SparkSession:
        return SparkSession.builder.config(conf=self.conf).getOrCreate()

    def cleanup(self, spark: SparkSession):
        spark.stop()

@task
def get_data(spark: SparkSession):
    return spark.createDataFrame([('look',), ('spark',), ('tutorial',), ('spark',), ('look', ), ('python', )], ['word'])

@task(log_stdout=True)
def analyze(df):
    word_count = df.groupBy('word').count()
    word_count.show()


with Flow("spark_flow") as flow:
    conf = SparkConf().setMaster('local[*]')
    with SparkCluster(conf) as spark:
        df = get_data(spark)
        analyze(df)

if __name__ == '__main__':
    flow.run()

Your setup() method returns the resource being managed and the cleanup() method accepts the same resource returned by setup(). In this case, we create and return a Spark session, and then stop it. You don't need spark-submit or anything (though I find it a bit harder managing dependencies this way).

Scaling it up gets harder and is something I'm still working on figuring out. For example, Prefect won't know how to serialize Spark DataFrames for output caching or persisting results. Also, you have to be careful about using the Dask executor with Spark sessions because they can't be pickled, so you have to set the executor to use scheduler='threads' (see here).

Related