ML Engineer MasterClass (October) | 4 seats left

Spark · Arrays
AmazonAmazon Analytics
Amazon · Work with Nested Data

Arrays

Inspect array elements and apply expressions within each row.

Step 1 of 6 · Learn

Inspect elements

The setup creates quantities with item_count followed by item_count + 1. size counts elements. element_at uses 1-based indexing: 1 means the first element and -1 the last; zero is invalid. Every array here contains two elements, so both requested positions exist.

Lesson reference: PySpark and Scala

Inspect elements

The setup creates quantities with item_count followed by item_count + 1. size counts elements. element_at uses 1-based indexing: 1 means the first element and -1 the last; zero is invalid. Every array here contains two elements, so both requested positions exist.

PySpark example

from pyspark.sql.functions import col, lit, struct, array, size, element_at, filter, transform

array_orders = orders.select(
    "order_id", array(col("item_count"), col("item_count") + 1).alias("quantities")
)

result = array_orders.select(
    col("order_id"), size(col("quantities")).alias("quantity_count"),
    element_at(col("quantities"), 1).alias("selected_quantity")
).orderBy("order_id")
result.show(truncate=False)

Scala example

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

val array_orders = orders.select(
  col("order_id"), array(col("item_count"), col("item_count") + 1).alias("quantities")
)

val result = array_orders.select(
    col("order_id"), size(col("quantities")).alias("quantity_count"),
    element_at(col("quantities"), 1).alias("selected_quantity")
).orderBy("order_id")
result.show(false)

Filter inside an array

The functions filter expression retains matching elements within each array. It does not remove the DataFrame row. An array with no matching elements becomes empty. The lambda builds a Spark expression, not a Python UDF; use Spark-supported expressions inside it.

PySpark example

from pyspark.sql.functions import col, lit, struct, array, size, element_at, filter, transform

array_orders = orders.select(
    "order_id", array(col("item_count"), col("item_count") + 1).alias("quantities")
)

result = array_orders.select(
    col("order_id"), filter(col("quantities"), lambda x: x >= 2).alias("kept_quantities")
).orderBy("order_id")
result.show(truncate=False)

Scala example

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

val array_orders = orders.select(
  col("order_id"), array(col("item_count"), col("item_count") + 1).alias("quantities")
)

val result = array_orders.select(
  col("order_id"), filter(col("quantities"), x => x >= 2).alias("kept_quantities")
).orderBy("order_id")
result.show(false)

Transform each element

transform applies an expression to every element while retaining array length and order. The example adds 1 to each quantity. Unlike explode, which appears in the next lesson, this operation still returns one DataFrame row per order.

PySpark example

from pyspark.sql.functions import col, lit, struct, array, size, element_at, filter, transform

array_orders = orders.select(
    "order_id", array(col("item_count"), col("item_count") + 1).alias("quantities")
)

result = array_orders.select(
    col("order_id"), transform(col("quantities"), lambda x: x + 1).alias("updated_quantities")
).orderBy("order_id")
result.show(truncate=False)

Scala example

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

val array_orders = orders.select(
  col("order_id"), array(col("item_count"), col("item_count") + 1).alias("quantities")
)

val result = array_orders.select(
  col("order_id"), transform(col("quantities"), x => x + 1).alias("updated_quantities")
).orderBy("order_id")
result.show(false)
example.pyPySpark
One order rowArray expressionOne order row
Source array_orders6 rows
order_idquantities
1001[2,3]
1002[3,4]
1003[1,2]
1004[4,5]
1005[1,2]
1006[2,3]
Inspect elements
Result6 rows
order_idquantity_countselected_quantity
100122
100223
100321
100424
100521
100622
Compare the source columns with the transformed result.

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.