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 (notconstof contents)- Type annotations come after the name:
x: Int - Lambdas:
x => x + 1or shorthand_ + 1 Unit≈None/void(returned by side-effecting code)
val≈ имя, которое нельзя переприсвоить (неconstсодержимого)- Аннотация типа идёт после имени:
x: Int - Лямбды:
x => x + 1или краткая форма_ + 1 Unit≈None/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
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/foreachis a bug — each executor mutates its own copy → use accumulators.
- Замыкания сериализуются и отправляются на executors; чистые трансформации без побочных эффектов безопасно повторять и безопасно параллелить.
- Ссылочная прозрачность = Spark может пересчитать потерянную партицию из lineage (fault recovery).
- Мутация общего состояния внутри
map/foreach— это баг; каждый executor мутирует свою копию → используй аккумуляторы.
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.val не делает объект неизменяемым — только ссылку. val xs = mutable.ArrayBuffer(1,2) — нельзя переприсвоить xs, но можно xs += 3. Назови это до того, как спросят.Case classes & pattern matchingCase-классы и pattern matching ⭐ workhorseрабочая лошадка
▼A case class auto-generates
Case-класс автоматически генерирует
- immutable
valfields equals/hashCode= value equalitytoString,copy- companion
apply(nonew) +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
}
sealed give you?" → compiler-checked exhaustiveness. "What enables matching on a case class?" → the generated unapply extractor.sealed?» → проверку полноты на уровне компилятора. «Что позволяет делать matching на case-классе?» → сгенерированный unapply extractor.Option / Either / Try — modelling absence & failureмоделирование отсутствия и сбоев ⭐
▼| Type | Cases | Use when |
|---|---|---|
Option[A] | Some(a) / None | value 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 is real. Idiom: wrap Java returns with Option(javaCall()) → null becomes None safely.null реален. Идиома: оборачивай Java-возвраты в Option(javaCall()) → null безопасно становится None.Option per row, but a UDF returning Option[T] produces a nullable column.Option на строку, но UDF, возвращающий Option[T], даёт nullable-колонку.Collections + functional ops (map / flatMap / fold)Коллекции + функциональные операции (map / flatMap / fold) ⭐
▼The big three
Большая тройка
map(f: A=>B)— 1:1, keeps cardinalityflatMap(f: A=>F[B])— 1:many then flatten (can expand or drop)foldLeft(z)((acc,x)=>…)— aggregate to one value with explicit zero;reducehas 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) prependVector— immutable, ~O(1) indexed (default)Array— mutable, JVM-native, Java interopMap/Set— immutable by default
List— неизменяемый связный; O(1) prependVector— неизменяемый, ~O(1) по индексу (по умолчанию)Array— мутабельный, JVM-нативный, Java interopMap/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
flatMap(_.split(" "))) and it's the engine behind for-comprehensions / monadic chaining of Option/Either/Future.flatMap(_.split(" "))) и это движок for-comprehensions / монадической цепочки Option/Either/Future.Implicits basicsосновы (Scala 3: given / using)
▼Three things you'll actually meet:
Три вещи, с которыми реально встретишься:
- Implicit parameters — compiler supplies a value: a
SparkSession, anExecutionContext, anEncoder[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»).
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).given/using.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.
Spark in Scalaна Scala: RDD / DataFrame / Dataset ⭐ heavily weightedвысокий вес
▼| API | Typed? | Catalyst/Tungsten? | When |
|---|---|---|---|
RDD[T] | ✅ compile-time | ❌ none | unstructured data, fine-grained control. Slow default. |
DataFrame = Dataset[Row] | ❌ runtime (Row) | ✅ full | relational 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-классах. |
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.Dataset API существует только на JVM. Он ловит ошибки схемы/опечатки на этапе компиляции и запускает UDF внутри JVM без Python↔JVM ser/de-пенальти.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.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)
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)
Serializable you get this error. Fix: copy the value into a local val before the closure, or broadcast it.Serializable, получишь эту ошибку. Лечение: скопируй значение в локальный val перед замыканием или broadcast.Common idiomsОбщие идиомы
▼for-comprehension = sugar overflatMap/map/withFilter; works on any monad; short-circuits.- companion object +
apply= factory, nonew. - 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 appadd(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 appadd(5) _. - чистая логика, отдельный I/O → unit-тест трансформаций без Spark.
SparkSession("local[*]") with in-memory DataFrames for Spark tests. Keep business logic pure → test without Spark; integration-test the I/O.SparkSession("local[*]") с in-memory DataFrames для Spark-тестов. Держи бизнес-логику чистой → тестируй без Spark; интеграционно тестируй I/O.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 KafkaConsumerRecords). - null: wrap with
Option(...). - Checked exceptions: Scala has none — Java checked exceptions propagate as unchecked.
- JavaBeans:
@BeanPropertygenerates getters/setters when Java callers expect them. - Mappings:
Unit≈void; primitives/String/arrays map directly;objectmembers ≈ statics.
- Коллекции:
.asScala/.asJava(например, итерация KafkaConsumerRecords). - null: оборачивай в
Option(...). - Checked exceptions: в Scala их нет — Java checked exceptions распространяются как unchecked.
- JavaBeans:
@BeanPropertyгенерирует getters/setters, когда Java-вызывающий их ожидает. - Маппинги:
Unit≈void; примитивы/String/массивы маппятся напрямую; членыobject≈ statics.
.asScala the records, wrap nullable returns in Option." Ties to your IU Group Kafka ingestion work..asScala для записей, оборачивай nullable-возвраты в Option.» Связано с твоей работой по Kafka ingestion в IU Group.Python → Scala Rosetta
▼| Concept | Python | Scala |
|---|---|---|
| Variable | x = 5 | val x = 5 / var x = 5 |
| Function | def f(a): return a+1 | def f(a: Int): Int = a + 1 |
| Lambda | lambda x: x+1 | x => 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 P | case class P(name: String, age: Int) |
| None / optional get | None · d.get(k, dflt) | None: Option[T] · m.getOrElse(k, dflt) |
| try/except | try/except | Try { ... } / try ... catch { case e => } |
| isinstance / switch | isinstance / match | x match { case _: T => ...; case Foo(a) => ... } |
| String format | f"{x}" | s"$x" |
| Static method | @staticmethod | method on object |
| enumerate | enumerate(xs) | xs.zipWithIndex |
| Aggregate | functools.reduce | xs.foldLeft(z)(f) / xs.reduce(f) |
| PySpark filter | df.filter(df.x > 0) | df.filter($"x" > 0) |
| Typed transform | (no equivalent) | ds.map(p => p.age): Dataset[Int] |
| Концепт | Python | Scala |
|---|---|---|
| Переменная | x = 5 | val x = 5 / var x = 5 |
| Функция | def f(a): return a+1 | def f(a: Int): Int = a + 1 |
| Лямбда | lambda x: x+1 | x => 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 P | case class P(name: String, age: Int) |
| None / optional get | None · d.get(k, dflt) | None: Option[T] · m.getOrElse(k, dflt) |
| try/except | try/except | Try { ... } / try ... catch { case e => } |
| isinstance / switch | isinstance / match | x match { case _: T => ...; case Foo(a) => ... } |
| Формат строки | f"{x}" | s"$x" |
| Static-метод | @staticmethod | метод на object |
| enumerate | enumerate(xs) | xs.zipWithIndex |
| Aggregate | functools.reduce | xs.foldLeft(z)(f) / xs.reduce(f) |
| PySpark filter | df.filter(df.x > 0) | df.filter($"x" > 0) |
| Типизир. трансформ | (нет эквивалента) | ds.map(p => p.age): Dataset[Int] |
Flip-to-reveal self-quizСамопроверка: карточки-переворотки
▼Tap a card to flip. Answer out loud first.
Нажми карточку, чтобы перевернуть. Сначала ответь вслух.