Qué priorizaría aprender en Data Engineering
2026-02-22
Data Engineer

Yendo por el MVP
Si arrancara de cero
¿qué priorizaría aprender en Data Engineering?
Para un perfil de ingeniero de datos
1
No estoy hablando de print("Hola mundo")
venv, pip, pyproject.tomlpytestEl 80% de lo que hacés en data es manipular estructuras.
from dataclasses import dataclass
from typing import List, Optional
@dataclass
class DataProduct:
name: str
owner: str
tables: List[str]
sla_hours: Optional[int] = None
def is_critical(self) -> bool:
return self.sla_hours is not None and self.sla_hours <= 4
# Leer config, parsear, validar — el 60% del laburo real
with open("products.json") as f:
products = [DataProduct(**p) for p in json.load(f)]
critical = [p for p in products if p.is_critical()]El problema de los objetos mutables:
El problema de los objetos mutables:
@dataclass
class Config:
host: str
port: int
timeout: int = 30
config = Config("db.prod.internal", 5432)
# Función B modifica la config sin que te des cuenta
def add_debug_settings(config):
config.host = "localhost" # ← cambió producción
connect(config) # → "db.prod.internal:5432" ✓
add_debug_settings(config)
connect(config) # → "localhost:5432" ← BUG@dataclass(frozen=True)
class Config:
host: str
port: int
timeout: int = 30
config = Config("db.prod.internal", 5432)
config.host = "localhost" # ← FrozenInstanceError
# Si necesitás una versión modificada, creás una nueva:
from dataclasses import replace
debug_config = replace(config, host="localhost", timeout=999)
print(config.host) # → "db.prod.internal"
print(debug_config.host) # → "localhost"Crear en vez de mutar — la forma más segura de trabajar con datos.
from pydantic import BaseModel, field_validator
class DataContract(BaseModel):
name: str
owner: str
sla_hours: int
columns: list[str]
@field_validator("sla_hours")
@classmethod
def sla_must_be_positive(cls, v):
if v <= 0:
raise ValueError("SLA debe ser mayor a 0")
return v
# Valida automáticamente al crear:
contract = DataContract(
name="transactions", owner="backend-team",
sla_hours=4, columns=["id", "amount", "date"]
)| Situación | Tipo |
|---|---|
| Modelo de datos simple (mayoría) | @dataclass |
| Config que no debe cambiar | @dataclass(frozen=True) |
| Resultado liviano, desempaquetable | NamedTuple |
| Parsing JSON/YAML con validación | Pydantic BaseModel |
| Lógica de negocio compleja | Clase tradicional |
| Dato descartable, se usa una vez | dict |
Arrancá con
@dataclass. Si necesitás inmutabilidad →frozen=True. Si necesitás validación → Pydantic.
2
+50 años y sigue siendo el más usado en datos
ROW_NUMBER(), LAG(), LEAD(), RANK()CROSS, ANTI, SEMI — no solo INNER y LEFTHAVING, GROUPING SETSCuando dominás window functions, resolvés en una query lo que antes te llevaba 3 scripts de Python.
WITH daily_sales AS (
SELECT
customer_id,
sale_date AS sale_day,
SUM(amount) AS daily_total,
COUNT(*) AS txn_count
FROM sales
WHERE status = 'COMPLETED'
GROUP BY 1, 2
),
ranked AS (
SELECT *,
ROW_NUMBER() OVER (
PARTITION BY customer_id
ORDER BY daily_total DESC
) AS rn,
LAG(daily_total) OVER (
PARTITION BY customer_id
ORDER BY sale_day
) AS prev_day_total
FROM daily_sales
)
SELECT customer_id, sale_day, daily_total,
prev_day_total,
ROUND((daily_total - prev_day_total)
/ prev_day_total * 100, 2) AS pct_change
FROM ranked WHERE rn <= 5
ORDER BY customer_id, rn3
No es “pandas grande” — es un paradigma distinto
collect(), count(), write()coalesce(1) en 500GB mata el clusterJOIN, GROUP BY puede generar unofrom pyspark.sql import functions as F
from pyspark.sql.window import Window
# Bien: funciones nativas de Spark
df = (
spark.read.table("sales.transactions")
.filter(F.col("status") == "completed")
.withColumn(
"running_total",
F.sum("amount").over(
Window.partitionBy("customer_id")
.orderBy("transaction_date")
.rowsBetween(Window.unboundedPreceding,
Window.currentRow)
)
)
)Serde overhead: serialización JVM ↔︎ Python fila por fila = lento.
# Wrapper: usa API nativa — Spark la entiende y optimiza
def add_revenue_flag(df, threshold=1000):
return df.withColumn(
"high_revenue",
F.when(F.col("amount") > threshold, True)
.otherwise(False)
)
# UDF: Python puro — fila por fila, sin optimización
@udf(returnType=BooleanType())
def high_revenue_udf(amount):
return amount > 1000Si no te queda otra, usá Pandas UDFs — procesan por batch, no fila por fila.
4
Modularidad, testing, documentación, linaje
SELECT, se referencian con { ref() }not_null, unique, relationships, tests custom-- models/staging/stg_transactions.sql
WITH source AS (
SELECT * FROM {{ source('raw', 'transactions') }}
),
cleaned AS (
SELECT
transaction_id,
customer_id,
CAST(amount AS DECIMAL(18,2)) AS amount,
CAST(transaction_date AS DATE) AS transaction_date,
UPPER(TRIM(status)) AS status
FROM source
WHERE transaction_id IS NOT NULL
)
SELECT * FROM cleaned# models/staging/stg_transactions.yml
version: 2
models:
- name: stg_transactions
description: "Transacciones limpias del core bancario"
columns:
- name: transaction_id
tests: [not_null, unique]
- name: amount
tests:
- not_null
- dbt_utils.accepted_range:
min_value: 0
max_value: 1000000
- name: status
tests:
- accepted_values:
values: ['COMPLETED', 'PENDING', 'FAILED']| Caso | Herramienta |
|---|---|
| Transformaciones SQL (limpieza, joins, agg) | dbt |
| Lógica con APIs externas | PySpark |
| ML pipelines | PySpark + MLflow |
| Archivos no estructurados (JSON, XML) | PySpark |
| Modelos dimensionales (Kimball) | dbt |
| Streaming / near-real-time | Structured Streaming |
Si tu transformación es SQL puro, usá dbt. Si necesitás lógica compleja, PySpark.
5
No es solo “Spark en la nube”
catalog.schema.table-- Crear catalog y schema
CREATE CATALOG IF NOT EXISTS analytics;
CREATE SCHEMA IF NOT EXISTS analytics.sales;
-- Tabla con Liquid Clustering
CREATE TABLE analytics.sales.transactions (
transaction_id BIGINT,
customer_id BIGINT,
amount DECIMAL(18,2),
transaction_date DATE
)
USING DELTA
CLUSTER BY (transaction_date, customer_id);
-- Gobernanza básica
GRANT USE CATALOG ON CATALOG analytics TO `analysts`;
GRANT USE SCHEMA ON SCHEMA analytics.sales TO `analysts`;
GRANT SELECT ON TABLE analytics.sales.transactions TO `analysts`;BONUS
El error más común — y el que más cuesta corregir
La regla: si algo corre más de una vez, no debería estar en un notebook.
# src/transformations/clean_sales.py
def clean_transactions(df):
return (
df
.filter(F.col("transaction_id").isNotNull())
.withColumn("amount",
F.col("amount").cast("decimal(18,2)"))
.withColumn("status",
F.upper(F.trim(F.col("status"))))
.filter(F.col("amount") > 0)
)
# tests/test_clean_sales.py
def test_clean_removes_nulls(spark):
input_df = spark.createDataFrame([
(1, 100.0, "completed"),
(None, 200.0, "pending"),
], ["transaction_id", "amount", "status"])
result = clean_transactions(input_df)
assert result.count() == 1Los conceptos se transfieren. No necesitás aprenderlos todos al mismo tiempo.
Construí algo.
Agarrá un dataset público, armá un pipeline de punta a punta, deployalo, rompelo, arreglalo.
Eso vale más que 10 cursos.

Spark de Ideas — Mauro Loprete