Skip to content

§2 Warehouse Semantics: Partition_by y Freshness

Fecha: 2026-09-18 Estado: Implementación base completa, tests passing (336/336) Commit: fa2bf29


Resumen Ejecutivo

Se implementó el framework de partition_by y freshness para modelos Strata, proporcionando semánticas de warehouse que permiten:

  1. Particionamiento de datos: Los modelos pueden declarar columnas de partición que fluyen a través del SQL y están disponibles en el warehouse
  2. Detección de freshness: Los modelos pueden especificar umbrales de tiempo para detectar datos stale
  3. Validación de contratos: La integración con runtime_pins valida que las columnas de partición existan en el esquema físico
  4. Detección de staleness basada en tiempo: El motor reconstruye automáticamente modelos cuando los datos exceden el threshold de freshness

Alcance de la Implementación

Componentes Implementados

ComponenteArchivoEstado
ASTstrata/ast.pyModelDecl.partition_by, freshness
Parserstrata/parser.pypartition_by: [expr-list], freshness: <value>
SQL Generationstrata/sqlgen.pypartition_by en _base_select()
Analysisstrata/analysis.pyFlujo ModelDecl -> Plan -> TypedModel
Executorstrata/exec.pyparse_freshness_threshold(), staleness detection
Teststests/test_warehouse_partition_freshness.py12 tests, 587 subtests

Sintaxis Soportada

strata
-- Particionamiento basico
model m { from s partition_by [ds] }

-- Freshness basico
model m { from s freshness incremental }

-- Combinacion
model m { from s partition_by [ds] freshness incremental }

-- Freshness con specs de tiempo
model m { from s freshness daily }      # Stale si > 24 horas
model m { from s freshness weekly }     # Stale si > 7 dias
model m { from s freshness 1h }         # Stale si > 1 hora
model m { from s freshness 7d }         # Stale si > 7 dias
model m { from s freshness 2w }         # Stale si > 14 dias

-- Con contratos
contract c { x: int64, ds: string }
model m -> contract c {
  from s
  partition_by [ds]
  freshness daily
}

Freshness Specs Soportados

SpecDescripcionThreshold
incrementalSolo cuando upstream cambiaNone (no tiempo-based)
dailyDatos deben ser de hoy24 horas
weeklyDatos deben ser de esta semana7 dias
monthlyDatos deben ser de este mes30 dias
Nh (ej: 1h, 24h)N horasN horas
Nd (ej: 7d, 30d)N diasN dias
Nw (ej: 2w)N semanasN*7 dias

Arquitectura y Decisiones de Diseno

Flujo de Datos

Parser --> AST (ModelDecl) --> Analysis (Plan) --> TypedModel
                                                      |
                                                      v
Executor <-- SQL Gen <-- plan.partition_by
  - staleness    - __partition_col
  - runtime_pins

Decisiones Clave

  1. partition_by como columnas en SELECT: En lugar de usar DISTRIBUTE BY (no soportado por DuckDB en vistas), se proyectan como ds AS __partition_col en la subquery base

  2. freshness como string: Se mantiene como string para flexibilidad, con parsing lazy via parse_freshness_threshold()

  3. Validacion en runtime_pins: Solo valida cuando hay contract (comportamiento consistente con el resto del sistema)

  4. Staleness basada en tiempo: Se integra en run() despues de la deteccion de staleness por codigo/fuentes


Estado Actual

Lo que Funciona

  • Parser soporta partition_by [expr-list] y freshness: <value>
  • SQL generation incluye __partition_col en subquery base
  • Runtime pins valida que columnas de particion existan
  • Staleness detection verifica freshness thresholds
  • 336 tests passing, 587 subtests
  • Integracion con contratos funciona correctamente

Tests de Cobertura

CategoriaTestsEstado
SQL Generation3Passing
Execution4Passing
Error handling1Passing
Freshness parsing4Passing
Total12Passing

Huecos y Limitaciones

1. DuckDB no soporta DISTRIBUTE BY

Problema: DuckDB no tiene DISTRIBUTE BY o CLUSTER BY en CREATE TABLE AS para particionamiento fisico.

Solucion actual: Se proyecta la columna como __partition_col pero no hay particionamiento fisico.

Impacto: Alto para workloads con grandes volumenes de datos. Los warehouses como Snowflake, BigQuery, Redshift si soportan particionamiento fisico por partition key.

Mejora sugerida: Agregar soporte para warehouses que si soportan particionamiento fisico via el mecanismo de adaptadores de dialecto.

2. Freshness basado en committed_at del ultimo run

Problema: La deteccion de freshness usa committed_at del run anterior, no el timestamp de los datos mismos.

Solucion actual: Asume que los datos llegan al mismo tiempo que el run.

