ETL with Python and geoinformation: production pipelines
A data pipeline is more than a script that moves files. It is a chain of technical decisions that determines whether an organization can trust its indicators, maps and final products. This article explains how to design one that works in production.
What makes an ETL pipeline robust?
Most pipelines I have seen in real projects fail for the same reason: they have no explicit validation rules and do not record what happened at each stage. When something fails, there is no way to know where or why.
A robust pipeline has three properties: verifiability (I can tell whether the data is correct), traceability (I can tell where each record came from) and maintainability (I can modify it without fear of breaking something).
Layered architecture
The most important decision in a pipeline is to never overwrite source data. The raw layer stores an exact copy of what arrived, with an ingestion timestamp. If something fails downstream, you can reprocess it without requesting the file from the source again.
The staging layer is where transformations happen: field normalization, encoding fixes, coordinate reference system conversions, joins and aggregations. Errors are recorded here too — they are documented rather than discarded.
The producción layer contains only records that passed every validation. It is the only layer consumed by analytics systems, maps and APIs.
Explicit validation in geoinformation
For geospatial data, validation means more than checking that a field is not null. A polygon can be syntactically correct but topologically invalid (self-intersection). A coordinate can be within the numeric range but outside the study area. An attribute can exist but fall outside the allowed domain.
Implementation with Python and PostGIS
The following pattern shows a real pipeline structure with explicit validation and structured logging:
import logging
from datetime import datetime
from sqlalchemy import create_engine, text
import geopandas as gpd
from shapely.validation import make_valid
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s"
)
log = logging.getLogger(__name__)
ENGINE = create_engine("postgresql://user:pass@localhost/geodata")
def ingest_raw(filepath: str, layer: str) -> gpd.GeoDataFrame:
"""Carga el archivo original sin modificar y registra la ingesta."""
gdf = gpd.read_file(filepath, layer=layer)
gdf["_ingested_at"] = datetime.utcnow()
gdf["_source_file"] = filepath
gdf.to_postgis("raw_parcels", ENGINE, if_exists="append", index=False)
log.info("raw: %d registros ingestados desde %s", len(gdf), filepath)
return gdf
def validate(gdf: gpd.GeoDataFrame) -> tuple[gpd.GeoDataFrame, list[dict]]:
"""Aplica reglas de validación y separa registros válidos de errores."""
errors = []
valid_idx = []
for idx, row in gdf.iterrows():
record_id = row.get("id", idx)
geom = row.geometry
# Regla 1: geometría no nula
if geom is None or geom.is_empty:
errors.append({"id": record_id, "field": "geometry", "reason": "null_or_empty"})
continue
# Regla 2: geometría válida (auto-intersecciones, etc.)
if not geom.is_valid:
fixed = make_valid(geom)
if not fixed.is_valid:
errors.append({"id": record_id, "field": "geometry", "reason": "invalid_geometry"})
continue
gdf.at[idx, "geometry"] = fixed # corregir si es posible
# Regla 3: CRS correcto
if gdf.crs is None or gdf.crs.to_epsg() != 4326:
errors.append({"id": record_id, "field": "crs", "reason": f"expected_4326_got_{gdf.crs}"})
continue
# Regla 4: atributos obligatorios
for field in ("parcel_id", "land_use", "area_m2"):
if row.get(field) is None:
errors.append({"id": record_id, "field": field, "reason": "required_null"})
break
else:
valid_idx.append(idx)
log.info("validate: %d válidos, %d errores", len(valid_idx), len(errors))
return gdf.loc[valid_idx].copy(), errors
def load_staging(gdf: gpd.GeoDataFrame, errors: list[dict]) -> None:
"""Carga staging con datos válidos y registra errores."""
gdf.to_postgis("staging_parcels", ENGINE, if_exists="replace", index=False)
if errors:
import pandas as pd
err_df = pd.DataFrame(errors)
err_df["logged_at"] = datetime.utcnow()
err_df.to_sql("etl_errors", ENGINE, if_exists="append", index=False)
log.warning("staging: %d errores registrados en etl_errors", len(errors))
# Ejecución
raw = ingest_raw("parcelas_2026.gpkg", layer="parcelas")
valid_gdf, errs = validate(raw)
load_staging(valid_gdf, errs)
Traceability: knowing what happened to each record
The logging above records errors, but traceability goes further: each production record must be traceable to its original source. To achieve this, the fields _ingested_at and _source_file are propagated from raw through to production.
In complex geospatial projects, where data undergoes multiple transformations (CRS changes, dissolves, spatial intersections), it is also useful to store the _etl_step — the name of the function or process that produced the record in its current state.
Rules for a pipeline that lasts
- Never overwrite raw. It is your data insurance. If something fails, you can reprocess it.
- Validation rules are code, not comments. Write them as clearly named functions with unit tests.
- Errors are logged, not ignored. A silent pipeline that incorrectly discards data is more dangerous than one that fails explicitly.
- Each stage has an expected input and a verifiable output. Define the output schema before writing the transformation.
- PostGIS is your ally. Use
ST_IsValid,ST_MakeValid,ST_Withinand CHECK constraints in the database to validate at load time.
A pipeline like this is no slower or more complex than one without validation. But when something fails — and eventually something always does — you have exactly the information you need to understand what happened and how to fix it.