【发布时间】:2021-09-28 01:48:59
【问题描述】:
总的来说,我是 PySpark 和 Spark 的新手。 我想对 DataFrame 中的给定列应用转换,本质上为该特定列上的每个值调用一个函数。
我的 DataFrame df 看起来像这样:
df.show()
+------------+--------------------+
|version | body |
+------------+--------------------+
| 1|9gIAAAASAQAEAAAAA...|
| 2|2gIAAAASAQAEAAAAA...|
| 3|3gIAAAASAQAEAAAAA...|
| 1|7gIAKAASAQAEAAAAA...|
+------------+--------------------+
我需要为version 为1 的每一行读取body 列的值,然后对其进行解密(我有自己的逻辑/函数,它接受一个字符串并返回一个解密的字符串)。最后,将解密后的值以 csv 格式写入 S3 存储桶。
def decrypt(encrypted_string: str):
# code that returns decrypted string
所以,当我执行以下操作时,我会得到相应的过滤值,我需要应用我的解密函数。
df.where(col('version') =='1')\
.select(col('body')).show()
+--------------------+
| body|
+--------------------+
|9gIAAAASAQAEAAAAA...|
|7gIAKAASAQAEAAAAA...|
+--------------------+
但是,我不清楚如何做到这一点。我尝试使用collect(),但它违背了使用 Spark 的目的。
我也尝试如下使用.rdd.map,但没有奏效。
df.where(col('version') =='1')\
.select(col('body'))\
.rdd.map(lambda x: decrypt).toDF().show()
OR
.rdd.map(decrypt).toDF().show()
有人可以帮忙吗。
【问题讨论】:
标签: amazon-s3 pyspark apache-spark-sql