Собирайте значения как словарь в родительском столбце с помощью Pyspark.Python

Программы на Python
Anonymous
Собирайте значения как словарь в родительском столбце с помощью Pyspark.

Сообщение Anonymous »

У меня есть код и данные, как показано ниже:

Код: Выделить всё

df_renamed = df.withColumnRenamed("id","steps.id").withColumnRenamed("status_1","steps.status").withColumnRenamed("severity","steps.error.severity")

df_renamed.show(truncate=False)

+----------+-------+------+-----------------------+------------+--------------------+
|apiVersion|expired|status|steps.id               |steps.status|steps.error.severity|
+----------+-------+------+-----------------------+------------+--------------------+
|2         |false  |200   |mexican-curp-validation|200         |null                |
+----------+-------+------+-----------------------+------------+--------------------+
Теперь я хочу преобразовать эти данные, как показано ниже:

Код: Выделить всё

+----------+-------+------+-----------------------+------------+--------------------+
|apiVersion|expired|status|steps                                                    |
+----------+-------+------+-----------------------+------------+--------------------+
|2         |false  |200   |{"id":"mexican-curp-validation", "status":200 ,"error":{"severity":null}}               |
+----------+-------+------+-----------------------+------------+--------------------+
где видно, что на основе точечной записи названий столбцов в данных формируется JSON-структура. По этой причине я использовал приведенный ниже код:

Код: Выделить всё

cols_list = [name for name in df_renamed.columns if "." in name]
df_new = df_renamed.withColumn("steps",F.to_json(F.struct(*cols_list)))
df_new.show()
Но выдает ошибку ниже, хотя столбец присутствует:

Код: Выделить всё

 df_new = df_renamed.withColumn("steps",F.to_json(F.struct(*cols_list)))
File "/Users/../IdeaProjects/pocs/venvsd/lib/python3.9/site-packages/pyspark/sql/dataframe.py", line 3036, in withColumn
return DataFrame(self._jdf.withColumn(colName, col._jc), self.sparkSession)
File "/Users/../IdeaProjects/pocs/venvsd/lib/python3.9/site-packages/py4j/java_gateway.py", line 1321, in __call__
return_value = get_return_value(
File "/Users/../IdeaProjects/pocs/venvsd/lib/python3.9/site-packages/pyspark/sql/utils.py", line 196, in deco
raise converted from None
pyspark.sql.utils.AnalysisException: Column 'steps.id' does not exist. Did you mean one of the following? [steps.id, expired, status, steps.status, apiVersion, steps.error.severity];
'Project [apiVersion#17, expired#18, status#19, steps.id#29, steps.status#37, steps.error.severity#44, to_json(struct(id, 'steps.id, status, 'steps.status, severity, 'steps.error.severity), Some(GMT+05:30)) AS steps#82]
+- Project [apiVersion#17, expired#18, status#19, steps.id#29, steps.status#37, severity#22 AS steps.error.severity#44]
+- Project [apiVersion#17, expired#18, status#19, steps.id#29, status_1#21 AS steps.status#37, severity#22]
+- Project [apiVersion#17, expired#18, status#19, id#20 AS steps.id#29, status_1#21, severity#22]
+- Relation [apiVersion#17,expired#18,status#19,id#20,status_1#21,severity#22] csv
Где я ошибаюсь? Любая помощь очень ценится.

Подробнее здесь: https://stackoverflow.com/questions/787 ... ng-pyspark

Вернуться в «Python»