【问题标题】:PySpark - Fill missing data with multiple columns as keysPySpark - 用多列作为键填充缺失的数据
【发布时间】:2021-06-02 01:28:15
【问题描述】:

我有一个包含 4 列的 spark 数据框。

ts -> long (unix timestamp)
col1 -> string
col2 -> string
value -> long

tscol1col2 的组合在我的数据中是独一无二的。我想通过创建缺少 tscol1col2 以及前 3 列的所有组合的行来填充缺失的数据(ts 具有特定范围,col1col2 具有离散值列表)。我能想到的唯一方法是创建一个包含 3 列的所有有效组合的新数据框,并将 value 列设置为 0,然后以某种方式合并 2 个数据框。 这就是我目前所拥有的

partial_data_df = spark.read.csv(my_path)

TS_DF = spark.range(min_ts, max_ts, 1000 * 3600).select(F.col('id').alias('ts')).orderBy('ts')
COL1_DF = spark.createDataFrame([..some data..], schema=['col1'])
COL2_DF = spark.createDataFrame([..some data..], schema=['col2'])
EMPTY_DF = TS_DF.crossJoin(COL1_DF).crossJoin(COL2_DF).withColumn('value', F.lit(0))

# now what?

我如何在 3 列上合并 partial_data_dfEMPTY_DF,这样如果存在组合,则从 partial_data_df 中取出 value 列,如果不存在则输入 0?是否有另一种方式(更优雅)来实现这一目标?

编辑

我尝试像这样进行左连接(我按照建议从EMPTY_DF 中删除了value 列)

merged_df = EMPTY_DF.join(partial_data_df, (
                     (partial_data_df.ts == EMPTY_DF.ets) & 
                     (partial_data_df.col1 == EMPTY_DF.ecol1) &
                     (partial_data_df.col2 == EMPTY_DF.ecol2)
                    ), how='left')
         .select(
           F.col('ets').alias('ts'), 
           F.col('ecol1').alias('col1'), 
           F.col('ecol2').alias('col2'),
           F.when(((F.col('value').isNull()) | (F.col('value') == 0)), 0).otherwise(F.col('value')).alias('value')
         )

但是行数没有加起来

row count in EMPTY_DF is 778176
row count in partial_data_df is 131709
row count in merged_df 778176
row count in merged_df that has non zero volume 100348
row count in partial_data_df that has non zero volume 131709
count distinct (ts, col1, col2) on partial_data_df 131709
count distinct (ts, col1, col2) on merged_df 778176

这里有什么问题?

【问题讨论】:

  • partial_data_df 具有非零值的行数是多少?
  • @ggordon - 请查看问题中更新的行数 - 有些东西没有加起来
  • 可能是某些实际值与生成的值不匹配,例如,如果您的实际 ts 值之一是 1 second off the hour 。您能否使用partial_data_df.select(F.expr("CASE WHEN mod(value,1000 * 3600) = 0 THEN 0 ELSE 1 END").alias('is_not_consistent')).agg(F.sum('is_not_consistent')) 确认这些值是否存在。

标签: python apache-spark pyspark


【解决方案1】:

您可以继续使用

方法 1

...表别名和when 函数使用原始数据帧中的值(如果可用)。使用这种方法,您不需要在内存中使用.withColumn('value', F.lit(0)) 生成所有可能的空值,因为交叉连接会产生我们可以使用的NULL

from pyspark.sql.functions import when, col

partial_data_df = partial_data_df.alias('original_df')
EMPTY_DF = EMPTY_DF.alias('empty_df')

final_df = EMPTY_DF.join(partial_data_df,(
    col('original_df.ts') == col('empty_df.ts') & 
    col('original_df.col1') == col('empty_df.col1') & 
    col('original_df.col2') == col('empty_df.col2') & 
),"left").select(
    col('original_df.ts'),
    col('original_df.col1'),
    col('original_df.col2'),
    when(col('original_df.value').isNull(),0).otherwise(col('original_df.value')).alias('value')
)

# Or Since both columns exist in both dataframes
final_df = EMPTY_DF.join(partial_data_df,['ts','col1','col2'],"left").select(
    col('empty_df.ts'),
    col('empty_df.col1'),
    col('empty_df.col2'),
    when(col('original_df.value').isNull(),0).otherwise(col('original_df.value')).alias('value')
)

方法2

另一种方法是使用 spark-sql,例如。

partial_data_df = partial_data_df.createOrReplaceTempView('original_df')
EMPTY_DF = EMPTY_DF.createOrReplaceTempView('empty_df')

final_df = sparkSession.sql("""
SELECT
    e.ts,
    e.col1,
    e.col2,
    CASE 
        WHEN o.value IS NULL THEN 0
        ELSE o.value
    END as value
FROM
    empty_df e
LEFT JOIN
    original_df o ON o.ts = e.ts AND 
                     o.col1 = e.col1 AND 
                     o.col2=e.COL2
""")

参考文献

【讨论】:

  • 谢谢 - 我试图这样做(我用我的代码编辑 OP)但最终的行数不匹配。我知道merged_df 具有empty_df 的确切行数(这是预期的并且很好),但value 列中非0 值的行数小于partial_data_df 的行数。知道这是为什么吗?
猜你喜欢
  • 2018-09-15
  • 1970-01-01
  • 1970-01-01
  • 2021-03-11
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多