Esecuzione di una pipeline Python Apache Beam su Spark

Oct 31 2020

Sto provando Apache Beam (con Python SDK) qui, quindi ho creato una semplice pipeline e ho provato a distribuirla su un cluster Spark.

from apache_beam.options.pipeline_options import PipelineOptions
import apache_beam as beam

op = PipelineOptions([
        "--runner=DirectRunner"
    ]
)


with beam.Pipeline(options=op) as p:
    p | beam.Create([1, 2, 3]) | beam.Map(lambda x: x+1) | beam.Map(print)

Questa pipeline funziona bene con DirectRunner. Quindi, per distribuire lo stesso codice su Spark (poiché la portabilità è un concetto chiave in Beam) ...

Per prima cosa ho modificato il PipelineOptionscome menzionato qui :

op = PipelineOptions([
        "--runner=PortableRunner",
        "--job_endpoint=localhost:8099",
        "--environment_type=LOOPBACK"
    ]
)

job_endpointè l'URL del contenitore docker del server di lavoro beam spark che eseguo utilizzando il comando:

docker run --net=host apache/beam_spark_job_server:latest --spark-master-url=spark://SPARK_URL:SPARK_PORT

Questo dovrebbe funzionare bene ma il lavoro non riesce su Spark con questo errore:

20/10/31 14:35:58 ERROR TransportRequestHandler: Error while invoking RpcHandler#receive() for one-way message.

java.io.InvalidClassException: org.apache.spark.deploy.ApplicationDescription; local class incompatible: stream classdesc serialVersionUID = 6543101073799644159, local class serialVersionUID = 1574364215946805297

Inoltre, ho questo WARN nei beam_spark_job_serverlog:

WARN org.apache.beam.runners.spark.translation.SparkContextFactory: Creating a new Spark Context.

Qualche idea su dove sia il problema qui? C'è un altro modo per eseguire Python Beam Pipelines su Spark senza passare da un servizio containerizzato?

Risposte

2 KennKnowles Nov 04 2020 at 01:12

Ciò potrebbe accadere a causa di una mancata corrispondenza di versione tra la versione del client Spark contenuta nel job server e la versione di Spark a cui si sta inviando il lavoro.