【发布时间】:2021-06-02 01:28:15
【问题描述】:
我有一个包含 4 列的 spark 数据框。
ts -> long (unix timestamp)
col1 -> string
col2 -> string
value -> long
ts、col1 和 col2 的组合在我的数据中是独一无二的。我想通过创建缺少 ts、col1 和 col2 以及前 3 列的所有组合的行来填充缺失的数据(ts 具有特定范围,col1 和 col2 具有离散值列表)。我能想到的唯一方法是创建一个包含 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_df 和 EMPTY_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