I am trying to modify a Spark dataframe such that depending on one filtered value, we subset to several conditions and flag/modify the variable based on those conditions.
Ideally, I would like the finalized dataframe to be the same size as the original, just with the necessary mods.
Example below:
data = [
[1, "soda", "LB", 1, "L", 20],
[2, "juice", "KG", 1, "GA", 12],
[3, "water", "LB", 1, "L", 35],
[4, "soda", "G", 1, "M2", 11],
]
df = pd.DataFrame(
data, columns=["ID", "Beverage", "Weight", "Sample", "Volume", "Amount"]
)
drink_dictionary = {'soda': {"LB/L": 100, "G/M2": 200,},
"juice": {"KG/GA": 500, "LB/L": 90,},
"water": {'LB/L': 1,}}
sdf = spark.createDataFrame(df)
for drink in ["soda", "juice", "water"]:
for mass_unit in drink_dictionary[drink].keys():
weight_unit = mass_unit.split("/")[0]
volume_unit = mass_unit.split("/")[1]
value = drink_dictionary[drink][weight_unit + "/" + volume_unit]
# create condition, i.e specific weight and volumes,
# IF this condition is met, modify a different variable.
# Loop to next condition under same Beverage type.
# Modify accordingly.
condition = (spark_fns.col("Weight") == weight_unit) & (
spark_fns.col("Volume") == volume_unit
)
new_sdf = (
sdf.filter(spark_fns.col("Beverage") == drink)
.withColumn("Flag", spark_fns.when((condition), True).otherwise(False))
.withColumn(
"corrected_amount",
spark_fns.when(
(condition),
spark_fns.expr(f"Amount / Sample * {value}"),
).otherwise(sdf["Amount"]),
)
)
The output of this is not correct. What I would ideally like the output to look like is below (this would involve looping through all beverages):
ID Beverage Weight Sample Volume Amount Corrected_Amount
1 soda LB 1 L 20 200
2 juice KG 1 GA 12 6000
3 water LB 1 L 35 35
4 soda G 1 M2 11 2200