Spark'ta bir python Apache Beam Ardışık Düzeni Çalıştırma

Oct 31 2020

Apache ışınını (python sdk ile) burada bir deneme yapıyorum, bu yüzden basit bir ardışık düzen oluşturdum ve onu bir Spark kümesine dağıtmaya çalıştım.

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)

Bu ardışık düzen DirectRunner ile iyi çalışıyor. Yani aynı kodu Spark'a dağıtmak için (çünkü taşınabilirlik Beam'de anahtar bir kavramdır) ...

İlk PipelineOptionsönce burada belirtildiği gibi düzenledim :

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

job_endpointşu komutu kullanarak çalıştırdığım beam spark iş sunucusunun docker konteynerinin url'sidir :

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

Bunun iyi çalışması gerekiyordu ancak iş Spark'ta şu hatayla başarısız oluyor:

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

Ayrıca, beam_spark_job_servergünlüklerde bu UYARI var :

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

Buradaki problemin nerede olduğuna dair bir fikriniz var mı? Python Beam Boru Hatlarını container mimarisine alınmış bir hizmetten geçmeden kıvılcım üzerinde çalıştırmanın başka bir yolu var mı?

Yanıtlar

2 KennKnowles Nov 04 2020 at 01:12

Bu, iş sunucusunda bulunan Spark istemcisinin sürümü ile işi gönderdiğiniz Spark sürümü arasındaki sürüm uyuşmazlığı nedeniyle olabilir.