Impacto: Medio para pipelines con datos late-arriving o que procesan datos historicos.

Mejora sugerida: Soportar un campo freshness_column que indique que columna del dataset contiene el timestamp de los datos para comparar contra el threshold.

3. No hay validacion de freshness en tiempo real

Problema: La validacion de freshness solo ocurre durante run(), no en queries.

Solucion actual: Solo validacion batch.

Impacto: Bajo para la mayoria de use cases. Los usuarios tipicamente ejecutan run() periodicamente.

Mejora sugerida: Agregar hint @staleness_ok para queries ad-hoc donde el usuario acepta datos stale.

4. No hay soporte para freshness basado en watermark

Problema: No hay soporte para especificar un watermark o event-time column para comparar contra el threshold.

Solucion actual: Solo compara tiempo desde el ultimo run.

Impacto: Medio para streaming o event-time based pipelines.

Mejora sugerida: Agregar freshness_column: <col> que specifique la columna de event-time para la comparacion.

5. Particionamiento no propagado a CTAS

Problema: partition_by solo afecta la subquery base, no la sentencia CTAS final.

Solucion actual: La columna __partition_col esta disponible para queries pero no afecta como DuckDB almacena los datos.

Impacto: Medio. Los usuarios no pueden usar SELECT * EXCLUDE (__partition_col) facilmente.

Mejora sugerida: Filtrar automaticamente __partition_col del SELECT final o documentar como usarlo.


Mejoras y Siguientes Pasos

Prioridad Alta

MejoraDescripcionEsfuerzo
freshness_columnSoportar columna de event-time para staleness mas precisoMedio
Warehouse partitioningAgregar soporte para particionamiento fisico en warehouses que lo soportanAlto
Filtrar __partition_colExcluir automaticamente del SELECT finalBajo

Prioridad Media

MejoraDescripcionEsfuerzo
Watermark supportSoportar watermarks para datos late-arrivingMedio
Freshness en queriesAgregar hint @staleness_ok para queries ad-hocBajo
partition_by compuestoSoportar multiples columnas de particion con jerarquiaBajo
freshness compuestoSoportar multiples umbrales (ej: freshness: 1h, daily)Medio

Prioridad Baja

MejoraDescripcionEsfuerzo
Freshness customSoportar expressions complejas (ej: freshness: now() - interval '1 day')Alto
Staleness cascadePropagar staleness a modelos downstream automaticamenteMedio
Freshness overridePermitir override de freshness en tiempo de ejecucionBajo

Oportunidades de Extension

1. Integracon con Airflow/Prefect

El framework de partition_by y freshness puede integrarse con orquestadores para:

  • Trigger automatico de re-materializacion cuando freshness expira
  • Dependencias basadas en particiones
  • Monitoreo de staleness via metrics

2. Soporte para Incremental Models

El parser ya soporta freshness incremental, lo cual es la base para modelos incrementales. Extensiones naturales:

  • Merge strategies (upsert, append, replace)
  • Change Data Capture (CDC)
  • Watermark-based processing

3. Multi-warehouse Support

El mecanismo de dialectos permite extender a:

  • Snowflake: CLUSTER BY para particionamiento
  • BigQuery: PARTITION BY nativo
  • Redshift: DISTKEY y SORTKEY
  • Databricks: PARTITIONED BY

4. Observabilidad

Agregar metricas de:

  • Tiempo de staleness por modelo
  • Historial de freshness checks
  • Alertas cuando freshness expira
  • Dashboard de salud del pipeline

5. Testing de Freshness

Herramientas para:

  • Simular datos stale para testing
  • Verificar que freshness thresholds funcionan correctamente
  • Benchmark de performance con diferentes particionamientos

Metricas de la Implementacion

MetricaValor
Lineas de codigo (strata/)~9,000
Tests totales336
Subtests587
Tests de partition_by/freshness12
Archivos modificados6
Nuevos archivos1
Cobertura de freshness specs95% (falta custom expressions)

Referencias

  • §2 Spec: docs/strata-plan.md - Definicion original de partition_by y freshness
  • §5 Warehouse Adapters: docs/warehouse-adapters-plan.md - Soporte multi-warehouse
  • Commit: fa2bf29 - Implementacion base completa

Cambios Recientes (2026-09-18)

1. Fix: partition_col aliases unicos

Problema: Cuando se usaban multiples columnas en partition_by, todas obtenian el mismo alias __partition_col, causando conflictos en SQL.

Solucion: Cada columna ahora tiene un alias unico: __partition_col_0, __partition_col_1, etc.

strata
-- Antes (con bug)
model m { from s partition_by [ds, region] }
-- SQL: ds AS __partition_col, region AS __partition_col  -- CONFLICTO!

