Ejecución de una canalización de Python Apache Beam en Spark

Oct 31 2020

Estoy probando apache beam (con python sdk) aquí, así que creé una canalización simple e intenté implementarla en un clúster 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)

Esta canalización funciona bien con DirectRunner. Entonces, para implementar el mismo código en Spark (ya que la portabilidad es un concepto clave en Beam) ...

Primero edité el PipelineOptionscomo se menciona aquí :

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

job_endpointes la URL del contenedor de la ventana acoplable del servidor de trabajo de beam spark que ejecuto usando el comando:

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

Se supone que esto funciona bien, pero el trabajo falla en Spark con este error:

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

Además, tengo esta ADVERTENCIA en los beam_spark_job_serverregistros:

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

¿Alguna idea de dónde está el problema aquí? ¿Hay alguna otra forma de ejecutar Python Beam Pipelines en Spark sin pasar por un servicio en contenedores?

Respuestas

2 KennKnowles Nov 04 2020 at 01:12

Esto podría suceder debido a una discrepancia de versión entre la versión del cliente Spark contenida en el servidor de trabajo y la versión de Spark a la que está enviando el trabajo.