Плохая проблема с kafka и Spark Streaming на Python

Oct 29 2020

NB. Это НЕ та же проблема, что и в моем первом посте на этом сайте, но это тот же проект.

Я загружаю некоторые файлы в PostgreSQL из kafka с помощью потоковой передачи искр. Вот мои шаги для проекта:

1- Создание скрипта для производителя kafka (готово, отлично работает)

2- Создание скрипта Python, который читает файлы от производителя kafka

3- Отправка файлов в PostgreSQL

Для связи между python и postgreSQL я использую psycopg2. Я также использую python 3 и java jdk1.8.0_261, и интеграция между kafka и потоковой передачей искры работает нормально. У меня есть kafka 2.12-2.6.0 и spark 3.0.1, и я добавил эти банки в свой каталог Spark jars:

  • postgresql-42.2.18 -spark-streaming-kafka-0-10-assembly_2.12-3.0.1
  • искра-токен-провайдер-кафка-0.10_2.12-3.0.1
  • кафка-клиенты-2.6.0
  • искра-sql-кафка-0-10-сборка_2.12-3.0.1

Мне также пришлось загрузить VC ++, чтобы исправить еще одну проблему, также связанную с моим проектом.

Это мой фрагмент кода Python, который берет файлы от производителя kafka и отправляет их в таблицу postgreSQL, которую я создал в postgreSQL, с которой у меня проблемы:

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 - это искровой фреймворк, который я создал из файлов от производителя kafka. process_row - это функция, которая вставляет каждую строку фрейма данных потоковой передачи в таблицу postgre. Вот:

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

Проблема, которую я получаю, когда я запускаю свой код, возникает, query = satelliteTable.writeStream.outputMode("append").foreachBatch(process_row) \ .option("checkpointLocation", "C:\\Users\\Vito\\Documents\\popo").start()и вкратце это следующая:

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

Интересно то, что тот же код отлично работает на ноутбуке моего друга с Spark 3.0.0. Итак, я думаю, что мне не хватает каких-то банок или чего-то еще, потому что код правильный.

Есть идеи? Благодарю.

Ответы

Matt Oct 29 2020 at 17:05

вам не хватает этой банки. https://mvnrepository.com/artifact/org.apache.commons/commons-pool2Попробуйте эту конкретную версию.https://mvnrepository.com/artifact/org.apache.commons/commons-pool2/2.6.2