Presentamos Pulsora: una base de datos de series temporales en Rust para datos de trading

Teníamos un torrente y ningún buen lugar donde ponerlo.

El problema interno era concreto: almacenar y consultar un flujo grande y continuo de datos de mercado — ticks, cotizaciones, barras de muchos símbolos — lo bastante rápido como para que la ingesta nunca se convirtiera en el cuello de botella y las consultas por rango de tiempo volvieran sin que tuviéramos que contener la respiración. Los datos de mercado son la versión más limpia y más exigente de una carga de trabajo de series temporales. Llegan en orden, llegan constantemente y no paran nunca. El volumen no es interesante por ser grande una vez. Es interesante porque es grande cada segundo, para siempre.

Primero probamos lo obvio: recurrir a una base de datos de series temporales existente. Esa suele ser la decisión correcta, y durante un tiempo estuvo bien. Luego los requisitos de rendimiento de escritura fueron subiendo, la huella en disco subió aún más rápido y las consultas por rango de tiempo que realmente ejecutábamos empezaron a notar el peso del índice. En algún momento la capa de almacenamiento hacía más contabilidad por fila que trabajo sobre la fila.

Así que construimos Pulsora — una base de datos de series temporales en Rust, optimizada para datos de mercado y conjuntos de datos ordenados por tiempo. Es código abierto en GitHub bajo Apache-2.0. Esta entrada es una introducción a qué es, por qué existe y cómo almacena y consulta datos en realidad — basada en el código, no en el folleto.


Por qué construir en lugar de adoptar

Construir una base de datos no es algo que se haga a la ligera. El listón que había que superar era: ¿qué necesitamos específicamente que no estábamos obteniendo a un coste razonable?

Tres cosas aparecían una y otra vez.

Rendimiento de escritura que escala con los bloques, no con las filas. Un flujo de ticks son millones de filas pequeñas y casi idénticas. La mayoría de los almacenes de propósito general quieren indexar cada una de ellas. Eso es una entrada de índice por fila, y con nuestros recuentos de filas el índice se convierte en el coste dominante — en amplificación de escritura, en disco, en compactación. Queríamos que el coste de almacenamiento por fila tendiera a cero y que el índice escalara con el número de bloques que escribimos, no con el número de filas.

Compresión que entiende los datos. Los datos de mercado se comprimen de maravilla si los tratas como columnas de un tipo conocido. Las marcas de tiempo llegan a intervalos casi regulares. Los precios cambian lentamente. Los volúmenes son enteros pequeños. La compresión genérica orientada a filas deja casi todo eso sin aprovechar. Queríamos compresión columnar consciente del tipo — del tipo que describió el artículo de Facebook sobre Gorilla — integrada, no añadida a posteriori.

Una interfaz aburrida y flexible en cuanto a formatos. Ingerimos desde varios productores distintos y consumimos hacia varias herramientas distintas. No queríamos un protocolo de cable a medida. Queríamos hacer POST de CSV desde un script, transmitir Apache Arrow desde una canalización y leer de vuelta JSON, Arrow o CSV según quién pregunte. HTTP y formatos columnares bien conocidos, nada exótico.

Ninguna de estas cosas es novedosa por sí sola. La combinación — para nuestra carga de trabajo, con nuestras restricciones de huella — bastó para justificar escribirla. Y una vez que ya la estás construyendo, hacerlo en Rust para obtener corrección en tiempo de compilación en la ruta crítica es lo fácil. Ya escribimos antes sobre por qué construimos nuestras propias herramientas; Pulsora encaja de lleno en ese linaje.


Qué es Pulsora hoy

Seré honesto sobre el alcance: Pulsora es 0.1.0. Es un motor que hace bien su trabajo, no una plataforma con clústeres, replicación y multitenencia con planificador de consultas. Es un almacén de series temporales columnar, de un solo nodo, con RocksDB embebido y una API REST. Eso es todo, y eso es precisamente el punto.

