TGViewer
Data Engineers Data Engineers @sql_engineer · 11.2K subscribers
Post #168 3.73K
Working with PySpark Aggregations

What are Aggregations?

Aggregations in PySpark allow you to transform large datasets by computing statistics across specified groups. PySpark offers built-in functions for common aggregations, such as sum, avg, min, max, count, and more.

Common Aggregation Methods in PySpark

1. groupBy(): Groups data by one or more columns and allows applying aggregation functions on each group.

2. agg(): Lets you apply multiple aggregation functions simultaneously.

3. count(): Counts the number of non-null entries.

4. sum(): Adds up the values in a column.

5. avg(): Computes the average of a column.

Example: Using groupBy() and Aggregations

Let’s say you have a DataFrame with sales data and want to calculate the total and average sales per salesperson.

from pyspark.sql import SparkSession
from pyspark.sql.functions import sum, avg

# Create Spark session
spark = SparkSession.builder.appName("AggregationExample").getOrCreate()

# Sample data
data = [("Alice", 100), ("Alice", 150), ("Bob", 200), ("Bob", 300)]
df = spark.createDataFrame(data, ["Salesperson", "Sales_Amount"])

# Aggregating data
agg_df = df.groupBy("Salesperson").agg(
sum("Sales_Amount").alias("Total_Sales"),
avg("Sales_Amount").alias("Avg_Sales")
)

agg_df.show()

In this example, we used groupBy("Salesperson") to group the data by each salesperson, and agg() to calculate the total and average sales for each.

Real-World Example: Aggregating Product Sales Data

Imagine you're analyzing sales data for a retail store. You might want to know the total sales per product category, the highest and lowest sales amounts, or the average sales per transaction. Aggregations allow you to gain these insights quickly:

# Group by product category and calculate total and average sales
sales_df.groupBy("Product_Category").agg(
sum("Sales_Amount").alias("Total_Sales"),
avg("Sales_Amount").alias("Avg_Sales")
).show()

Advanced Aggregation Functions

countDistinct(): Counts unique values in a column.

df.groupBy("Salesperson").agg(countDistinct("Product_ID").alias("Unique_Products_Sold")).show()

approx_count_distinct(): Uses an approximate algorithm to count distinct values, useful for very large datasets.

from pyspark.sql.functions import approx_count_distinct
df.agg(approx_count_distinct("Product_ID")).show()

Windowed Aggregations

Sometimes, aggregations are performed over a “window” rather than over the entire dataset or specific groups. We’ve covered window functions, but it’s useful to know they can be combined with aggregations for tasks like rolling averages.
  • 👍 6
More from @sql_engineer
  1. Aug 29, 2026Example: Source Database → CDC → Only Changed Records → Data Platform CDC is especially us…
  2. Aug 29, 2026🚀 Data Engineering Fundamentals – Part 7 📥 Data Ingestion: How Data Enters a Data Platfo…
  3. Aug 18, 2026🚀 Data Engineering Fundamentals – Part 6 📌 ETL vs ELT: How Data Moves from Source to Des…
  4. Aug 11, 2026📊 The 90-Minutes Business Analytics Masterclass Learn how to transform raw data into powe…
  5. Aug 8, 2026Data Warehouse Stores: Cleaned sales data Customer KPIs Revenue reports Historical busines…
  6. Aug 8, 2026🚀 Data Engineering Fundamentals – Part 4 📌 Databases vs Data Warehouses vs Data Lakes vs…
Threads Profile ViewerView any public Threads profile without an account.Open ThreadLook →Writing with AI? Make it sound human.Metric37 rewrites AI drafts so they read naturally. Free AI detector, 1,500 words free.Try Metric37 →