Scala for Big Dataдля Big Data

A dense, skimmable cheatsheet for a Python/SQL data engineer bridging to JVM big-data work (Spark / Flink / Kafka / Iceberg). Fast-revision cheatsheet for Data Engineering interviews. Tailored for Senior/Lead DE.
Компактная шпаргалка для Python/SQL дата-инженера, переходящего на JVM big-data (Spark / Flink / Kafka / Iceberg). Быстрое повторение перед Data Engineering интервью. Заточена под Senior/Lead DE.
Your anchor: real production Scala + Spark on Hadoop (Front Tier, bank default model, 50% loss reduction). Lead with it; say you're actively refreshing. Never bluff.
Якорь: реальный production Scala + Spark на Hadoop (Front Tier, модель дефолтов банков, снижение потерь на 50%). Лидируй с этого; говори, что активно освежаешь. Никогда не блефуй.
⭐ 5 top topics drilled ⭐ 5 ключевых тем Flip-cards for self-quiz ↓ Карточки-переворотки ↓ Click section headers to collapse Клик на заголовок сворачивает секцию
01

Syntax essentials for a Python devОсновы синтаксиса для Python-разработчика

Static types, expression-oriented (almost everything returns a value), semicolons optional, braces or significant indentation (Scala 3). The last expression is the return value — no return needed.

Статическая типизация, ориентация на выражения (почти всё возвращает значение), точки с запятой опциональны, фигурные скобки или значимые отступы (Scala 3). Последнее выражение — это возвращаемое значение, return не нужен.

object Demo {
  val name: String = "amal"      // immutable, type inferred if omitted
  var count: Int = 0             // reassignable (avoid)

  def square(x: Int): Int = x * x   // = body, last expr returned

  def grade(s: Int): String =          // if is an expression
    if (s >= 90) "A" else "B"

  val doubled = List(1,2,3).map(_ * 2)  // _ = positional arg
  val msg = s"hello $name, n=$count"     // string interpolation
}

Mental model vs Python

Ментальная модель vs Python

  • val ≈ a name you can't rebind (not const of contents)
  • Type annotations come after the name: x: Int
  • Lambdas: x => x + 1 or shorthand _ + 1
  • UnitNone/void (returned by side-effecting code)
  • val ≈ имя, которое нельзя переприсвоить (не const содержимого)
  • Аннотация типа идёт после имени: x: Int
  • Лямбды: x => x + 1 или краткая форма _ + 1
  • UnitNone/void (возвращается побочно-эффектным кодом)

Gotchas coming from Python

Подводные камни для Python-разработчиков

  • No truthiness — conditions must be Boolean
  • Indentation is cosmetic in Scala 2 (use braces); Scala 3 allows significant indentation
  • Everything is an object; no primitives at the language level (boxing handled by JVM)
  • == is value equality (calls .equals), unlike Java
  • Нет truthiness — условия должны быть Boolean
  • Отступы косметические в Scala 2 (нужны фигурные скобки); Scala 3 допускает значимые отступы
  • Всё — объект; на уровне языка нет примитивов (boxing обрабатывает JVM)
  • == — это равенство по значению (вызывает .equals), в отличие от Java
02

Immutability & why distributed code loves itНеизменяемость и почему распределённый код её любит probedдожимают

val vs var

