Guía de ingeniería

Cómo empezar a construir un ETL personalizado en Scala

Scala es una gran opción para ETL cuando buscas desarrollo nativo de Spark, seguridad de tipos y una implementación que pueda crecer desde un solo job batch hasta una plataforma de datos más amplia.

Cuándo Scala tiene sentido para ETL

Scala resulta especialmente valioso cuando la lógica ETL es lo suficientemente compleja como para que las abstracciones sólidas, las bibliotecas reutilizables y las verificaciones en compilación reduzcan el mantenimiento real. Con Apache Spark, además, el encaje es natural.

Usa Scala cuando quieras que el código de la pipeline se comporte como software y no como una colección de scripts.

What Spark, Hadoop, and HDFS Actually Are

These technologies are related, but they are not the same layer of the stack.

Spark

Spark is a distributed compute engine. It parallelizes reads, joins, aggregations, and writes across many workers. For ETL, think of Spark as the execution layer rather than the long-term storage layer.

Hadoop

Hadoop is a broader ecosystem that historically bundled distributed storage, cluster resource management, and batch processing. Modern teams often borrow concepts from Hadoop without running a classic Hadoop cluster end to end.

HDFS

HDFS is Hadoop Distributed File System. It stores files across many machines and was the standard storage layer for self-hosted big data clusters before object storage became the default in many cloud environments.

Catalog and table metadata

As soon as multiple jobs and analysts rely on the same curated datasets, you also need a catalog or metastore. It keeps schemas, partitions, and table definitions consistent across your ETL and query tools.

Should You Use HDFS?

Usually not as a default. HDFS still has a place, but many modern ETL stacks use object storage instead.

  • - Choose HDFS mainly when you operate your own cluster, want data locality, and already accept the operational cost of managing distributed storage nodes.
  • - Prefer S3, ADLS, or GCS when your compute can be ephemeral, your storage should scale independently, or you want simpler disaster recovery and cross-service integration.
  • - Do not adopt Hadoop just because you use Spark. Spark can run on Kubernetes, YARN, Databricks, EMR, Synapse, or other managed runtimes without HDFS being the center of the design.

Empieza con la forma correcta

Capas principales

  • - Configuración: rutas, secretos, entornos y tablas destino
  • - Readers: JDBC, archivos, colas o adaptadores API
  • - Transformaciones: normalización, enriquecimiento, deduplicación y reglas de negocio
  • - Validadores: controles de esquema, volumen y dominio
  • - Writers: salidas hacia lake, warehouse o servicios destino

Objetivos tempranos de diseño

  • - Cargas idempotentes para evitar duplicados en re-ejecuciones
  • - Logging estructurado con contexto de batch o partición
  • - Límites claros de fallo entre extract, transform y load
  • - Transformaciones pequeñas y testeables en lugar de jobs monolíticos
  • - Un modelo de despliegue alineado desde el primer día con la runtime

Do You Have to Use Spark?

No. Spark is powerful, but it is not the mandatory starting point for every ETL pipeline.

Use Spark when

You need distributed joins, large historical scans, partitioned batch outputs, or multi-hundred-gigabyte to multi-terabyte processing windows that no longer fit comfortably on one machine.

Skip Spark when

Your ETL is mostly SQL against Postgres, the data volume is still modest, one server can finish the workload inside the SLA, and operational simplicity matters more than horizontal scale.

Common alternatives

Start with Postgres SQL, dbt, Airflow plus Python, DuckDB, or Polars when the pipeline is small or medium. Move to Spark when the workload proves it needs a distributed compute engine.

Esqueleto del proyecto

Un proyecto ETL mínimo basado en Spark puede empezar con build.sbt, una única clase principal y la estructura clásica src/main/scala. Ajusta las versiones de ejemplo a tu clúster objetivo.

name := "custom-etl"

version := "0.1.0"

scalaVersion := "2.13.17"

val sparkVersion = "4.1.2" // Replace to match your target cluster

libraryDependencies ++= Seq(
  "org.apache.spark" %% "spark-sql" % sparkVersion % "provided",
  "com.typesafe" % "config" % "1.4.3",
  "org.scalatest" %% "scalatest" % "3.2.19" % Test
)

Esqueleto mínimo del job

