Zły problem z kafką i Spark Streaming w Pythonie

Oct 29 2020

Uwaga: NIE jest to ten sam problem, który miałem w moim pierwszym poście na tej stronie, jednak jest to ten sam projekt.

Wchodzę do PostgreSQL z kafki za pomocą przesyłania strumieniowego Spark. Oto moje kroki w projekcie:

1- Tworzenie skryptu dla producenta kafki (gotowe, działa dobrze)

2- Stworzenie skryptu w Pythonie, który czyta pliki od producenta kafka

3- Wysyłanie plików do PostgreSQL

Do połączenia między pythonem i postgreSQL używam psycopg2. Używam również pythona 3 i java jdk1.8.0_261, a integracja między kafką a strumieniowaniem iskra działa dobrze. Mam kafka 2.12-2.6.0 i Spark 3.0.1 i dodałem te słoiki do mojego katalogu Spark jars:

  • 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

Musiałem również pobrać VC ++, aby naprawić inny problem, również związany z moim projektem.

To jest mój fragment kodu Pythona, który pobiera pliki od producenta kafki i wysyła je do tabeli postgreSQL, którą utworzyłem w postgreSQL, z którą mam problemy:

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

SatelliteTable to iskra dataframe, którą utworzyłem z plików od producenta kafka. process_row to funkcja, która wstawia każdy wiersz ramki danych strumieniowych do tabeli postgre. Oto ona:

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

Problem, który pojawia się po uruchomieniu kodu, pojawia się w query = satelliteTable.writeStream.outputMode("append").foreachBatch(process_row) \ .option("checkpointLocation", "C:\\Users\\Vito\\Documents\\popo").start()iw skrócie jest następujący:

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

Zabawne jest to, że ten sam kod działa dobrze na laptopie mojego przyjaciela, z iskrą 3.0.0. Myślę więc, że brakuje mi kilku słoików lub innych rzeczy, ponieważ kod jest poprawny.

Dowolny pomysł? Dzięki.

Odpowiedzi

Matt Oct 29 2020 at 17:05

brakuje ci tego słoika https://mvnrepository.com/artifact/org.apache.commons/commons-pool2Wypróbuj tę konkretną wersjęhttps://mvnrepository.com/artifact/org.apache.commons/commons-pool2/2.6.2