Saltar a contenido

La capa canónica: Polars, Parquet y NumPy contiguo

Entre el CSV del proveedor y el indicador hay una capa que casi nadie ve y de la que depende todo lo demás: la capa canónica. Es el sitio donde las barras dejan de ser «lo que dijo el proveedor» y pasan a ser «lo que el proyecto garantiza». Este capítulo explica su esquema, cómo se particiona, cómo se limpia, cómo se lee y cómo se demuestra que dos corridas usaron exactamente los mismos datos.

Qué vas a aprender

  • El esquema canónico de una barra y sus tres relojes.
  • Cómo se organiza el Parquet en particiones (154 de Twelve Data + 480 de Yahoo).
  • La política de limpieza clean.v1: excluir, reparar, contar… y nunca rellenar.
  • Qué es el data_hash, el manifest.json y el validation_report.json.
  • Cómo se elige la base de precio (price_basis) en cada corrida.
  • Qué son BarArrays y el panel multi-activo, y cómo se leen desde finazbench/data/.
  • Cómo se usa BASELINE_MANIFEST.json para reproducir por hashes.

1. La idea: una sola verdad, tres herramientas

flowchart LR
    CSV[CSV congelado] -->|tools/prep_data.py<br/>Polars| VAL[validate + clean.v1]
    VAL --> PQ[(Parquet zstd-9<br/>particionado)]
    PQ --> MAN[manifest.json<br/>validation_report.json]
    PQ -->|loader.load_bars<br/>scan_parquet lazy| NP[BarArrays<br/>NumPy C-contiguo]
    NP --> IND[Indicadores<br/>VectorTA / Nautilus / canónico]

Cada herramienta hace lo que mejor sabe:

Herramienta Papel
Polars Transformar y leer: lectura perezosa (scan_parquet), proyección de columnas, filtros con predicate pushdown.
Parquet Almacenar: columnar, comprimido (zstd nivel 9), con particiones por fuente/mercado/activo/marco/año.
NumPy contiguo Calcular: vector_ta exige float64 contiguo y los wranglers de Nautilus consumen columnas.

Analogía

Parquet es el almacén ordenado por pasillos; Polars es la carretilla que trae solo las cajas pedidas; NumPy es la mesa de trabajo donde se opera. Nada se mueve dos veces si no hace falta.


2. El esquema canónico de una barra

Columna Tipo Significado
asset_id Utf8 Activo (NVDA, BBVA.MC, EURUSD=X…)
venue Utf8 Plaza (XNAS, XNYS…)
market_open_ns Int64 Apertura económica del intervalo, UTC ns
market_close_ns Int64 Fin económico del intervalo, UTC ns
data_available_ns Int64 Desde cuándo puede usarse el dato, UTC ns
session_id Int32 Sesión de bolsa, AAAAMMDD
open, high, low, close Float64 OHLC
volume Float64 Volumen
source_row_id Utf8 Identidad de la fila en el fichero fuente
is_tradable Boolean Si se puede operar sobre la barra
is_synthetic Boolean Si la barra fue reparada por la limpieza

Yahoo añade adj_close y adj_factor (columnas opcionales).

2.1 Los tres relojes

Para la barra t el intervalo es semiabierto: [market_open_ns, market_close_ns).

sequenceDiagram
    participant M as Mercado
    participant D as Dato
    Note over M: market_open_ns[t]
    M->>M: se negocia la barra t
    Note over M: market_close_ns[t]
    M->>D: data_available_ns[t] = market_close_ns[t]<br/>(latencia de publicación cero, supuesto declarado)
    D->>M: la decisión tomada al cierre de t<br/>se ejecuta en open[t+1]
  • data_available_ns = market_close_ns es un supuesto sintético declarado en el manifiesto (availability_policy: zero_publication_latency.v1). Twelve Data no documenta su latencia real; asumir cero es lo que permite el carril SIM-S next-open.
  • El lector filtra por data_available_ns, nunca por market_open_ns: pedir «desde el 1 de enero» significa «desde que se pudo saber algo».