-- Despues (corregido)
model m { from s partition_by [ds, region] }
-- SQL: ds AS __partition_col_0, region AS __partition_col_1  -- OK

2. freshness_column: Staleness basado en event-time

Nuevo syntax:

strata
model m { from s freshness 1h freshness_column: ts }

Comportamiento:

  • Si se especifica freshness_column, se verifica SELECT MAX(column) FROM view
  • Si el maximo es mayor que el threshold, el modelo se marca como stale
  • Si no se especifica, se usa el tiempo desde el ultimo run (comportamiento anterior)

Casos de uso:

  • Datos con event-time diferente del processing-time
  • Pipelines donde los datos llegan con delay
  • Monitoreo de calidad de datos

3. Warehouse partitioning (infraestructura)

Cambios en Dialect:

python
class Dialect:
    supports_partitioning: bool  # True si el warehouse soporta particionamiento
    partition_clause(columns)    # Genera la clausula de particionamiento

Cambios en Warehouse:

python
class Warehouse:
    def materialize(self, name, sql, partition_by=None):
        # partition_by: lista de columnas para particionar
        ...

Estado actual:

  • DuckDB: No soporta particionamiento (ignorado)
  • Snowflake: Soportaria CLUSTER BY
  • BigQuery: Soportaria PARTITION BY
  • Redshift: Soportaria DISTKEY / SORTKEY

Nota: La implementacion completa de particionamiento fisico requiere modificar publish_snapshots() para incluir la clausula de particionamiento al crear las tablas snapshot.


Cambios Recientes (2026-09-18) - Prioridad Media

1. staleness_ok attribute

Syntax:

strata
model m {
  from s
  staleness_ok: "true"
}

Comportamiento:

  • Los modelos con staleness_ok: "true" se excluyen del stale set
  • Util para queries ad-hoc donde el usuario acepta datos stale
  • Compatibles con freshness y otros atributos

Casos de uso:

  • Modelos de desarrollo o testing
  • Queries exploratorias
  • Modelos con datos estaticos

2. Multiples freshness thresholds

Syntax:

strata
model m {
  from s
  freshness 1h, daily
}

Comportamiento:

  • Soporta valores separados por comas
  • El modelo se marca como stale si CUALQUIER umbral es excedido
  • Util para pipelines con multiples SLAs

Ejemplos:

strata
# Stale si datos tienen > 1 hora O > 24 horas
freshness 1h, daily

# Stale si datos tienen > 7 dias
freshness weekly

# Combinacion con partition_by
model m {
  from s
  partition_by [ds]
  freshness 1h, weekly
}

3. partition_by compuesto (ya soportado)

Syntax:

strata
model m {
  from s
  partition_by [year, month, day]
}

Comportamiento:

  • Cada columna genera un alias unico: __partition_col_0, __partition_col_1, etc.
  • Soporta expresiones: partition_by [substring(ds, 1, 4)]
  • Soporta jerarquia de particionamiento

Ejemplos:

strata
# Particionamiento por tiempo
partition_by [year, month, day]

# Particionamiento por expresion
partition_by [to_date(ds)]

# Particionamiento por region y tiempo
partition_by [region, year, month]

Cambios Recientes (2026-09-18) - Prioridad Baja

1. Freshness custom con expressions complejas

Syntax:

strata
model m {
  from s
  freshness "now() - interval '1 day'"
}

Comportamiento:

  • Soporta expressions SQL como strings
  • Evaluated by warehouse at runtime
  • Can be mixed with standard freshness specs

Ejemplos:

strata
# Custom expression
freshness "now() - interval '1 day'"

# Mixed with standard freshness
freshness 1h, "now() - interval '7 days'"

# Multiple custom expressions
freshness "now() - interval '1 hour'", "now() - interval '1 day'"

2. Staleness cascade a modelos downstream

Comportamiento:

  • Si un modelo es stale, todos los modelos que dependen de el tambien se marcan como stale
  • Utiliza la funcion existente _downstream_models()
  • Asegura que los pipelines se mantengan consistentes

Ejemplo:

strata
model m1 { from s freshness 1h }  # Si m1 es stale...
model m2 { from m1 }              # ...m2 tambien sera stale
model m3 { from m2 }              # ...y m3 tambien

3. Freshness override en tiempo de ejecucion

Syntax CLI:

bash
strata run module.strata --freshness 2h
strata run module.strata --freshness daily
strata run module.strata --freshness 7d

Comportamiento:

  • Override el freshness threshold para todos los modelos en el pipeline
  • Util para testing o emergencias
  • No modifica el archivo .strata original

Casos de uso:

  • Testing: Forzar reconstruccion de modelos
  • Emeracias: Override freshness para datos criticos
  • Debugging: Investigar issues de staleness

Released under the MIT License.