A data engineer is working with a PySpark DataFrame named `transactions_df` that contains `transaction_id` (string), `product_category` (string), and `amount` (double). The engineer needs to calculate the total `amount` for each `product_category` and then filter out categories where the total `amount` is less than 1000. Which PySpark operation sequence correctly achieves this?
- Atransactions_df.filter(F.col('amount') >= 1000).groupBy('product_category').agg(F.sum('amount').alias('total_amount'))
- Btransactions_df.where(F.col('amount') >= 1000).groupBy('product_category').agg(F.sum('amount').alias('total_amount'))
- Ctransactions_df.groupBy('product_category').agg(F.sum('amount').alias('total_amount')).filter(F.col('total_amount') >= 1000)
- Dtransactions_df.groupBy('product_category').agg(F.sum('amount').alias('total_amount')).having(F.col('total_amount') >= 1000)
Show answer & explanationAnswer & explanation
Correct answer: C. transactions_df.groupBy('product_category').agg(F.sum('amount').alias('total_amount')).filter(F.col('total_amount') >= 1000)
To filter based on an aggregated value, the aggregation must occur first. After grouping and summing, the `filter()` method can be applied to the resulting DataFrame, using the alias of the aggregated column (`total_amount`). PySpark's `filter()` serves the purpose of both SQL `WHERE` and `HAVING` clauses depending on when it's applied.
Why the other options are wrong
- A. Filtering `amount` before aggregation would remove individual transactions, not categories based on their total. The filter needs to happen after the sum.
- B. `where()` is an alias for `filter()`. Similar to B, filtering individual rows before aggregation is incorrect for this requirement.
- D. PySpark DataFrames do not have a `.having()` method directly. The `filter()` method is used for both pre-aggregation (SQL `WHERE`) and post-aggregation (SQL `HAVING`) filtering.
PySpark Filtering Aggregated Data
In PySpark, filtering based on an aggregated value (equivalent to SQL's `HAVING` clause) is achieved by first performing the `groupBy()` and `agg()` operations, and then applying a `filter()` transformation on the resulting DataFrame using the aggregated column.
- Aggregation (`groupBy().agg()`) must precede the filter.
- The `filter()` method is used for both pre- and post-aggregation filtering.
- Refer to the aggregated column by its alias within the `filter()` clause.
- There is no direct `.having()` method in PySpark DataFrames.
Memory trick: Group, then aggregate, then filter the aggregates!