Esecuzione di una pipeline Python Apache Beam su Spark
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
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.