mapInPandas

Pandas DataFrame'lerde hem giriş hem de çıkış olarak gerçekleştirilen Python yerel bir işlev kullanarak geçerli DataFrame'deki toplu işlemlerin yineleyicisini eşler ve sonucu DataFrame olarak döndürür.

Sözdizimi

mapInPandas(func: "PandasMapIterFunction", schema: Union[StructType, str], barrier: bool = False, profile: Optional[ResourceProfile] = None)

Parametreler

Parametre Türü Açıklama
func function pandas.DataFrames yineleyicisini alan ve pandas.DataFrames yineleyicisi veren Python yerel işlevi.
schema DataType veya str PySpark'ta değerinin func dönüş türü. Değer bir pyspark.sql.types.DataType nesne veya DDL biçimli bir tür dizesi olabilir.
barrier bool, isteğe bağlı, varsayılan False Aşamadaki tüm Python çalışanlarının eşzamanlı olarak başlatılmasını sağlamak için engel modu yürütmeyi kullanın.
profile ResourceProfile, isteğe bağlı mapInPandas için kullanılacak isteğe bağlı ResourceProfile.

İadeler

DataFrame

Örnekler

df = spark.createDataFrame([(1, 21), (2, 30)], ("id", "age"))

def filter_func(iterator):
    for pdf in iterator:
        yield pdf[pdf.id == 1]

df.mapInPandas(filter_func, df.schema).show()
# +---+---+
# | id|age|
# +---+---+
# |  1| 21|
# +---+---+

def mean_age(iterator):
    for pdf in iterator:
        yield pdf.groupby("id").mean().reset_index()

df.mapInPandas(mean_age, "id: bigint, age: double").show()
# +---+----+
# | id| age|
# +---+----+
# |  1|21.0|
# |  2|30.0|
# +---+----+

df.mapInPandas(filter_func, df.schema, barrier=True).collect()
# [Row(id=1, age=21)]