Как записать результаты функций в pyspark во внешние переменные?

Работаю на pyspark. Нужно пройтись по колонке датафрейма, применить некоторую функцию для каждой строки и записать вывод этой функции в стороннюю переменную.

Например:

есть DataFrame с названием df и колонкой Text

есть словарь dict, который имеет в качестве ключа гласную, а в качестве значения -сколько раз эта гласная встречалась - {'a':0,'e':4 и т.д.}

есть функция func, которая добавляет в словарь dict частотность гласных в этой строке. Например, dict['a']+=1

я пробовал выполнить задачу следующим образом:

Сделал так, чтобы func возвращал 0

инициализировал func_udf, которая возвращает IntegerType()

Далее использовал withColumn

df.withColumn("any",func_udf("Text",dict))

предполагалось, что функция пройдет по датафрейму, создаст фиктивную колонку с нулями и результаты операции запишет в dict

но в итоге dict не изменилась

еще был вариант записывать для каждой строки словарь в отдельную колонку т.е.

df.withColumn("column_with_dict",func_udf("Text"))

но как потом объединить результаты column_with_dict в один словарь, не используя collect, я не знаю

Вопросы такие:

  1. Можно ли в pyspark записать результаты во внешние переменные, не используя collect?

  2. Если нельзя, то как объединить результаты колонки column_with_dict в одно значение не используя collect

  3. Какие есть методы решения данной задачи?


Ответы (0 шт):