NATIVO EN GO • SUB-MICROSEGUNDO • SIN IMPUESTO DE JVM STREAMFLOW

Procesamiento de Eventos de Alto Rendimiento para Desarrolladores Go

StreamFlow es un framework ligero y con tipado seguro para procesamiento de streams en Go. Ingiere, agrega y une flujos de eventos desde Kafka, NATS y RabbitMQ con latencia inferior a un microsegundo, estado embebido en árbol LSM y cero sobrecarga operativa de JVM.

Comenzar en GitHub →
go install github.com/santiagolertora/streamflow/cmd/streamflow@latest
CONTENEDOR ULTRA-LIGERO
15 MB
Arranque Ultra-Rápido (< 100ms)

Binario estático en Go sin dependencias externas. Elimina la configuración de heap de Java, el ajuste de Garbage Collection y contenedores pesados de más de 1GB.

pipeline.go DSL CON TIPADO SEGURO
orders := runtime.From[OrderEvent](source, pipeline)
engaged := orders.Filter("engaged", func(o OrderEvent) bool {
    return o.Amount > 100
})
engaged.To("premium-orders")
PROBLEMAS OPERATIVOS

Los Problemas Que Resolvemos

El procesamiento de eventos en tiempo real está plagado de complejidad operativa, latencias y código frágil.

⚠️ El Impuesto de Memoria de la JVM

El Problema: Ejecutar microservicios con Apache Flink o Kafka Streams requiere de 1GB a 4GB de RAM base por contenedor, inflando los costos cloud.

✨ La Solución: StreamFlow se ejecuta nativamente en Go usando menos de 50MB de RAM base por instancia, permitiéndole empaquetar docenas de servicios en hardware mínimo.

⚠️ Estado Externo de Alta Latencia

El Problema: Depender de Redis, DynamoDB o PostgreSQL para agregaciones con estado (promedios de ventana, contadores) introduce de 2ms a 10ms de latencia de red por evento.

✨ La Solución: PebbleDB embebido (almacén clave-valor LSM-tree) mantiene el estado localmente dentro del proceso para 0ms de latencia de red.

⚠️ Panics Silenciosos en Tiempo de Ejecución

El Problema: Los pipelines de streams de bytes crudos ([]byte) fallan silenciosamente o provocan panics en producción ante cambios de esquema inesperados.

✨ La Solución: El DSL de Generics con tipado estricto de StreamFlow garantiza que las incompatibilidades de esquema fallen en tiempo de compilación, no en producción a las 3 AM.

⚠️ Pruebas Frágiles e Inestables

El Problema: Probar timeouts de ventanas de sesión o joins deslizantes requiere configurar Docker, usar stubs/mocks o insertar instrucciones time.Sleep() frágiles.

✨ La Solución: El arnés de pruebas streamflow/testing ofrece un almacén en memoria determinista y simulación de viajes en el tiempo (AdvanceTime), permitiendo probar ventanas temporales de forma instantánea en pruebas unitarias estándar.

⚠️ Despliegues Políglotas Complejos

El Problema: Ejecutar modelos personalizados en Python o Rust dentro de su pipeline requiere desplegar microservicios separados o usar sidecars gRPC lentos e inseguros.

✨ La Solución: Cargue lógica políglota dinámicamente como plugins WASM aislados en sandbox impulsados por wazero con invocación en proceso a velocidad de microsegundos.

PROPUESTA DE VALOR

¿Por qué elegir StreamFlow?

🚀 Máxima Eficiencia de Costos

Reduzca sus costos de infraestructura de streaming hasta un 90%. Al reemplazar entornos pesados en Java por binarios estáticos compilados en Go, puede condensar más microservicios en máquinas virtuales de menor tamaño.

🚀 SLAs Sub-Milisegundo

Las pausas de recolección de basura (GC) en frameworks JVM provocan picos de latencia impredecibles. El bucle de procesamiento en Go de baja asignación de memoria garantiza latencias p99 predecibles por debajo del milisegundo bajo tráfico pico.

🚀 Velocidad de Desarrollo Excepcional

Construya pipelines más rápido mediante un DSL fluido en Go. Aproveche el autocompletado en IDE, la seguridad de tipos, soporte de refactorización y retroalimentación de compilación instantánea.

🚀 Simplicidad Operativa

