Menjalankan python Apache Beam Pipeline di Spark

Oct 31 2020

Saya memberikan apache beam (dengan python sdk) untuk dicoba di sini, jadi saya membuat pipeline sederhana dan saya mencoba menerapkannya pada 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)

Pipa ini bekerja dengan baik dengan DirectRunner. Jadi untuk menerapkan kode yang sama pada Spark (karena portabilitas adalah konsep utama di Beam) ...

Pertama saya mengedit PipelineOptionsseperti yang disebutkan di sini :

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

job_endpointadalah url ke kontainer buruh pelabuhan dari server pekerjaan percikan berkas yang saya jalankan menggunakan perintah:

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

Ini seharusnya berfungsi dengan baik tetapi pekerjaan gagal di Spark dengan kesalahan ini:

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

Juga, saya memiliki PERINGATAN ini di beam_spark_job_serverlog:

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

Ada ide di mana masalahnya di sini? Apakah ada cara lain untuk menjalankan python Beam Pipelines pada percikan api tanpa melewati layanan dalam peti kemas?

Jawaban

2 KennKnowles Nov 04 2020 at 01:12

Hal ini dapat terjadi karena ketidakcocokan versi antara versi klien Spark yang terdapat dalam server pekerjaan dan versi Spark yang Anda kirimkan pekerjaan tersebut.