Si algún día se midiera la latencia real

data_available_ns dejaría de coincidir con el cierre y el dataset pasaría a ser SIM-R, no SIM-S. No se mezclan.


3. Particionado

data/canonical/
  source=twelvedata/market=XNYS/asset=NVDA/freq=1min/year=2024/part.parquet
                                           freq=5min/year=2024/...
                                           freq=1d_long/year=1999/...
  source=yahoo/market=XNYS/asset=AMZN/freq=1d/year=2024/part.parquet
              /market=XMAD/asset=BBVA.MC/freq=1d/year=2024/...
              /market=FX/asset=EURUSD=X/freq=1d/year=2024/...
  manifest.json            validation_report.json          (twelvedata)
  manifest.yahoo.json      validation_report.yahoo.json    (yahoo)
Fuente Particiones Manifiesto
Twelve Data (NVDA, KO; 1min, 2, 5, 15, 30min, 1h, 1d y 1d_long) 154 manifest.json
Yahoo (40 series diarias) 480 manifest.yahoo.json

Tres decisiones que conviene entender:

  1. El año se calcula en hora local del mercado, no en UTC: la barra de las 15:59 del 31 de diciembre pertenece a ese año aunque en UTC ya sea 1 de enero.
  2. 1d y 1d_long están separados: 1d se deriva del intradía (2020→) y 1d_long viene del diario del proveedor (1999→ o 1970→). Mezclarlos daría dos valores distintos para la misma sesión.
  3. Un manifiesto por fuente: las fuentes no se mezclan y dos corridas de prep_data no se pisan.
El asset_id con = en la ruta

asset=EURUSD=X funciona con glob y pl.scan_parquet, pero la inferencia hive de Polars deja caer la columna asset en silencio (DEV-01B-02). No afecta porque el lector nunca usa hive_partitioning y asset_id es también una columna de datos.


4. La política de limpieza clean.v1

Versionada en finazbench/data/validate.py y copiada al manifiesto. Cambiarla cambia el data_hash aunque el CSV sea idéntico.

# Regla NVDA 1min KO 1min
(a) Barras fuera de sesión regular → se excluyen y se cuentan 360 337
(b) Invariante OHLC rota → se repara y is_synthetic=True 20 464
(c) Timestamps duplicados → se colapsa, se queda la última 0 0
(d) Huecos dentro de sesión → se cuentan, nunca se rellenan 453 342
flowchart TD
    B[Barra del CSV] --> Q1{¿Dentro de sesión<br/>según calendario?}
    Q1 -- no --> X[Excluir y contar<br/>out_of_session]
    Q1 -- sí --> Q2{¿low ≤ open,close ≤ high?}
    Q2 -- no --> R[Reparar: high=max, low=min<br/>conservar open y close<br/>is_synthetic=True]
    Q2 -- sí --> OK[Barra canónica]
    R --> OK
    OK --> G[Huecos: contar gaps_in_session<br/>NUNCA rellenar]

4.1 Reparar en vez de excluir

Las 20 filas rotas de NVDA lo están por menos de un tick (desviación máxima 0,0049 USD) y high < low no ocurre nunca. Excluirlas abriría huecos de un minuto que obligan a pausar la cadena de indicadores, para «arreglar» un error más pequeño que la unidad mínima de precio. Se conservan open y close porque son precios de transacción observados.

4.2 No rellenar nunca

El forward-fill hace trampas a tu favor

Una vela inventada tiene high == low == close: un ATR la lee como volatilidad cero y el retorno de ese día es cero exacto. El backtest sale más suave que la realidad por construcción y el Sharpe mejora. Por eso los huecos se cuentan y se declaran (gaps_in_session, gap_sessions_before), pero no se tapan.

4.3 Así se ve el informe de validación

Extracto real de data/canonical/validation_report.json para NVDA 1min:

