6 dépôts
High-throughput data processing architectures that handle continuous streams of data using concurrent workers to manage memory and maintain responsiveness.
Distinct from Multi-Threaded Packet Processing: Candidates focus on network packets, educational models, or UI thread messaging, whereas this is about high-volume data ingestion pipelines for databases.
Explore 6 awesome GitHub repositories matching data & databases · Stream Processing Pipelines. Refine with filters or upvote what's useful.
Apache Storm is a distributed stream processing framework and real-time data processing engine. It functions as a fault-tolerant distributed computing system designed to analyze data in motion across a cluster of machines for continuous stream computation. The system enables the creation of fault-tolerant data pipelines and scalable event processing by distributing workloads across a network of computing nodes. This architecture ensures low latency and high throughput for live data while allowing the system to recover automatically from individual node failures. The framework provides capabi
Passes discrete data records through asynchronous message streams using a high-throughput pipeline architecture.
Hazelcast is a distributed data platform that combines an in-memory data grid with a stream processing engine to support real-time analytics and event-driven applications. It functions as a partitioned, distributed key-value store that replicates data across cluster nodes to provide low-latency access and high availability. The platform also serves as a distributed SQL query engine, allowing users to execute standard SQL statements against both in-memory datasets and external data sources. What distinguishes Hazelcast is its use of a distributed consensus subsystem to maintain strongly consis
Ships a high-throughput stream processing engine for building real-time, event-driven data pipelines.
Lazy.js is a JavaScript library that implements a lazy evaluation model for processing collections and data streams. It defers all computation until iteration begins, building chains of transformations that execute only when values are consumed, avoiding intermediate arrays and buffering. The library wraps data sources into a uniform sequence interface, enabling operations like map and filter to be chained together without materializing intermediate results. The library extends lazy processing beyond simple collections to handle asynchronous data sources, DOM events, strings, and Node.js stre
Drives computation by pulling values on demand from the sequence, avoiding buffering and intermediate storage.
RxGo est une bibliothèque de programmation réactive fonctionnelle et une implémentation de ReactiveX pour le langage Go. Elle sert de boîte à outils de traitement de flux asynchrone conçue pour coordonner les programmes basés sur les événements et les flux de données en utilisant le pattern observable. La bibliothèque permet la construction de pipelines de traitement asynchrones qui transforment, filtrent et combinent des séquences d'événements. Elle se distingue par l'utilisation d'opérateurs fonctionnels pour composer ces pipelines et fournit des mécanismes pour gérer l'exécution concurrente. La boîte à outils couvre un large éventail de capacités d'orchestration de flux, incluant l'agrégation de données, la combinaison de flux multiples et la conversion de flux en structures de données statiques. Elle inclut un support intégré pour la récupération d'erreurs, le contrôle de contre-pression (backpressure) pour réguler les vitesses de production de données, et le pooling de travailleurs pour paralléliser le traitement sur les cœurs CPU.
Implements high-throughput data processing pipelines that handle continuous streams using concurrent worker pools.
ZIO is a functional effect system for the JVM that models asynchronous and concurrent programs as pure, composable values with typed error handling and dependency injection. Its core identity is built on fiber-based concurrency, where lightweight, non-blocking fibers execute millions of concurrent tasks with structured lifecycle management, and a dual-channel error model that separates expected business failures from unexpected system defects at compile time. The system provides effect-typed dependency injection through a layer-based dependency graph, pull-based reactive stream processing with
Streams emit elements on demand with integrated backpressure, using a pull model where consumers control the flow and producers respond to downstream demand.
Maxwell est un outil de capture de données de changement (CDC) MySQL et une application de streaming de binlog qui convertit les modifications de base de données en événements JSON structurés. Il fonctionne comme un pipeline de données qui lit les logs binaires MySQL pour synchroniser les changements à travers des index externes, des moteurs de recherche et des systèmes de messagerie distribués tels que Kafka. Le projet fournit des capacités pour maintenir des pistes d'audit persistantes en enregistrant un historique chronologique de toutes les modifications de base de données. Il permet la synchronisation des données en temps réel et l'intégration d'architecture pilotée par événements en diffusant les changements de base de données vers des plateformes externes pour déclencher des flux de travail et notifier des microservices. Le système couvre de larges domaines fonctionnels incluant l'amorçage de données via des instantanés initiaux, la gestion de version de schéma et le filtrage d'événements. Il intègre la gestion du trafic via le routage par clé de partition et fournit une surveillance via des vérifications de santé et des métriques de performance exposées via un point de terminaison HTTP. Les connexions aux bases de données et aux producteurs de streaming sont sécurisées en utilisant SSL et une communication chiffrée.
Implements a high-throughput data pipeline that feeds database changes into streaming platforms with batching and partitioning.