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, elmanifest.jsony elvalidation_report.json. - Cómo se elige la base de precio (
price_basis) en cada corrida. - Qué son
BarArraysy el panel multi-activo, y cómo se leen desdefinazbench/data/. - Cómo se usa
BASELINE_MANIFEST.jsonpara 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_nses 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 pormarket_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:
- 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.
1dy1d_longestán separados:1dse deriva del intradía (2020→) y1d_longviene del diario del proveedor (1999→ o 1970→). Mezclarlos daría dos valores distintos para la misma sesión.- Un manifiesto por fuente: las fuentes no se mezclan y dos corridas de
prep_datano 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).
- 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).
- Ningún cubo cruza el límite de sesión.
- 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):
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.v1excluye 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. BarArraysentrega columnas NumPy contiguas; el panel T×S lleva una máscara que impide el forward-fill.- El
data_hashidentifica el dataset entero; sin el mismo hash no hay comparación válida.
Para practicar¶
- Abre
data/canonical/validation_report.jsony compararows_inyrows_outde NVDA 1min. Explica la diferencia con las reglas (a)–(d). - Carga BBVA.MC con las dos bases de precio y comprueba que
high >= close >= lowse cumple en ambas. ¿Por qué fallaría si solo escalaras el cierre? - Enumera tres cambios que moverían el
data_hashsin tocar ningún CSV. - Construye un panel AMZN + BBVA.MC para 2024 y cuenta las celdas con máscara
False. ¿Qué festivos explican cada una?