Transformación de datos: estructuremos los datos con Dataform
Antes de iniciar este post debo resaltar que esta es la continuación de Cómo crear un pipeline ELT robusto donde explico la capa de ingesta. En esta entrada hablaremos sobre la etapa de transformación; los conceptos son transversales a diferentes proveedores de la nube. Sin embargo, los detalles de implementación se harán dentro del ecosistema de Google Cloud. Arranquemos.
El Stack
Para lograr implementar un proceso de transformación de datos de grado productivo debemos apoyarnos de las siguientes herramientas:
- Google Cloud Workflows (mencionado en la primera publicación).
- Google Dataform: el protagonista de hoy, pues nos permitirá implementar el paradigma DaC (Data as Code). Es decir, podremos hacer lo siguiente: - Versionar las consultas conforme las vayamos desarrollando. - Definir tests o assertions; si algo falla, nos enteraremos dónde y por qué. - Metadatos y catalogación. Es obligatorio para trabajar con equipos y gobierno de datos.
- BigQuery: nunca falla, es el corazón de todo. Está integrado con Dataform.
El código
Todas las consultas del proyecto, la arquitectura (utilizando Terraform) y los servicios desplegados se encuentran en github.com/JoseChavezUriarte/coc_elt.
Cómo funciona
Por detrás podemos decir que Dataform es el orquestador de consultas. En este caso, no utilizaremos archivos SQL convencionales, sino sqlx. Son bastante similares, pero únicamente deben realizarse cláusulas SELECT. Las consultas DDL se gestionan a través de metadados. Una consulta se ve más o menos así, la cual consta principalmente de dos partes:
Config
Nos dice dónde, cómo y de qué manera se ensamblan los datos. Veamos más de cerca:
config {
type: "incremental",
schema: "coc_silver",
uniqueKey: ["extracted_date", "ptag", "troop_name"],
description: "Denormalized and deduplicated silver table containing player troop statistics.",
tags: ["silver", "daily"],
bigquery: {
partitionBy: "extracted_date",
clusterBy: ["ptag", "troop_village", "troop_name", "troop_level"],
updatePartitionFilter: "extracted_date >= DATE_SUB(CURRENT_DATE(), INTERVAL 2 DAY)",
labels: {
environment: "production",
domain: "clash-of-clans",
layer: "silver"
}
},
assertions: {
uniqueKey: ["extracted_date", "ptag", "troop_name"],
nonNull: ["extracted_date", "ptag", "troop_name"]
},
columns: {
extracted_at: "Timestamp when the player profile raw payload was retrieved from Clash of Clans API.",
extracted_date: "Partitioning date derived from extracted_at.",
ptag: "Unique identifier tag of the player.",
troop_name: "Name of the troop (e.g. Barbarian, Archer).",
troop_level: "Current level of the troop upgraded by the player.",
troop_max_level: "Maximum possible level for the troop.",
troop_village: "Village where the troop is active (home or builderBase)."
}
}
El fragmento previo nos indica:
- La tabla resultante es de tipo incremental, es decir: la transformación se restringe a registros creados o modificados desde la última ejecución. En nuestro caso, se modifica parcialmente porque añadí una ventana deslizante con el atributo
updatedPartitionFilter. En lugar de la última ejecución, consideramos dinámicamente los últimos dos días; esto es útil por si hay reprocesos. - Se define la granularidad con
uniqueKey. - Se definen las particiones de las tablas y clusters en la clave
bigquery. - En
assertionsse valida la granularidad y que no haya nulos. columnsañade descripciones de cada una de las columnas.
Select Statement
Esta parte del archivo es una consulta tradicional, salvo que algo cambia en el primer FROM:
WITH parsed_members AS (
SELECT
extracted_at,
DATE(extracted_at) AS extracted_date,
JSON_VALUE(payload.tag) AS ptag,
payload.troops AS troops
FROM
${ref("coc_members")}
${when(incremental(), "WHERE extracted_at >= TIMESTAMP(DATE_SUB(CURRENT_DATE(), INTERVAL 2 DAY))")}
),
ranked_members AS (
SELECT
*,
ROW_NUMBER() OVER (PARTITION BY extracted_date, ptag ORDER BY extracted_at DESC) AS row_num
FROM
parsed_members
),
deduped_members AS (
SELECT
extracted_at,
extracted_date,
ptag,
troops
FROM
ranked_members
WHERE
row_num = 1
)
SELECT
extracted_at,
extracted_date,
ptag,
JSON_VALUE(troop.name) AS troop_name,
SAFE_CAST(JSON_VALUE(troop.level) AS INT64) AS troop_level,
SAFE_CAST(JSON_VALUE(troop.maxLevel) AS INT64) AS troop_max_level,
JSON_VALUE(troop.village) AS troop_village
FROM
deduped_members,
UNNEST(JSON_QUERY_ARRAY(troops)) AS troop
QUALIFY ROW_NUMBER() OVER (PARTITION BY extracted_date, ptag, troop_name ORDER BY troop_level DESC) = 1
Este cambio es fundamental: estas inyecciones de “JavaScript” hacen que podamos referenciar tablas, pero también nos permiten detallar el DAG. El grafo acíclico nos permite entender visualmente la secuencia lógica de ejecución. También es fundamental para la gobernanza de datos, pero esto lo veremos en otra publicación. El DAG de este proyecto el día de hoy se ve así:

Buenas prácticas
Cuando vayas a utilizar Dataform o cualquier servicio que te permita implementar DaC, ten en cuenta lo siguiente:
- Arquitectura incremental inteligente: filtra únicamente por particiones para optimizar recursos, tanto financieros como computacionales.
- Particionamiento: las tablas deben estar particionadas por un atributo de temporalidad.
- Clusterización: tienes hasta cuatro columnas para utilizar. Esto es como crear un índice; hazlo sabiamente. Básate en la frecuencia de uso de los atributos dentro de tus WHERE statements.
- Deduplicación: ten en cuenta estrategias de deduplicación cuando sea necesario. En el ejemplo anterior nos aseguramos únicamente de tener el registro más fresco; esta fue una decisión de diseño que trasciende desde el post anterior.
- Contrato de calidad de datos: utiliza assertions; es mejor que haga ruido para accionar.
- Casting y parseo: convierte tus datos; utiliza SAFE_CAST para evitar errores en atributos no formateados, especialmente si ingieres de una fuente que no controlas, como una API de terceros.
- Catalogación: cada campo debe tener una descripción. Si lo aplicas en todos los archivos, tendrás un diccionario de datos actualizado.
Referencias
- https://www.databricks.com/blog/what-is-medallion-architecture
- https://cloud.google.com/blog/products/data-analytics/transform-sql-into-sqlx-for-dataform/
- https://docs.cloud.google.com/dataform/docs/overview
- https://docs.cloud.google.com/bigquery/docs/partitioned-tables
- https://docs.cloud.google.com/bigquery/docs/clustered-tables