Incluso un job pequeño debería aceptar rutas de entrada y salida por argumentos, crear una SparkSession, aplicar una cadena de transformación estrecha y escribir en un formato determinista como Parquet.

package dev.ryware.etl

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{col, to_timestamp}

object CustomEtlJob {
  def main(args: Array[String]): Unit = {
    val inputPath = args(0)
    val outputPath = args(1)

    val spark = SparkSession.builder()
      .appName("custom-etl-job")
      .getOrCreate()

    val curated = spark.read
      .option("header", "true")
      .csv(inputPath)
      .filter(col("transaction_id").isNotNull)
      .withColumn("amount", col("amount").cast("decimal(18,2)"))
      .withColumn("updated_at", to_timestamp(col("updated_at")))

    curated.write
      .mode("overwrite")
      .parquet(outputPath)

    spark.stop()
  }
}

When Postgres Is Enough and When to Move to a Warehouse

Raw database size matters, but query concurrency, data shape, retention, and the cost of long analytical scans matter more.

Approximate decision guide for analytics workloads. These are planning bands, not hard limits.
Workload stage When Postgres is still reasonable When you should plan the move Typical target
Early analytics Up to roughly 50-100 GB of curated analytical data, a few internal dashboards, light concurrency, and daily batch updates. Stay put unless queries already compete with transactional traffic or refresh windows are missing the SLA. Postgres plus SQL, dbt, and simple orchestration.
Growing reporting stack Roughly 100-300 GB, moderate joins, fewer than 10-20 active analysts, and mostly structured tables. Plan a warehouse when refresh jobs grow brittle, vacuum and index tuning turns into constant work, or semi-structured data starts piling up. Postgres can still work, but start evaluating BigQuery, Snowflake, Redshift, Synapse, or a lakehouse pattern.
Heavy analytics Possible only with careful tuning, but the tradeoff gets worse once scans span hundreds of gigabytes, concurrency climbs, and retention grows fast. Move when BI workloads, historical backfills, and ML feature extraction all want the same data at the same time. Columnar warehouse or lakehouse with separated storage and compute.
Platform scale Rarely the right long-term shape once you are in multi-terabyte territory, ingesting many domains, or supporting self-service analytics. At this point the question is usually not whether to move, but which warehouse, lakehouse, governance model, and cost controls fit the business. Warehouse or lakehouse platform with catalog, partitioning, and workload isolation.

What If the Source Files Are XML and You Need XPath?

Spark can process XML, but the best approach depends on how irregular the documents are and how much XPath logic you need.

Spark can read XML

On open-source Spark, teams commonly add the spark-xml package and define an explicit row tag and schema. That works well when the XML documents are repetitive enough to map into tabular records.

XPath is possible, but expensive

XPath-style extraction is reasonable for a few known fields, but deeply nested or dynamic XPath rules can become CPU-heavy and awkward to test at scale. Flatten the documents as early as you can.

Sometimes pre-processing is better

If the XML is highly irregular, namespace-heavy, or validated against complex XSD rules, it can be cleaner to parse it first with a dedicated XML library in Scala or Java and then land normalized JSON or Parquet for downstream Spark work.

Operational advice

Keep the raw XML, version the mapping rules, capture malformed files separately, and build test fixtures from real samples. Special formats fail at the edges, so sample coverage matters more than happy-path demos.

Lo siguiente que debes añadir para producción

Validación y calidad

  • - Añadir validaciones de esquema antes de escribir salidas curadas
  • - Guardar conteos de filas y tasas de nulos por ejecución
  • - Rechazar o poner en cuarentena particiones inválidas en lugar de corregirlas en silencio

Pruebas y empaquetado

  • - Pruebas unitarias de transformaciones con DataFrames pequeños en memoria
  • - Pruebas de integración para conectores de origen y destino
  • - Empaquetado con sbt y ejecución vía spark-submit o el launcher del clúster

Un primer hito sensato

El primer objetivo no debe ser construir toda la plataforma de datos. Debe ser leer una fuente de forma fiable, aplicar una ruta de transformación limpia, escribir una salida curada y hacer observable la ejecución.

Hito 1

Un solo job batch con salidas deterministas.

Hito 2

Entornos guiados por configuración y reglas de validación.

Hito 3

Orquestación, alertas, lineage y líneas base de calidad.

© 2026 - Ryware.