ML Engineer MasterClass (October) | 4 seats left

Spark · Transformations and actions
AmazonAmazon Analytics
Amazon · Meet Spark

Transformations and actions

Build a customer-directory plan and execute it with show or count.

Step 1 of 6 · Learn

Build a plan

filter and select describe transformations; assigning planned does not materialize rows. show executes enough of the plan to display a preview. This example selects IDs from 105 onward. An ID threshold is a selection rule, not a measure of customer age or value.

Lesson reference: PySpark and Scala

Build a plan

filter and select describe transformations; assigning planned does not materialize rows. show executes enough of the plan to display a preview. This example selects IDs from 105 onward. An ID threshold is a selection rule, not a measure of customer age or value.

PySpark example

from pyspark.sql.functions import col, when

planned = customers.filter(col("customer_id") >= 105)
result = planned.select("customer_id", "city").orderBy("customer_id")
result.show()

Scala example

import org.apache.spark.sql.functions._

val planned = customers.filter(col("customer_id") >= 105)
val result = planned.select("customer_id", "city").orderBy("customer_id")
result.show(false)

Trigger a count

count is an action that returns a number, not a DataFrame. The example counts the two customers selected by the plan. Displaying a preview does not cache the plan: a later count can execute it again.

PySpark example

from pyspark.sql.functions import col

planned = customers.filter(col("customer_id") >= 105)
customer_count = planned.count()
print(customer_count)

Scala example

import org.apache.spark.sql.functions._

val planned = customers.filter(col("customer_id") >= 105)
val customer_count = planned.count()
println(customer_count)

Extend the plan

Add orderBy and limit before show so Spark can optimize the whole plan. The example previews two customers alphabetically by city, breaking ties by customer_id. The Studio also collects result to render its table; without caching this can trigger more work.

PySpark example

from pyspark.sql.functions import col, when

planned = customers.select("customer_id", "city")
result = planned.orderBy(col("city").asc(), col("customer_id").asc()).limit(2)
result.show()

Scala example

import org.apache.spark.sql.functions._

val planned = customers.select("customer_id", "city")
val result = planned.orderBy(col("city").asc(), col("customer_id").asc()).limit(2)
result.show(false)
example.pyPySpark
1. Read customersStart with 6 records; preserve the original DataFrame.
2. Build a planTransformations build a plan; the action triggers execution.
3. Inspect the resultEli and Fran pass the example filter. The directory still contains six customers.
Source customers6 rows
customer_idcustomer_namecity
101AriSeattle
102BoAustin
103CamChicago
104DeeBoston
105EliDenver
106FranPortland
Build a plan
Result2 rows
customer_idcity
105Denver
106Portland
Eli and Fran pass the example filter. The directory still contains six customers.

The example is loaded in the editor. Run it as written, then try a small change.

Runs on the Spark backend. First startup may take a moment.

Run your code to see the result.