Ausführen einer Python-Apache-Beam-Pipeline unter Spark
Ich probiere hier Apache Beam (mit Python SDK) aus, also habe ich eine einfache Pipeline erstellt und versucht, sie auf einem Spark-Cluster bereitzustellen.
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)
Diese Pipeline funktioniert gut mit DirectRunner. Bereitstellen des gleichen Codes auf Spark (da die Portabilität ein Schlüsselkonzept in Beam ist) ...
Zuerst habe ich das PipelineOptionswie hier beschrieben bearbeitet :
op = PipelineOptions([
"--runner=PortableRunner",
"--job_endpoint=localhost:8099",
"--environment_type=LOOPBACK"
]
)
job_endpointist die URL zum Docker-Container des Beam-Spark-Job-Servers , den ich mit dem folgenden Befehl ausführe:
docker run --net=host apache/beam_spark_job_server:latest --spark-master-url=spark://SPARK_URL:SPARK_PORT
Dies soll gut funktionieren, aber der Job schlägt bei Spark mit folgendem Fehler fehl:
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
Außerdem habe ich diese WARN in den beam_spark_job_serverProtokollen:
WARN org.apache.beam.runners.spark.translation.SparkContextFactory: Creating a new Spark Context.
Irgendeine Idee, wo das Problem hier liegt? Gibt es eine andere Möglichkeit, Python Beam Pipelines mit Funken zu betreiben, ohne an einem Containerdienst vorbeizukommen?
Antworten
Dies kann auf eine Versionsinkongruenz zwischen der auf dem Jobserver enthaltenen Version des Spark-Clients und der Version von Spark zurückzuführen sein, an die Sie den Job senden.