9 repositorios
Toolkits for composing asynchronous data pipelines using non-blocking operators and backpressure.
Distinct from Stream Combinators: The candidates are too narrow (combinators or specific traits); this represents the framework's primary identity.
Explore 9 awesome GitHub repositories matching programming languages & runtimes · Asynchronous Stream Processing Frameworks. Refine with filters or upvote what's useful.
Reactor Core es un kit de herramientas de programación reactiva y una base no bloqueante para componer pipelines de datos asíncronos en la JVM. Sirve como framework de procesamiento de flujos asíncronos y sistema de gestión de contrapresión (backpressure), permitiendo a los desarrolladores transformar, filtrar y combinar secuencias de eventos mientras regulan el flujo de datos entre productores y consumidores para evitar el agotamiento de recursos. La biblioteca se diferencia por un sofisticado sistema de planificación de concurrencia y control de flujo basado en la demanda. Desacopla el procesamiento de señales de hilos específicos utilizando un registro de planificadores y proporciona mecanismos para la propagación de metadatos inmutables conscientes del contexto a través de límites asíncronos. También cuenta con herramientas especializadas para la captura de trazas en tiempo de ensamblaje y planificación de tiempo virtual para facilitar la prueba de operadores basados en el tiempo. El proyecto cubre una amplia gama de capacidades, incluyendo procesamiento funcional de datos para agregación y ventanas de secuencias, una variedad de estrategias de recuperación de errores como reintentos con retroceso exponencial y utilidades para conectar API de callback heredadas o síncronas en flujos reactivos. Además, proporciona instrumentación para el monitoreo de pipelines y un conjunto de herramientas de prueba para verificar secuencias de señales.
Acts as a comprehensive non-blocking foundation for composing asynchronous data pipelines on the JVM.
RxGo es una biblioteca de programación reactiva funcional y una implementación de ReactiveX para el lenguaje Go. Sirve como un kit de herramientas de procesamiento de flujos asíncronos diseñado para coordinar programas basados en eventos y flujos de datos utilizando el patrón observable. La biblioteca permite la construcción de pipelines de procesamiento asíncrono que transforman, filtran y combinan secuencias de eventos. Se distingue por el uso de operadores funcionales para componer estos pipelines y proporciona mecanismos para gestionar la ejecución concurrente. El kit de herramientas cubre una amplia gama de capacidades de orquestación de flujos, incluyendo agregación de datos, combinación de múltiples flujos y la conversión de flujos en estructuras de datos estáticas. Incluye soporte integrado para recuperación de errores, control de contrapresión (backpressure) para regular las velocidades de producción de datos y agrupación de workers para paralelizar el procesamiento a través de núcleos de CPU.
Provides a comprehensive framework for composing asynchronous data pipelines with built-in backpressure and non-blocking operators.
RxPY es una librería de programación reactiva funcional y una librería de observables ReactiveX para Python. Funciona como un procesador de flujos asíncronos y un framework de coordinación basado en eventos, utilizado para construir pipelines de datos que reaccionan a cambios de estado o flujos de eventos a lo largo del tiempo. La librería proporciona un kit de herramientas para componer programas asíncronos y basados en eventos mediante secuencias observables y operadores. Se distingue por el uso de planificadores (schedulers) configurables para gestionar la concurrencia, el timing y los ciclos de vida de las suscripciones. El proyecto cubre una amplia gama de capacidades de procesamiento de flujos, incluyendo agregación, filtrado y combinación de datos. Proporciona mecanismos para la difusión de eventos, almacenamiento en búfer de secuencias y gestión de errores, así como herramientas para coordinar flujos observables con bucles de eventos asíncronos. Las pruebas y el aseguramiento de la calidad se apoyan en la simulación de tiempo virtual, el modelado con diagramas de mármol y la verificación de emisiones.
Serves as a comprehensive framework for composing asynchronous data pipelines using non-blocking operators.
Este proyecto proporciona una especificación formal y un conjunto de interfaces estándar de Java para el procesamiento de flujos asíncronos. Define un protocolo estandarizado para pasar secuencias de elementos entre editores (publishers) y suscriptores a través de diferentes hilos, centrándose en una especificación de flujos reactivos para la JVM. El proyecto se centra en la interoperabilidad al proporcionar una API común que permite que diferentes bibliotecas de streaming asíncrono trabajen juntas. Esto se logra mediante un conjunto estándar de interfaces y mecanismos de puente que traducen entre especificaciones de streaming incompatibles. La especificación cubre un protocolo de contrapresión (backpressure) sin bloqueo para regular el flujo de datos y evitar la sobrecarga del sistema, requiriendo que los suscriptores señalen la demanda. También define el ciclo de vida de los flujos, incluyendo la gestión de suscripciones, el procesamiento de elementos y la terminación basada en señales para la limpieza de recursos. El proyecto incluye un framework para verificar el comportamiento del flujo y validar la lógica de procesamiento frente a las reglas de contrapresión y eventos asíncronos.
Provides the foundational toolkit for composing asynchronous data pipelines using non-blocking operators and backpressure.
Sofa-rpc es un framework de llamada a procedimiento remoto de alto rendimiento diseñado para construir aplicaciones Java distribuidas. Funciona como un kit de herramientas para gestionar un service mesh distribuido, proporcionando una capa de comunicación gRPC y un sistema para registrar y localizar instancias de servicios remotos. El framework cuenta con una capa de seguridad de red que implementa cifrado TLS y comprobaciones de autorización para proteger los datos transmitidos entre servicios. Utiliza una capa de protocolo conectable para admitir múltiples estándares de comunicación, asegurando una conectividad punto a punto flexible. La fiabilidad y la gestión del tráfico se manejan a través de disyuntores, balanceo de carga del lado del cliente y monitoreo de la salud del servicio. El sistema también incluye herramientas de observabilidad para el rastreo de solicitudes distribuidas y streaming remoto reactivo para aumentar la eficiencia de los recursos. El framework proporciona utilidades para la serialización de datos JSON y gestiona la conectividad remota a través de la agrupación de conexiones y un registro de descubrimiento de servicios.
Supports remote invocations using asynchronous streams to increase throughput and resource efficiency.
Whisper streaming is an automated speech recognition engine designed to convert live audio into text. It functions as a network-based transcription server that accepts raw audio data from remote clients and returns incremental text results in real-time. The system distinguishes itself through its ability to process audio streams incrementally, allowing for immediate transcription and translation as speech is captured. It incorporates voice activity detection to isolate human speech from background noise and utilizes sliding-window buffering to manage incoming audio segments, ensuring that pro
Decouples audio ingestion from transcription tasks using non-blocking queues to ensure continuous data flow.
Js-csp is a concurrency library that implements communicating sequential processes with channels and generator-based routines for asynchronous programming in JavaScript. It provides a framework for building data processing pipelines and managing concurrent workflows, bringing Go-style channels and process coordination primitives to applications. The library coordinates message passing through buffered channels, stream mixing, multi-channel selection, and pub-sub broadcasting. Communication pathways support unbuffered rendezvous, fixed buffers, sliding windows, and dropping overflow. Operation
Provides an asynchronous framework for building data processing pipelines and managing concurrent workflows.
Level is a database library that provides a unified interface for managing sorted key-value data. It functions as an abstraction layer that allows applications to store and retrieve binary information consistently across server-side environments and web browsers. The project utilizes a modular architecture that supports pluggable storage backends, enabling the system to adapt to different host environments while maintaining identical behavior. By organizing data in lexicographical order, it facilitates efficient range queries and ordered retrieval. The library handles large datasets through a
Implements asynchronous stream processing to handle large datasets without blocking the main execution thread.
Este proyecto sirve como un recurso educativo completo para aprender programación paralela y computación de alto rendimiento utilizando unidades de procesamiento gráfico (GPU). Proporciona guía técnica sobre los paradigmas fundamentales requeridos para delegar tareas computacionalmente intensivas desde un sistema host a aceleradores de hardware especializados. Los materiales cubren las metodologías centrales para gestionar operaciones de datos paralelos, incluyendo la orquestación de memoria entre espacios de host y dispositivo y la organización de hilos en rejillas y bloques estructurados. Detalla los modelos de ejecución necesarios para distribuir cargas de trabajo a través de múltiples núcleos de procesamiento, permitiendo a los desarrolladores escalar aplicaciones pesadas en datos de manera efectiva. Más allá de la implementación básica, el recurso incluye prácticas de diagnóstico para analizar métricas de ejecución e identificar cuellos de botella de rendimiento. Ofrece estrategias para optimizar la ejecución de kernels y depurar errores lógicos dentro de bases de código concurrentes para asegurar el máximo rendimiento y eficiencia en entornos de computación acelerada.
Overlaps data transfers and kernel execution using non-blocking queues to maximize hardware utilization.