【问题标题】:How to use chaining in pyspark?如何在pyspark中使用链接?
【发布时间】:2021-05-18 15:58:30
【问题描述】:

我有一个名为 Incito 的数据框,在该数据框的 Supplier Inv Nocolumn 中,由逗号分隔的值组成。我需要通过使用 pyspark 适当地重复这些逗号分隔值来重新创建数据框。我正在使用以下 python 代码。我可以将其转换为 pyspark 吗?可以通过 pyspark 吗?

from itertools import chain
def chainer(s):
    return list(chain.from_iterable(s.str.split(',')))
incito['Supplier Inv No'] = incito['Supplier Inv No'].astype(str)

# calculate lengths of splits
lens = incito['Supplier Inv No'].str.split(',').map(len)

# create new dataframe, repeating or chaining as appropriate
dfnew = pd.DataFrame({'Supplier Inv No': chainer(incito['Supplier Inv No']),
                      'Forwarder': np.repeat(incito['Forwarder'], lens),
                      'Mode': np.repeat(incito['Mode'], lens),
                      'File No': np.repeat(incito['File No'], lens),
                      'ETD': np.repeat(incito['ETD'], lens),
                      'Flight No': np.repeat(incito['Flight No'], lens),
                      'Shipped Country': np.repeat(incito['Shipped Country'], lens),
                      'Port': np.repeat(incito['Port'], lens),
                      'Delivered_Country': np.repeat(incito['Delivered_Country'], lens),
                      'AirWeight': np.repeat(incito['AirWeight'], lens),
                      'FREIGHT CHARGE': np.repeat(incito['FREIGHT CHARGE'], lens)})

这是我在 pyspark 中尝试过的。但我没有得到预期的结果。

from pyspark.context import SparkContext, SparkConf
from pyspark.sql.session import SparkSession
from pyspark.sql import functions as F
import pandas as pd

conf = SparkConf().setAppName("appName").setMaster("local")
sc = SparkContext(conf=conf)

spark = SparkSession(sc)
ddf = spark.createDataFrame(dfnew)


exploded = ddf.withColumn('d', F.explode("Supplier Inv No"))
exploded.show()

【问题讨论】:

  • 欢迎来到 SO,问题必须至少有一点尝试。请提供您迄今为止尝试过的内容。

标签: python-3.x pyspark itertools chaining


【解决方案1】:

类似这样的东西,使用repeat?

from pyspark.sql import functions as F

df = (spark
    .sparkContext
    .parallelize([
        ('ABCD',),
        ('EFGH',),
    ])
    .toDF(['col_a'])
)

(df
    .withColumn('col_b', F.repeat(F.col('col_a'), 2))
    .withColumn('col_c', F.repeat(F.lit('X'), 10))
    .show()
)
# +-----+--------+----------+
# |col_a|   col_b|     col_c|
# +-----+--------+----------+
# | ABCD|ABCDABCD|XXXXXXXXXX|
# | EFGH|EFGHEFGH|XXXXXXXXXX|
# +-----+--------+----------+

【讨论】:

    猜你喜欢
    • 2016-04-13
    • 2018-10-14
    • 1970-01-01
    • 2021-01-30
    • 2018-06-20
    • 1970-01-01
    • 2019-09-21
    相关资源
    最近更新 更多