Esta es su forma:

  • Backend de almacenamiento: RocksDB, con una capa columnar propia encima.
  • Formatos de ingesta: CSV, flujo Apache Arrow IPC y Protocol Buffers.
  • Formatos de salida de consultas: JSON (predeterminado), flujo Arrow IPC, Protobuf y CSV — negociados mediante la cabecera Accept.
  • Esquema: inferido automáticamente en la primera ingesta. No declaras tablas; haces POST de los datos y Pulsora deduce las columnas y los tipos.
  • Interfaz: una API HTTP/REST servida por Axum sobre Tokio.
  • Durabilidad: un registro de escritura anticipada (WAL) opcional para que las filas en búfer sobrevivan a una caída.

Todo el conjunto es un único binario llamado pulsora. Hay un solo flag de línea de comandos.

# Ejecutar con valores predeterminados
cargo run

# O apuntarlo a un archivo de configuración
cargo run -- -c pulsora.toml

Ese flag -c/--config es toda la superficie de línea de comandos. Todo lo demás vive en TOML. Esto es deliberado — la configuración es la API para operarlo, y la API en tiempo de ejecución es HTTP.


Cómo almacena los datos en disco

Esta es la parte que más nos importaba, así que es la parte que merece la pena explicar.

Bloques columnares, no filas

Las filas entran, pero no se almacenan como filas. Pulsora acumula las filas entrantes en un búfer en memoria, y cuando el búfer alcanza su umbral de tamaño (buffer_size, 1000 por defecto) o su umbral de tiempo (flush_interval_ms), las agrupa en un ColumnBlock — un fragmento de filas almacenado columna por columna, cada columna comprimida de forma independiente con un algoritmo elegido para su tipo.

El bloque es la unidad de almacenamiento. Un ColumnBlock es, a grandes rasgos:

pub struct ColumnBlock {
    pub row_count: usize,                        // filas en este bloque
    pub columns: HashMap<String, Vec<u8>>,       // datos de columna comprimidos
    pub null_bitmaps: HashMap<String, Vec<u8>>,  // seguimiento de nulos por columna
}

Almacenar por columna en lugar de por fila es lo que hace que la compresión funcione, porque todos los valores de una columna tienen el mismo tipo y una magnitud similar. Esa es toda la razón por la que existen los almacenes columnares, y los datos de series temporales son la carga de trabajo para la que se inventaron.

Compresión específica por tipo

Cada tipo de columna recibe una estrategia de compresión hecha para su forma. Según la documentación del repositorio:

Tipo de dato Compresión Ratio típico (cifras del repo) Por qué funciona
Timestamp Delta-of-delta + varint 5–10x Intervalos casi regulares
Float XOR (Gorilla) + varfloat 2–5x Valores que cambian despacio
Integer Delta + varint 3–8x Datos secuenciales / contador
String Codificación por diccionario 2–4x Valores repetidos (símbolos)
Boolean Codificación por longitud (RLE) 10–50x Datos dispersos

La ruta de las marcas de tiempo es el caso de manual. Los ticks llegan a intervalos casi constantes, así que la diferencia de segundo orden (delta-of-delta) suele ser cero o diminuta, y un varint codifica un número diminuto en un byte. Los floats usan el empaquetado de bits XOR-contra-el-valor-anterior del algoritmo Gorilla — cuando un precio apenas se mueve, la mayoría de los bits del XOR son cero y se empaquetan en nada. Las cadenas como los símbolos de cotización se repiten sin cesar, así que un diccionario asigna cada símbolo distinto a un entero pequeño una sola vez.

Por debajo de todo eso, RocksDB aplica una segunda capa de compresión de bloque genérica — lz4 por defecto, con none, snappy y zstd disponibles. Así obtienes primero codificación consciente del tipo, y luego una pasada de propósito general por encima.

Estos ratios son las propias afirmaciones de benchmark del repositorio, medidas con las suites de Criterion bajo benches/. Los reporto como las cifras del proyecto, no como un veredicto independiente — ejecuta cargo bench contra tus propios datos y confía en eso.

El índice que no crece por fila

Esta es la decisión de diseño sobre la que gira todo.

La mayoría de los almacenes de series temporales mantienen un índice por fila para poder encontrar cualquier fila individual por tiempo. Pulsora no. Mantiene un índice solo a nivel de bloque — metadatos sobre cada bloque, nunca sobre filas individuales. La disposición de claves en disco, según los documentos de arquitectura, se ve así:

Block index: [table_hash:u32]['B'][min_ts:i64][block_id:u64]
             value: [block_id][min_ts][max_ts][rows][min_id][max_id]   — una entrada por bloque
