Menjelajahi kekuatan Scala di dunia data: Apache Spark dengan DataFrame dan RDD untuk pemrosesan dataset besar, serta Kafka dengan producer dan consumer — lewat kafka-clients dan fs2-kafka — untuk streaming event yang andal.

Di episode 13 kalian membangun API; sekarang kita memasuki ranah di mana Scala benar-benar tak tertandingi: big data. Dua infrastruktur paling penting dalam data engineering modern — Apache Spark dan Apache Kafka — keduanya ditulis di Scala, dan penguasaan Scala memberi kalian akses langsung ke API mereka.
Mengapa episode ini penting? Karena data pipeline adalah mesin di balik hampir semua produk modern: rekomendasi, analitik, log processing, dan AI. Data Engineer yang menguasai Scala bisa bekerja langsung dengan Spark dan Kafka tanpa lapisan bahasa perantara — efisiensi dan keleluasaan yang tidak dimiliki pengguna Python murni.
Spark adalah kerangka kerja untuk pemrosesan data terdistribusi: satu dataset besar dipecah dan diproses paralel di banyak mesin. Dua abstraksi utamanya: DataFrame dan RDD.
DataFrame adalah tabel terdistribusi dengan skema — kolom bertipe. API-nya terasa seperti SQL/collection:
import org.apache.spark.sql.SparkSession
val spark = SparkSession.builder
.appName("belajar-scala")
.master("local[*]")
.getOrCreate()
import spark.implicits.*
val df = spark.read
.option("header", true)
.csv("data/orders.csv")
df.printSchema()
df.filter($"amount" > 100_000).groupBy($"city").count().show()Perhatikan kemiripan dengan pelajaran collections (episode 6): filter, groupBy, count — tetapi kini berjalan di atas ribuan core. $"column" adalah sintaks khusus Spark untuk referensi kolom.
RDD (Resilient Distributed Dataset) adalah abstraksi yang lebih rendah — koleksi terdistribusi dengan fungsi transformasi seperti collections Scala:
val rdd = spark.sparkContext.parallelize(1 to 100_000)
val total = rdd
.filter(_ % 2 == 0)
.map(_ * 2)
.sum()
println(total)RDD memberi kontrol granular (partisi, akumulator, broadcast), sementara DataFrame dioptimalkan otomatis oleh Catalyst (query optimizer) dan Tungsten (execution engine). Aturan praktis: mulai dari DataFrame, turun ke RDD hanya jika butuh kontrol tingkat rendah.
Transformasi Spark lazy: df.filter(...) tidak langsung memproses apa pun — ia membangun DAG (Directed Acyclic Graph) eksekusi yang baru dijalankan saat aksi (seperti .count() atau .show()). Ini memungkinkan Spark mengoptimalkan seluruh pipeline sebelum eksekusi. Konsep ini paralel dengan "program sebagai nilai" dari effect system di episode 11.
Kafka adalah sistem message streaming terdistribusi: producer mengirim event ke topic, consumer membaca dari topic. Ia dirancang untuk throughput tinggi dan durabilitas.
Library kafka-clients dari Apache adalah pilihan paling langsung:
libraryDependencies += "org.apache.kafka" % "kafka-clients" % "3.9.0"import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord}
import java.util.Properties
val props = Properties()
props.put("bootstrap.servers", "localhost:9092")
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
val producer = KafkaProducer[String, String](props)
producer.send(ProducerRecord("orders", "order-1", "{\"amount\": 250000}"))
producer.flush()
producer.close()import org.apache.kafka.clients.consumer.{KafkaConsumer, ConsumerConfig}
import java.time.Duration
import java.util
val props = Properties()
props.put("bootstrap.servers", "localhost:9092")
props.put("group.id", "order-processor")
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")
props.put("auto.offset.reset", "earliest")
val consumer = KafkaConsumer[String, String](props)
consumer.subscribe(util.List.of("orders"))
while true do
val records = consumer.poll(Duration.ofMillis(100))
records.forEach { r =>
println(s"terima: ${r.value()}")
}Kelompok consumer (group.id) membuat load-balancing: beberapa instance consumer berbagi partisi topic secara otomatis — itu cara Kafka men-skalakan pemrosesan.
Note
kafka-clients adalah Java API — bekerja sempurna di Scala lewat interop. Namun jika kalian ingin pemrosesan streaming yang composable dengan effect system, fs2-kafka (untuk cats-effect) adalah pilihan yang lebih idiomatis. Kita lihat sekilas di bawah.
fs2-kafka mengintegrasikan Kafka dengan cats-effect dan fs2 streams — event diperlakukan sebagai stream yang bisa dikomposisi, dibatalkan, dan diuji:
import cats.effect.{IO, IOApp}
import fs2.kafka.*
object Main extends IOApp.Simple:
val consumerSettings = ConsumerSettings[IO, String, String]
.withBootstrapServers("localhost:9092")
.withGroupId("order-processor")
val stream =
KafkaConsumer.stream(consumerSettings)
.subscribeTo("orders")
.records
.mapAsync(10) { record =>
IO.println(s"proses: ${record.record.value()}").as(record)
}
val run: IO[Unit] = stream.compile.drainSetiap event menjadi elemen stream; mapAsync(10) memproses 10 event paralel. Pekko Streams (penerus Akka Streams) menawarkan pendekatan serupa dengan model graph yang berbeda. Pilihannya ditentukan ekosistem yang sudah dipakai tim.
| Skenario | Pilihan |
|---|---|
| ETL batch besar | Spark (Scala) — bahasa utama dan API paling penuh |
| Streaming event | Kafka + fs2-kafka / pekko-streams |
| Data pipeline ringan | Collections Scala (episode 6) sudah cukup |
| Analisis ad-hoc | Spark SQL / notebooks (Databricks) |
Kelebihan Scala di data: performa JVM tanpa lapisan tambahan, tipe kuat untuk skema kompleks, dan interop penuh dengan library Java. Kekurangan yang jujur: kurva belajar lebih curam daripada Python.
option("header", true) membuat kolom _c0, _c1; pastikan skema terdefinisi..count(), .show(), atau .write.enable.auto.commit dan commitSync.Inti yang harus dibawa pulang:
kafka-clients untuk kontrol langsung.Di episode 15 selanjutnya, kita menghubungkan aplikasi ke penyimpanan: database dan persistence — JDBC, Doobie, Quill, Slick, connection pooling dengan HikariCP, dan migrasi dengan Flyway. Sampai jumpa di episode 15!