Executando um pipeline do Apache Beam em Python no Spark

Oct 31 2020

Estou experimentando o Apache beam (com python sdk) aqui, então criei um pipeline simples e tentei implantá-lo em um 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)

Este pipeline está funcionando bem com o DirectRunner. Então, para implantar o mesmo código no Spark (já que a portabilidade é um conceito-chave no Beam) ...

Primeiro editei o PipelineOptionsconforme mencionado aqui :

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

job_endpointé o url para o contêiner do docker do servidor de jobs beam spark que executo usando o comando:

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

Isso deveria funcionar bem, mas a tarefa falha no Spark com este erro:

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

Além disso, tenho este AVISO nos beam_spark_job_serverlogs:

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

Alguma ideia de onde está o problema aqui? Existe alguma outra maneira de executar Python Beam Pipelines no Spark sem passar por um serviço em contêineres?

Respostas

2 KennKnowles Nov 04 2020 at 01:12

Isso pode acontecer devido a uma incompatibilidade de versão entre a versão do cliente Spark contida no servidor de trabalho e a versão do Spark para a qual você está enviando o trabalho.