Sin clústeres JVM, sin dependencias de ZooKeeper ni configuraciones complejas que mantener. Construya, compile y despliegue un único binario estático (< 30MB) en cualquier contenedor.

BENCHMARKS

Las Métricas Que Importan

Medido ejecutando pasos de pipeline en memoria sobre hardware cloud estándar.

1.4M+
Throughput / núcleo
msg/seg particionado por clave
< 0.7 μs
Latencia P50
Cero saltos de red
< 50 MB
RAM Base
Sin sobrecarga JVM
0 ms
Latencia de Estado
PebbleDB LSM embebido
CAPACIDADES

Características Principales

🔹 DSL en Go con Generics y Tipado Seguro

Sin Panics en Producción: Eleve flujos binarios a structs fuertemente tipadas en Go. Los cambios de esquema fallan en tiempo de compilación.
Autocompletado en IDE: Soporte completo para autocompletado y análisis estático en filtros, maps y ventanas.

🔹 Motor de Estado Embebido LSM (PebbleDB)

Sin Saltos a Bases de Datos Externas: Calcule agregaciones y promedios de sesión localmente sin consultas a Redis o Postgres.
Compactación Activa en Disco: Purga automática de estados expirados e índices de join en almacenamiento.

🔹 SQL Declarativo para Streams

Consultas SQL en Proceso: Ejecute SELECT category, amount FROM orders WHERE amount BETWEEN 100 AND 200 directamente sobre eventos en vivo.
Enrutamiento Dinámico de Fuentes: Configura automáticamente las suscripciones a brokers según la cláusula FROM.

🔹 Plugins WASM Aislados en Sandbox

Pipelines Políglotas: Ejecute lógica escrita en Rust, C++, Python o Go compilada a WebAssembly.
Pooling de Instancias WASM: Pre-instancia módulos para escalar con hilos paralelos sin contención de locks.

🔹 Semántica Effectively-Once (EOS)

Consistencia transaccional garantizada mediante upserts idempotentes en base de datos y commits de offset alineados.

🔹 DLQ Tolerante a Fallos

Etiquete mensajes con fallos incluyendo metadatos de error y enrútelos a una Dead Letter Queue (DLQ) sin interrumpir la ingestión.

PARADIGMAS DE CÓDIGO

Construya Pipelines a su Manera

Seleccione un paradigma para ver lo sencillo que es construir con StreamFlow.

[1] DSL Go Tipado [2] Streaming SQL [3] Operador WASM (Rust)
// Transformación de stream fuertemente tipada
clicks := runtime.From[ClickEvent](source, pipeline)
engaged := clicks.Filter("filter_bounces", func(c ClickEvent) bool {
    return c.DurationS > 2
})
engaged.To("engaged-clicks")
ARQUITECTURA INTERNA

Diseño de Ingeniería

Pools de Workers Particionados por Clave

Para escalar el throughput entre núcleos de CPU garantizando el orden de mensajes por clave, StreamFlow calcula hashes usando FNV-1a. Los mensajes con la misma clave se asignan a la misma goroutine worker, evitando bloqueos mutex y contención de CPU.

Pruebas de Viaje en el Tiempo Sin Mocks

StreamFlow incluye un arnés de pruebas dedicado (streamflow/testing). Escriba pruebas unitarias para ventanas de sesión o joins deslizantes y avance el tiempo instantáneamente (AdvanceTime) sin llamadas sleep inestables.

INICIO RÁPIDO

Comience en 60 Segundos

1. Instale el CLI
go install github.com/santiagolertora/streamflow/cmd/streamflow@latest
2. Configure su Pipeline (pipeline.toml)
[pipeline]
  [pipeline.source]
  type = "kafka"
  brokers = "localhost:9092"
  topic = "raw-events"
  group_id = "analytics-consumer"

  [pipeline.sink]
  type = "postgres"
  dsn = "postgres://user:pass@localhost:5432/db"
  table = "analytics_metrics"
3. Ejecute el Pipeline
streamflow run pipeline.toml

¿Listo para Diseñar Pipelines de Streaming a Velocidad Sub-Microsegundo?

Consulte el repositorio público de StreamFlow en GitHub o póngase en contacto con los ingenieros de datos de Binlogic para diseñar su arquitectura.

Explorar Repositorio en GitHub Hablar con Ingenieros de Binlogic