val = immutable binding (can't reassign the reference). var = reassignable. Default to val; use var rarely + locally.

val = неизменяемая привязка (нельзя переприсвоить ссылку). var = переприсваиваемая. По умолчанию val; var редко и локально.

immutable vs mutable collections

immutable vs mutable коллекции

scala.collection.immutable.* return a new collection on every change; the original is untouched. mutable.* change in place.

scala.collection.immutable.* возвращают новую коллекцию при каждом изменении; оригинал не затронут. mutable.* меняют на месте.

Why it matters in Spark

Почему это важно в Spark

  • Closures are serialized + shipped to executors; pure, side-effect-free transforms are safe to retry and safe to parallelize.
  • Referential transparency = Spark can recompute a lost partition from lineage (fault recovery).
  • Mutating shared state inside map/foreach is a bug — each executor mutates its own copy → use accumulators.
  • Замыкания сериализуются и отправляются на executors; чистые трансформации без побочных эффектов безопасно повторять и безопасно параллелить.
  • Ссылочная прозрачность = Spark может пересчитать потерянную партицию из lineage (fault recovery).
  • Мутация общего состояния внутри map/foreach — это баг; каждый executor мутирует свою копию → используй аккумуляторы.
⚠ The #1 trap val does not make the object immutable — only the reference. val xs = mutable.ArrayBuffer(1,2) — you can't reassign xs, but you can xs += 3. Call this out before they ask.
⚠ Ловушка №1 val не делает объект неизменяемым — только ссылку. val xs = mutable.ArrayBuffer(1,2) — нельзя переприсвоить xs, но можно xs += 3. Назови это до того, как спросят.
Amal anchor "In my Front Tier Spark pipeline, transformations were pure functions over RDDs/DataFrames — that's what kept retries and partition recomputation correct at scale."
Якорь Amal «В моём Spark-пайплайне во Front Tier трансформации были чистыми функциями над RDD/DataFrame — это и обеспечивало корректность ретраев и пересчёта партиций на масштабе.»
03

Case classes & pattern matchingCase-классы и pattern matching workhorseрабочая лошадка

A case class auto-generates

Case-класс автоматически генерирует

  • immutable val fields
  • equals/hashCode = value equality
  • toString, copy
  • companion apply (no new) + unapply (enables matching/destructuring)
  • неизменяемые поля val
  • equals/hashCode = равенство по значению
  • toString, copy
  • companion-объект с apply (без new) + unapply (позволяет pattern matching/деструктуризацию)

Value equality is why case classes are great Spark keys — same fields → equal → same hash for reduceByKey, joins, distinct.

Равенство по значению — вот почему case-классы отличные ключи для Spark: одинаковые поля → равны → один и тот же хеш для reduceByKey, joins, distinct.

Sealed ADT = exhaustiveness

Sealed ADT = полнота проверки

A sealed trait + case classes restricts subtypes to one file, so the compiler warns on non-exhaustive matches. Model your event types this way and the compiler tells you when you forgot a case.

sealed trait + case-классы ограничивают подтипы одним файлом, поэтому компилятор предупреждает о неполных match. Моделируй типы событий так, и компилятор скажет, когда забыл case.

sealed trait Event
case class Click(userId: Long, ts: Long) extends Event
case class Purchase(userId: Long, amount: Double, ts: Long) extends Event
case object Heartbeat extends Event

def revenue(e: Event): Double = e match {
  case Purchase(_, amt, _) if amt > 0 => amt   // destructure + guard
  case _: Click                          => 0.0   // type match
  case Heartbeat                         => 0.0   // singleton match
}
Interviewer will probe "What does sealed give you?" → compiler-checked exhaustiveness. "What enables matching on a case class?" → the generated unapply extractor.
О чём спросит интервьюер «Что даёт sealed?» → проверку полноты на уровне компилятора. «Что позволяет делать matching на case-классе?» → сгенерированный unapply extractor.
04

Option / Either / Try — modelling absence & failureмоделирование отсутствия и сбоев

TypeCasesUse when
Option[A]Some(a) / Nonevalue may be absent; no reason needed
Either[L,R]Left(err) / Right(ok)can fail and you need the reason (Right = success)
Try[A]Success(a) / Failure(t)wrapping code that throws: Try(parse(s))
ТипСлучаиКогда использовать
Option[A]Some(a) / Noneзначение может отсутствовать; причина не нужна
Either[L,R]Left(err) / Right(ok)может упасть, и нужна причина (Right = успех)
Try[A]Success(a) / Failure(t)оборачивание кода, который бросает исключения: Try(parse(s))
def parseAge(s: String): Option[Int] = s.toIntOption.filter(_ >= 0)

val age = parseAge(raw).getOrElse(-1)        // explicit default, no NPE

// chaining short-circuits on the first None / Left / Failure
val total = for {
  a <- parseAge(x)
  b <- parseAge(y)
} yield a + b                              // Option[Int]
⚠ null still exists Scala interops with Java, so null is real. Idiom: wrap Java returns with Option(javaCall())null becomes None safely.
⚠ null всё ещё существует Scala взаимодействует с Java, поэтому null реален. Идиома: оборачивай Java-возвраты в Option(javaCall())null безопасно становится None.
Probe "Option vs Either?" → Option = present/absent; Either/Try = carries the failure reason. Spark angle: nullability lives in the DataFrame schema, not as Option per row, but a UDF returning Option[T] produces a nullable column.
Проба «Option vs Either?» → Option = есть/нет; Either/Try = несёт причину сбоя. Угол Spark: nullability живёт в схеме DataFrame, а не как Option на строку, но UDF, возвращающий Option[T], даёт nullable-колонку.
05

Collections + functional ops (map / flatMap / fold)Коллекции + функциональные операции (map / flatMap / fold)

The big three

Большая тройка

  • map(f: A=>B) — 1:1, keeps cardinality
  • flatMap(f: A=>F[B]) — 1:many then flatten (can expand or drop)
  • foldLeft(z)((acc,x)=>…) — aggregate to one value with explicit zero; reduce has no zero (fails on empty)
  • map(f: A=>B) — 1:1, сохраняет количество
  • flatMap(f: A=>F[B]) — 1:много, затем flatten (может расширить или отбросить)
  • foldLeft(z)((acc,x)=>…) — свёртка в одно значение с явным zero; reduce без zero (падает на пустой)

Sequence types

Типы последовательностей

  • List — immutable linked; O(1) prepend
  • Vector — immutable, ~O(1) indexed (default)
  • Array — mutable, JVM-native, Java interop
  • Map/Set — immutable by default
  • List — неизменяемый связный; O(1) prepend
  • Vector — неизменяемый, ~O(1) по индексу (по умолчанию)
  • Array — мутабельный, JVM-нативный, Java interop
  • Map/Set — immutable по умолчанию
List(1,2).flatMap(x => List(x, x*10))   // List(1,10,2,20)
List(1,2,3,4).foldLeft(0)(_ + _)      // 10
xs.collect { case Purchase(_, amt, _) => amt }  // map+filter via partial fn
events.groupBy(_.userId)                  // Map[Long, List[Event]]
xs.zipWithIndex                           // like Python enumerate
xs.sliding(3)                            // windowing over a Seq
Why flatMap matters twice It tokenizes in Spark word-count (flatMap(_.split(" "))) and it's the engine behind for-comprehensions / monadic chaining of Option/Either/Future.
Почему flatMap важен дважды Он токенизирует в Spark word-count (flatMap(_.split(" "))) и это движок for-comprehensions / монадической цепочки Option/Either/Future.
06

Implicits basicsосновы (Scala 3: given / using)

Three things you'll actually meet:

Три вещи, с которыми реально встретишься:

  • Implicit parameters — compiler supplies a value: a SparkSession, an ExecutionContext, an Encoder[T].
  • Type classes — ad-hoc polymorphism: Ordering[T], Encoder[T].
  • Extension methods — add methods to existing types (old "implicit class" pattern).
  • Implicit parameters — компилятор подставляет значение: SparkSession, ExecutionContext, Encoder[T].
  • Type classes — ad-hoc полиморфизм: Ordering[T], Encoder[T].
  • Extension methods — добавление методов к существующим типам (старый паттерн «implicit class»).
The one line to remember import spark.implicits._ supplies Encoders for case classes/primitives and the $"col" / .toDF / .toDS syntax. Without it, ds.map(...) won't compile (no Encoder in scope).
Одна строка для запоминания import spark.implicits._ даёт Encoders для case-классов/примитивов и синтаксис $"col" / .toDF / .toDS. Без этого ds.map(...) не скомпилируется (нет Encoder в scope).
⚠ Senior take Implicits are powerful for libraries/type classes but hurt readability ("where did this come from?"). Use sparingly in app code. Scala 3 makes them explicit via given/using.
⚠ Уровень Senior Implicits мощны для библиотек/type classes, но вредят читаемости («откуда это взялось?»). Используй скупо в прикладном коде. Scala 3 делает их явными через given/using.

Encoders (the why behind typed Datasets)

Encoders (почему за типизированными Datasets)

An Encoder[T] converts a JVM object ↔ Spark's Tungsten binary/columnar format without slow Java/Kryo serialization. Auto-derived for case classes via implicits._. This is the bridge that makes Dataset[T] both typed and fast — and why typed .map lambdas cost ser/de.

Encoder[T] конвертирует JVM-объект ↔ бинарный/колоночный формат Tungsten без медленной Java/Kryo-сериализации. Авто-выводится для case-классов через implicits._. Это мост, который делает Dataset[T] и типизированным, и быстрым — и почему типизированные .map-лямбды стоят ser/de.

07

Spark in Scalaна Scala: RDD / DataFrame / Dataset heavily weightedвысокий вес

APITyped?Catalyst/Tungsten?When
RDD[T]✅ compile-time❌ noneunstructured data, fine-grained control. Slow default.
DataFrame = Dataset[Row]❌ runtime (Row)✅ fullrelational heavy lifting. PySpark's only option.
Dataset[T] (JVM only)✅ compile-time✅ (limited thru lambdas)typed domain model on case classes.
APIТипизирован?Catalyst/Tungsten?Когда
RDD[T]✅ compile-time❌ нетнеструктурированные данные, детальный контроль. Медленный по умолчанию.
DataFrame = Dataset[Row]❌ runtime (Row)✅ полностьюреляционная тяжёлая работа. Единственный вариант PySpark.
Dataset[T] (только JVM)✅ compile-time✅ (ограничено через lambdas)типизированная доменная модель на case-классах.
"Why Scala over PySpark?" (concrete) The typed Dataset API only exists on the JVM. It catches schema/typo errors at compile time and runs UDFs in-JVM with no Python↔JVM ser/de penalty.
«Почему Scala, а не PySpark?» (конкретно) Типизированный Dataset API существует только на JVM. Он ловит ошибки схемы/опечатки на этапе компиляции и запускает UDF внутри JVM без Python↔JVM ser/de-пенальти.
⚠ The senior trade-off Typed Dataset.map/.filter(fn) lambdas are opaque to Catalyst — no predicate pushdown / column pruning through a closure, plus object ser/de cost. So "Dataset > DataFrame" is wrong. Mix: DataFrame column ops for the heavy relational work, drop to typed Dataset only where compile-time safety earns its keep.
⚠ Компромисс уровня Senior Типизированные Dataset.map/.filter(fn) лямбды непрозрачны для Catalyst — нет predicate pushdown / column pruning через замыкание, плюс цена object ser/de. Поэтому «Dataset > DataFrame» неверно. Микс: DataFrame column ops для тяжёлой реляционной работы, спускайся к типизированному Dataset только там, где compile-time safety окупается.

How a Spark job runs (ASCII)

Как выполняется Spark-джоба (ASCII)

Driver (your code, builds the DAG, lazy) │ action triggers a Job ▼ Job ──► split at shuffles into Stages ──► each Stage = many Tasks (1 per partition) │ ┌─────────────────────────────────────┴───────────────┐ ▼ ▼ ▼ Executor Executor Executor (tasks on its partitions, cache, shuffle write/read) transformation (map/filter/join) = lazy, builds lineage action (count/collect/write) = triggers execution shuffle (groupBy/join/repartition) = stage boundary, network move
Driver (твой код, строит DAG, ленивый) │ action запускает Job ▼ Job ──► разделяется на Stages по shuffles ──► каждый Stage = много Tasks (1 на партицию) │ ┌─────────────────────────────────────┴───────────────┐ ▼ ▼ ▼ Executor Executor Executor (задачи на своих партициях, cache, shuffle write/read) transformation (map/filter/join) = ленивый, строит lineage action (count/collect/write) = запускает выполнение shuffle (groupBy/join/repartition) = граница stage, сетевая пересылка

Word count + top-K (be able to type these)

Word count + top-K (умей набрать это)

val spark = SparkSession.builder.appName("wc").getOrCreate()
import spark.implicits._

// RDD: reduceByKey does a MAP-SIDE COMBINE before the shuffle -> cheaper than groupByKey
val counts = spark.sparkContext.textFile("path")
  .flatMap(_.split("\\s+")).filter(_.nonEmpty)
  .map(w => (w, 1)).reduceByKey(_ + _)

// DataFrame: top-3 per region via window fn (mirrors your SQL skill)
import org.apache.spark.sql.expressions.Window
val w = Window.partitionBy($"region").orderBy($"streams".desc)
df.withColumn("rn", row_number().over(w)).filter($"rn" <= 3)
⚠ Task not serializable A lambda that references an instance field serializes the whole enclosing object; if it's not Serializable you get this error. Fix: copy the value into a local val before the closure, or broadcast it.
⚠ Task not serializable Лямбда, которая ссылается на поле инстанса, сериализует весь объемлющий объект; если он не Serializable, получишь эту ошибку. Лечение: скопируй значение в локальный val перед замыканием или broadcast.
Probes you'll get reduceByKey vs groupByKey (map-side combine); transformation vs action; narrow vs wide; skew fix (salting / broadcast / AQE); cache vs checkpoint; repartition vs coalesce.
Пробы, которые будут reduceByKey vs groupByKey (map-side combine); transformation vs action; narrow vs wide; лечение skew (salting / broadcast / AQE); cache vs checkpoint; repartition vs coalesce.
Amal anchor "At Front Tier I was RDD / early-DataFrame era and had to manage shuffles + partitioning to make large joins tractable. Today I'd default to DataFrame for relational work, typed Datasets where compile-time safety on the domain model pays off."
Якорь Amal «Во Front Tier я был в эпоху RDD / ранних DataFrame и приходилось управлять shuffles + partitioning, чтобы сделать большие join управляемыми. Сегодня я бы по умолчанию брал DataFrame для реляционной работы, типизированные Datasets там, где compile-time safety на доменной модели окупается.»
08

Common idiomsОбщие идиомы

  • for-comprehension = sugar over flatMap/map/withFilter; works on any monad; short-circuits.
  • companion object + apply = factory, no new.
  • singleton object = Scala's "static" utility holder.
  • lazy val = compute once on first access (memoized).
  • partial functions in collect = map+filter in one.
  • for-comprehension = синтаксический сахар над flatMap/map/withFilter; работает на любой монаде; короткое замыкание.
  • companion object + apply = фабрика, без new.
  • singleton object = «статический» держатель утилит в Scala.
  • lazy val = вычисляется один раз при первом доступе (мемоизация).
  • partial functions в collect = map+filter в одном.
  • tuple destructuring: rdd.map { case (k,v) => (k, v*2) }.
  • string interpolation: s"$x", f"$x%.2f".
  • traits as mix-ins for multiple inheritance (linearization).
  • currying: def add(a:Int)(b:Int) → partial app add(5) _.
  • pure logic, separate I/O → unit-test transforms without Spark.
  • деструктуризация кортежей: rdd.map { case (k,v) => (k, v*2) }.
  • интерполяция строк: s"$x", f"$x%.2f".
  • traits как mix-ins для множественного наследования (линеаризация).
  • каррирование: def add(a:Int)(b:Int) → partial app add(5) _.
  • чистая логика, отдельный I/O → unit-тест трансформаций без Spark.
Testing (you'll be probed) ScalaTest / specs2; ScalaCheck for property-based tests of pure transforms; local SparkSession("local[*]") with in-memory DataFrames for Spark tests. Keep business logic pure → test without Spark; integration-test the I/O.
Тестирование (тебя спросят) ScalaTest / specs2; ScalaCheck для property-based тестов чистых трансформаций; локальная SparkSession("local[*]") с in-memory DataFrames для Spark-тестов. Держи бизнес-логику чистой → тестируй без Spark; интеграционно тестируй I/O.
09

Java interop

Scala compiles to JVM bytecode → call any Java library directly. Spark, Hadoop, and the Kafka client are all JVM/Java.

Scala компилируется в JVM-байткод → вызывай любую Java-библиотеку напрямую. Spark, Hadoop и Kafka-клиент — всё JVM/Java.

import scala.jdk.CollectionConverters._   // the bridge import

val scalaList = javaList.asScala               // Java -> Scala
val javaList2 = scalaSeq.asJava                // Scala -> Java

val safe = Option(javaApi.getThing())          // null -> None, no NPE
  • Collections: .asScala / .asJava (e.g. iterate Kafka ConsumerRecords).
  • null: wrap with Option(...).
  • Checked exceptions: Scala has none — Java checked exceptions propagate as unchecked.
  • JavaBeans: @BeanProperty generates getters/setters when Java callers expect them.
  • Mappings: Unitvoid; primitives/String/arrays map directly; object members ≈ statics.
  • Коллекции: .asScala / .asJava (например, итерация Kafka ConsumerRecords).
  • null: оборачивай в Option(...).
  • Checked exceptions: в Scala их нет — Java checked exceptions распространяются как unchecked.
  • JavaBeans: @BeanProperty генерирует getters/setters, когда Java-вызывающий их ожидает.
  • Маппинги: Unitvoid; примитивы/String/массивы маппятся напрямую; члены object ≈ statics.
Amal anchor "Consuming Kafka from Scala: use the Java Kafka client directly, .asScala the records, wrap nullable returns in Option." Ties to your IU Group Kafka ingestion work.
Якорь Amal «Потребление Kafka из Scala: используй Java Kafka-клиент напрямую, .asScala для записей, оборачивай nullable-возвраты в Option.» Связано с твоей работой по Kafka ingestion в IU Group.
10

Python → Scala Rosetta

ConceptPythonScala
Variablex = 5val x = 5 / var x = 5
Functiondef f(a): return a+1def f(a: Int): Int = a + 1
Lambdalambda x: x+1x => x + 1 / _ + 1
List[1,2,3]List(1, 2, 3)
Dict{"a":1}Map("a" -> 1)
Tuple(a, b)(a, b) · ._1 ._2
Comprehension[x*2 for x in xs]xs.map(_ * 2)
Filter + map[x*2 for x in xs if x>0]xs.filter(_ > 0).map(_ * 2)
Flatten[y for x in xss for y in x]xss.flatten / xss.flatMap(identity)
Dataclass@dataclass class Pcase class P(name: String, age: Int)
None / optional getNone · d.get(k, dflt)None: Option[T] · m.getOrElse(k, dflt)
try/excepttry/exceptTry { ... } / try ... catch { case e => }
isinstance / switchisinstance / matchx match { case _: T => ...; case Foo(a) => ... }
String formatf"{x}"s"$x"
Static method@staticmethodmethod on object
enumerateenumerate(xs)xs.zipWithIndex
Aggregatefunctools.reducexs.foldLeft(z)(f) / xs.reduce(f)
PySpark filterdf.filter(df.x > 0)df.filter($"x" > 0)
Typed transform(no equivalent)ds.map(p => p.age): Dataset[Int]
КонцептPythonScala
Переменнаяx = 5val x = 5 / var x = 5
Функцияdef f(a): return a+1def f(a: Int): Int = a + 1
Лямбдаlambda x: x+1x => x + 1 / _ + 1
Список[1,2,3]List(1, 2, 3)
Словарь{"a":1}Map("a" -> 1)
Кортеж(a, b)(a, b) · ._1 ._2
Comprehension[x*2 for x in xs]xs.map(_ * 2)
Filter + map[x*2 for x in xs if x>0]xs.filter(_ > 0).map(_ * 2)
Flatten[y for x in xss for y in x]xss.flatten / xss.flatMap(identity)
Датакласс@dataclass class Pcase class P(name: String, age: Int)
None / optional getNone · d.get(k, dflt)None: Option[T] · m.getOrElse(k, dflt)
try/excepttry/exceptTry { ... } / try ... catch { case e => }
isinstance / switchisinstance / matchx match { case _: T => ...; case Foo(a) => ... }
Формат строкиf"{x}"s"$x"
Static-метод@staticmethodметод на object
enumerateenumerate(xs)xs.zipWithIndex
Aggregatefunctools.reducexs.foldLeft(z)(f) / xs.reduce(f)
PySpark filterdf.filter(df.x > 0)df.filter($"x" > 0)
Типизир. трансформ(нет эквивалента)ds.map(p => p.age): Dataset[Int]
11

Flip-to-reveal self-quizСамопроверка: карточки-переворотки

Tap a card to flip. Answer out loud first.

Нажми карточку, чтобы перевернуть. Сначала ответь вслух.