Я новичок в Spark.
Я запускаю Spark на локальном компьютере без кластера.
Я разработал следующий скрипт .
Цель состоит в том, чтобы получить уникальные значения ключей «лицензии».
Все идет далеко за рамки
licenses_collected = rdd.map(lambda x: x["license"]).collect()
unique_licenses = rdd.map(lambda x: x["license"]).distinct()
но когда я пытаюсь применить метод .distinct() к сопоставленному rdd, а затем применить метод .collect(), возникает Traceback.
Это строка, в которой возникает ошибка
unique_licenses_collected = rdd.map(lambda x: x["license"]).distinct().collect()
В чем проблема?
скрипт
from pyspark.sql import SparkSession
spark = SparkSession.builder.master("local").\
appName("prima applicazione spark").\
config("spark.ui.port", "4040").\
getOrCreate()
df_spark = spark.read.json(path_json)
rdd = df_spark.rdd
rdd_collected = rdd.collect()
"""Get the first two records"""
for line in rdd.take(2):
print(line)
print("\n\n")
""" Get the name of the licenses"""
licenses_collected = rdd.map(lambda x: x["license"]).collect()
unique_licenses = rdd.map(lambda x: x["license"]).distinct()
unique_licenses_collected = rdd.map(lambda x: x["license"]).distinct().collect()
Обратная трассировка
AttributeError: невозможно получить атрибут «PySparkRuntimeError» в
полная трассировка здесь
версии программ
pyspark --version
version 3.4.2
python --version
Python 3.10.12
Подробнее здесь: https://stackoverflow.com/questions/782 ... -apply-col