Код: Выделить всё
from pyspark.sql import SparkSession
spark = SparkSession.builder.config("spark.driver.host", "localhost").appName('SparkByExamples.com').getOrCreate()
data = [
('James','Smith','M',3000),
('James','Smith','M',3000),
('James','Smith','M',3000),
('James','Smith','M',3000),
('Anna','Rose','F',4100),
('Anna','Rose','F',4100),
('Anna','Rose','F',4100),
('Anna','Rose','F',4100),
('Robert','Williams','M',6200),
('Robert','Williams','M',6200),
('Robert','Williams','M',6200),
('Robert','Williams','M',6200),
]
columns = ["firstname","lastname","gender","salary"]
df = spark.createDataFrame(data=data, schema = columns)
df.show()
#Example 1 mapPartitions()
def reformat(partitionData):
for row in partitionData:
yield [row.firstname+","+row.lastname,row.salary*10/100]
df2=df.repartition(4).rdd.mapPartitions(reformat).toDF(["name","bonus"])
df2.cache()
df2.show()
Я тестировал приведенный выше пример с более простой логикой, но не обнаружил проблемы.
Подробнее здесь: https://stackoverflow.com/questions/793 ... n-expected