Block data:  [table_hash:u32]['D'][block_id:u64]                       — el bloque columnar comprimido
Overrides:   [table_hash:u32]['O'][block_id:u64]                       — posiciones de filas muertas

El table_hash es un hash FNV-1a de 4 bytes del nombre de la tabla, usado como prefijo de clave para que las tablas estén aisladas sin arrastrar nombres de cadena por cada clave. Tras el prefijo, un solo byte ('B', 'D', 'O') separa el índice, los datos y el conjunto de invalidaciones.

La consecuencia: el coste de almacenamiento escala con el número de bloques, no con el número de filas. Escribe mil millones de filas y no añades ninguna entrada de índice más allá de los bloques que las contienen. Eso era el requisito que no podíamos comprar a buen precio, y es el núcleo de por qué existe Pulsora.

Las actualizaciones de filas se manejan sin reescribir bloques. Un REPLACE marca la posición de la copia sustituida en el conjunto de invalidaciones del bloque antiguo — un conjunto, por bloque, de posiciones de filas muertas, fusionado mediante un operador de merge de RocksDB (unión de conjuntos) en el mismo write batch que los datos nuevos. La validez es una lectura puntual por bloque, no una búsqueda por fila.


La ruta de escritura

La ingesta es un POST. El esquema se infiere en el primero.

curl -X POST http://localhost:8080/tables/stocks/ingest \
  -H "Content-Type: text/csv" \
  --data-binary @ticks.csv

La fila de cabecera de ese CSV se convierte en el esquema. Pulsora muestrea los valores, detecta el tipo de cada columna (Integer, Float, String, Boolean, Timestamp) e identifica la columna de marca de tiempo automáticamente — entiende RFC 3339, YYYY-MM-DD HH:MM:SS, solo fecha y marcas de tiempo Unix tanto en segundos como en milisegundos. A partir de ahí, el esquema de esa tabla queda fijado y los datos entrantes se validan contra él.

Para canalizaciones de mayor rendimiento, sáltate por completo el parseo de CSV y transmite Apache Arrow:

curl -X POST http://localhost:8080/tables/stocks/ingest \
  -H "Content-Type: application/vnd.apache.arrow.stream" \
  --data-binary @ticks.arrow

Arrow llega con su esquema ya adjunto, así que no hay paso de inferencia ni parseo de texto — las columnas ya son columnas. La ingesta de Protobuf funciona igual para los productores que ya lo hablan.

Internamente, el flujo es: parsear la entrada, validar o inferir el esquema, almacenar las filas en un búfer en memoria (deduplicando al entrar, gana la última escritura) y — si el WAL está activado — añadir cada fila a un archivo .wal antes de que toque la memoria, para que una caída no pierda los datos en búfer. Cuando el búfer se vacía, las filas se convierten en un ColumnBlock comprimido, el bloque y su entrada de índice se escriben en RocksDB en un solo batch, y el WAL se trunca.


La ruta de lectura

Las consultas son GET con un rango de tiempo y paginación.

# Todo (limitado por el límite predeterminado)
curl "http://localhost:8080/tables/stocks/query"

# Una ventana de tiempo
curl "http://localhost:8080/tables/stocks/query?start=2024-01-01T09:30:00&end=2024-01-01T16:00:00"

# Paginado
curl "http://localhost:8080/tables/stocks/query?limit=1000&offset=0"

start y end aceptan el mismo conjunto flexible de formatos de marca de tiempo que la ingesta. limit es 1000 por defecto, sin máximo impuesto; offset salta filas para la paginación.

El motor de consultas resuelve la petición primero contra el índice de bloques: construye un rango de claves binario a partir de los límites de tiempo e itera solo los bloques cuyo [min_ts, max_ts] se solapa con la ventana. Recopila los bloques candidatos, obtiene cada bloque único una vez, lo descomprime, cachea el resultado descomprimido para que el acceso repetido dentro de una consulta no vuelva a descomprimir, aplica el conjunto de invalidaciones para omitir filas muertas, extrae las filas que necesita y las serializa en el formato que haya pedido la cabecera Accept.

¿Quieres el resultado de vuelta como Arrow para una canalización posterior en lugar de JSON? Cambia una cabecera:

curl "http://localhost:8080/tables/stocks/query?limit=10000" \
  -H "Accept: application/vnd.apache.arrow.stream" > out.arrow

curl "http://localhost:8080/tables/stocks/query?limit=10000" \
  -H "Accept: text/csv" > out.csv

La misma consulta, cuatro posibles formatos de salida, sin paso de conversión por tu parte. Hay un puñado de otros endpoints de lectura — GET /tables para listar tablas, GET /tables/{name}/schema para el esquema inferido, GET /tables/{name}/count para un recuento de filas, GET /tables/{name}/row/{id} para obtener una sola fila por id y GET /health.


Pruébalo

Pulsora compila con un toolchain actual de Rust y se ejecuta como un único binario.

# Clonar y ejecutar desde el código fuente
git clone https://github.com/muvon/pulsora.git
cd pulsora
cargo run

# O instalar el binario directamente
cargo install --git https://github.com/muvon/pulsora.git

# Ingerir un CSV
curl -X POST http://localhost:8080/tables/stocks/ingest \
  -H "Content-Type: text/csv" \
  -d "timestamp,symbol,price,volume
2024-01-01 09:30:00,AAPL,150.00,1000
2024-01-01 09:31:00,AAPL,150.25,1500
2024-01-01 09:32:00,AAPL,149.75,2000"

# Leerlo de vuelta
curl "http://localhost:8080/tables/stocks/query?limit=10"

# Mirar el esquema inferido
curl "http://localhost:8080/tables/stocks/schema"

El ajuste vive en pulsora.toml. Los parámetros que conviene conocer primero:

[storage]
data_dir = "./data"
buffer_size = 1000          # filas en búfer antes de escribir un bloque
flush_interval_ms = 1000    # tiempo máximo que una fila espera en el búfer
wal_enabled = true          # las filas en búfer sobreviven a una caída

[ingestion]
batch_size = 10000          # filas por batch de escritura en la base de datos

[performance]
compression = "lz4"         # none | snappy | lz4 | zstd
cache_size_mb = 256         # caché de bloques para lecturas

Si te importa la máxima compresión por encima de la latencia, pon flush_interval_ms = 0 para un modo solo por lotes que retiene las filas hasta que el búfer se llena, produciendo bloques más grandes y mejor comprimidos. Si te importa la latencia de lectura en datos calientes, sube cache_size_mb. El conjunto completo de parámetros está documentado en doc/CONFIGURATION.md.


Lo que no es, y qué viene después

Te ahorraré la decepción de descubrirlo por las malas. Pulsora hoy es de un solo nodo. No hay clústeres, ni replicación, ni SQL, ni planificador general de consultas, ni pushdown de predicados más allá del rango de tiempo, ni limitación de tasa. Las consultas filtran por tiempo y paginan; no hacen GROUP BY. El archivo doc/ARCHITECTURE.md lista las direcciones que estamos mirando — poda de columnas y pushdown de predicados, una caché de bloques respaldada en disco, almacenamiento por niveles caliente/templado/frío, y clústeres — pero eso es hoja de ruta, no realidad. Prefiero que conozcas los bordes a que choques con ellos.

Lo que es, hoy, es un motor enfocado que ingiere un torrente de datos ordenados por tiempo, los comprime con fuerza usando algoritmos que los entienden y sirve consultas por rango de tiempo sin un índice que se infla por fila. Ese era el problema que teníamos. Es el problema que Pulsora resuelve.


Código abierto

Pulsora está en GitHub bajo Apache-2.0. Es Rust — RocksDB para el almacenamiento, Axum y Tokio para la API, Arrow y Prost para los formatos columnar y protobuf — y compila con un simple cargo build. Las suites de benchmark están en benches/ por si quieres medir el rendimiento de ingesta y consulta contra tus propios datos en lugar de confiar en los nuestros.

Si te estás ahogando en filas ordenadas por tiempo — ticks de mercado, lecturas de sensores, cualquier cosa que llegue en orden y no pare nunca — clónalo, apúntale un flujo y dinos dónde se rompe. Esa es la forma más rápida de que mejore.

— Don

Pulsora es código abierto en GitHub en github.com/muvon/pulsora bajo la licencia Apache-2.0. ¿Encontraste un error o quieres una función? Abre un issue.