Ein schlechtes Problem mit Kafka und Spark Streaming auf Python

Oct 29 2020

NB Dies ist NICHT das gleiche Problem, das ich in meinem ersten Beitrag auf dieser Site hatte, es ist jedoch das gleiche Projekt.

Ich nehme einige Dateien von Kafka mithilfe von Spark-Streaming in PostgreSQL auf. Dies sind meine Schritte für das Projekt:

1- Erstellen eines Skripts für den Kafka-Produzenten (fertig, es funktioniert einwandfrei)

2- Erstellen eines Python-Skripts, das Dateien vom Kafka-Produzenten liest

3- Senden von Dateien an PostgreSQL

Für die Verbindung zwischen Python und PostgreSQL verwende ich psycopg2. Ich benutze auch Python 3 und Java jdk1.8.0_261 und die Integration zwischen Kafka und Spark Streaming funktioniert gut. Ich habe kafka 2.12-2.6.0 und spark 3.0.1 und habe diese Gläser in mein Spark-Gläserverzeichnis aufgenommen:

  • postgresql-42.2.18 -spark-Streaming-Kafka-0-10-Assembly_2.12-3.0.1
  • spark-token-provider-kafka-0.10_2.12-3.0.1
  • kafka-clients-2.6.0
  • spark-sql-kafka-0-10-assembly_2.12-3.0.1

Ich musste auch VC ++ herunterladen, um ein anderes Problem zu beheben, das auch mit meinem Projekt zusammenhängt.

Dies ist mein Teil des Python-Codes, der Dateien vom kafka-Produzenten nimmt und sie in eine Tabelle von postgreSQL sendet, die ich in postgreSQL erstellt habe und bei der ich Probleme habe:

query = satelliteTable.writeStream.outputMode("append").foreachBatch(process_row) \
.option("checkpointLocation", "C:\\Users\\Vito\\Documents\\popo").start()
print("Starting")
print(query)
query.awaitTermination()
query.stop()

satellitentabelle ist der Funken-Datenrahmen, den ich mit Dateien vom kafka-Produzenten erstellt habe. process_row ist die Funktion, die jede Zeile des Streaming-Datenrahmens in die Postgre-Tabelle einfügt. Hier ist es:

def process_row(df, epoch_id):
for row in df.rdd.collect():
    cursor1.execute(
        'INSERT INTO satellite(filename,satellite_prn_number, date, time,crs,delta_n, m0, 
                   cuc,e_eccentricity,cus,'
        'sqrt_a, toe_time_of_ephemeris, cic, omega_maiusc, cis, i0, crc, omega, omega_dot, idot) 
                    VALUES (%s,%s,%s,'
        '%s,%s,%s, %s, %s, %s, %s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)', row)
connection.commit()
pass

Das Problem, das beim Ausführen meines Codes auftritt, tritt bei query = satelliteTable.writeStream.outputMode("append").foreachBatch(process_row) \ .option("checkpointLocation", "C:\\Users\\Vito\\Documents\\popo").start()und kurz gesagt wie folgt auf:

py4j.protocol.Py4JJavaError: An error occurred while calling 
z:org.apache.spark.api.python.PythonRDD.collectAndServe.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 0.0 failed 1 
times, most recent failure: Lost task 0.0 in stage 0.0 (TID 0, DESKTOP- 
D600TY.homenet.telecomitalia.it, executor driver): java.lang.NoClassDefFoundError: 
org/apache/commons/pool2/impl/GenericKeyedObjectPoolConfig

=== Streaming Query ===
Identifier: [id = 599f75a7-5db6-426e-9082-7fbbf5196db9, runId = 67693586-27b1-4ca7-9a44-0f69ad90eafe]
Current Committed Offsets: {}
Current Available Offsets: {KafkaV2[Subscribe[bogi2890.20n]]: {"bogi2890.20n":{"0":68}}}

Current State: ACTIVE
Thread State: RUNNABLE

Die lustige Tatsache ist, dass der gleiche Code auf dem Laptop meines Freundes mit Spark 3.0.0 einwandfrei funktioniert. Ich denke also, dass mir einige Gläser oder anderes fehlen, weil der Code korrekt ist.

Irgendeine Idee? Vielen Dank.

Antworten

Matt Oct 29 2020 at 17:05

Sie vermissen dieses Glas https://mvnrepository.com/artifact/org.apache.commons/commons-pool2Probieren Sie diese spezielle Version aushttps://mvnrepository.com/artifact/org.apache.commons/commons-pool2/2.6.2