{
  "asset_id": "NVDA",
  "calendar_version": "XNYS-1999-2026.v1:4f047091cf176d61",
  "counts": {
    "gaps_in_session": 453,
    "high_below_low": 0,
    "ohlc_repaired": 20,
    "out_of_session": 360,
    "out_of_session_zero_volume": 359,
    "rows_in": 632095,
    "rows_out": 631735,
    "sessions": 1627
  },
  "ok": true,
  "policy_version": "clean.v1",
  "timeframe": "1min"
}

5. El remuestreo dentro de sesión

finazbench/data/resample.py es la única implementación de agregación del proyecto (open = primero, high = máximo, low = mínimo, close = último, volume = suma).

  1. Los cubos se anclan a la apertura de sesión del calendario, no a medianoche ni a la primera barra presente (corrección DEV-01B-09: afectaba a 7 sesiones con desfases de hasta 241 min).
  2. Ningún cubo cruza el límite de sesión.
  3. El intervalo es semiabierto.

Comparación con la vía antigua (pandas sin calendario)

Sobre NVDA 5min, la vía antigua producía 126.539 barras y la canónica 126.467. En las 126.148 barras comunes, open y volume coinciden exactamente; solo 4 difieren en high/low/close, todas explicadas (relleno de volumen cero en sesiones cortas y filas reparadas). Un test exige que toda diferencia caiga en una categoría conocida.


6. Base de precio por corrida (price_basis)

El Parquet guarda siempre el OHLC del proveedor (y en Yahoo, adj_close y adj_factor). El ajuste se aplica al leer:

price_basis Qué devuelve Disponible en
raw_split_adjusted (por defecto) OHLC tal cual todas las series
total_return_adjusted open, high, low y close multiplicados por adj_factor series Yahoo de acciones

Se escalan los cuatro precios (no solo el cierre) porque el factor es un escalar positivo por barra y así se conservan las invariantes OHLC. Pedir total_return_adjusted a NVDA, KO o un par FX lanza un error: nunca se cae en silencio a la base cruda.

from finazbench.data.loader import load_bars

bars = load_bars(
    "BBVA.MC", "1d",
    market="XMAD", source="yahoo",
    price_basis="total_return_adjusted",
)
bars.validate()          # contrato: dtypes, contigüidad, monotonía estricta
print(bars.n, bars.close[:3])

7. De Parquet a memoria: BarArrays y el panel

7.1 BarArrays

finazbench/data/contracts.py::BarArrays es la estructura en memoria de un activo y un marco: un ndarray por columna (layout columnar), no un array de estructuras.

Campo dtype
market_open_ns, market_close_ns, data_available_ns int64
open, high, low, close, volume float64 C-contiguo
is_tradable bool
session_id int32
is_synthetic (opcional) bool
gap_sessions_before (opcional) int32

validate() comprueba longitudes, dtypes exactos, contigüidad en C y monotonía estricta de los tres relojes.

7.2 Cómo lee load_bars

finazbench/data/loader.py: scan_parquet perezoso, proyección explícita de 10 columnas, filtro por data_available_ns con predicate pushdown, collect, rechunk y conversión a NumPy intentando to_numpy(allow_copy=False).

Activo Marco Barras Lectura (mediana, 1 hilo)
NVDA 1min 631.735 51,6 ms
NVDA 5min 126.467 13 ms
NVDA 1d 1.627 3,7 ms
KO 1min 526.049 ~60 ms

La única copia inevitable

Las nueve columnas numéricas no copian. La única que copia es is_tradable: Arrow guarda los booleanos en un bit y NumPy en un byte, y no existe vista zero-copy entre ambos.

La caché BytesLRU limita por bytes, no por entradas: una entrada puede ser NVDA 1min (~44 MB) o KO 1d (~100 kB).

7.3 El panel multi-activo

finazbench/data/panel.py::align_panel(asset_ids, freq, start, end, price_basis) -> Panel devuelve matrices T×S de close y open y una máscara de disponibilidad. El eje de filas es la unión de sesiones (un festivo de BME no borra un día en que AMZN cotizó), y Panel.validate() exige que la máscara sea True exactamente donde el precio es finito: un forward-fill rompe la validación.

Sobre el nombre PanelArrays

Algunos documentos de arquitectura (docs/v2/) hablan de «BarArrays/PanelArrays». En el código la estructura de un activo es BarArrays (finazbench/data/contracts.py) y la del panel se llama Panel (finazbench/data/panel.py); no existe ninguna clase PanelArrays.


8. Reproducibilidad por hashes

8.1 El data_hash

Es el SHA-256 de la lista ordenada ruta|sha256_fichero|filas, más las versiones de esquema, calendario, política de limpieza y metadatos de instrumento: el hash de los hashes.

Cambia si… No cambia si…
cambia un byte de una partición se regenera el dataset sin tocar nada (created_at no entra)
se añade o quita una partición
cambia un festivo o la política de limpieza
cambia el tick de un instrumento

data_hash vigente de Twelve Data (el que citan la paridad y los benchmarks):

e8b5fc3612fd944a692af21807706f665fb5bec4f1251d1652fbeaabea6595b0

Regla de oro

Dos cifras solo se comparan si comparten data_hash. Si se regenera data/canonical, la matriz de paridad queda anticuada y hay que repetirla, no mezclarla.

8.2 BASELINE_MANIFEST.json

migration/baseline/BASELINE_MANIFEST.json congela la base de referencia de la plataforma: data_hash, manifest_sha256, validation_report_sha256, las 154 particiones, las dependencias fijadas (vector-ta==0.2.8, nautilus-trader==1.231.0, numpy==2.5.3, polars==1.44.2…), los hashes de ficheros clave (finazbench/features/canonical.py, finazbench/ledger/reference_ledger.py, catalog/strategy_catalog.json…) y las desviaciones declaradas (p. ej., numpy 2.5.1 local frente a 2.5.3 fijado).

# comprobar a mano que el manifiesto no ha cambiado
sha256sum data/canonical/manifest.json
# debe coincidir con data.manifest_sha256 del BASELINE_MANIFEST:
# bb2d249300fd2cb54b22a9695984cd4071859bdf49c8da74bc752253a775ac19

8.3 Regenerar

docker compose --profile dev run --rm -T dev \
    python tools/prep_data.py --symbol all --sources data/raw/twelvedata --out data/canonical
docker compose --profile dev run --rm -T dev \
    python tools/prep_data.py --source yahoo --out data/canonical
docker compose --profile dev run --rm -T dev python -m pytest tests/data -q

La escritura es atómica (temporal + rename): otros procesos leen data/canonical en paralelo y un Parquet a medio escribir sería un fichero corrupto.


Resumen

  • La capa canónica convierte CSV en Parquet particionado con un esquema fijo y tres relojes.
  • clean.v1 excluye lo que está fuera de sesión, repara y marca OHLC rotos y cuenta los huecos sin rellenarlos.
  • La base de precio se decide al leer (price_basis) y nunca se sustituye en silencio.
  • BarArrays entrega columnas NumPy contiguas; el panel T×S lleva una máscara que impide el forward-fill.
  • El data_hash identifica el dataset entero; sin el mismo hash no hay comparación válida.

Para practicar

  1. Abre data/canonical/validation_report.json y compara rows_in y rows_out de NVDA 1min. Explica la diferencia con las reglas (a)–(d).
  2. Carga BBVA.MC con las dos bases de precio y comprueba que high >= close >= low se cumple en ambas. ¿Por qué fallaría si solo escalaras el cierre?
  3. Enumera tres cambios que moverían el data_hash sin tocar ningún CSV.
  4. Construye un panel AMZN + BBVA.MC para 2024 y cuenta las celdas con máscara False. ¿Qué festivos explican cada una?