A minimal, `asyncio`-native event sourcing library
Project description
py-event-sourcing - A minimal, asyncio-native event sourcing library
This library provides core components for building event-sourced systems in Python. It uses SQLite for persistence and offers a simple API for write, read, and watch operations on event streams.
For a deeper dive into the concepts and design, please see:
- Concepts: An introduction to the Event Sourcing pattern.
- Design: An overview of the library's architecture and components.
Key Features
- Simple & Serverless: Uses SQLite for a file-based, zero-dependency persistence layer.
- Global Event Stream: Query all events across all streams in their global sequence using the special
'@all'stream, perfect for building cross-stream read models or for auditing. - Idempotent Writes: Prevents duplicate events in distributed systems by using an optional
idon each event. - Optimistic Concurrency Control: Ensures data integrity by allowing writes only against a specific, expected stream version.
- Efficient Watching: A centralized notifier polls the database once to serve all watchers, avoiding the "thundering herd" problem and ensuring low-latency updates.
- Snapshot Support: Accelerate state reconstruction for long-lived streams by saving and loading state snapshots. This is also supported for the
'@all'stream, which is ideal for complex, system-wide projections. - Fully Async API: Built from the ground up with
asynciofor high-performance, non-blocking I/O. - Extensible by Design: Core logic is decoupled from storage implementation via
Protocol-based adapters. While this package provides a highly-optimized SQLite backend, you can easily create your own adapters for other databases (e.g., PostgreSQL, Firestore).
Installation
This project uses uv for dependency management.
-
Create and activate the virtual environment:
uv venv source .venv/bin/activate
-
Install the package in editable mode with dev dependencies:
uv pip install -e ".[dev]"
Quick Start
Here’s a quick example of writing to and reading from a stream.
import asyncio
import os
import tempfile
from pysource import sqlite_stream_factory, CandidateEvent
async def main():
# Use a temporary file for the database to keep the example self-contained.
with tempfile.TemporaryDirectory() as tmpdir:
db_path = os.path.join(tmpdir, "example.db")
# The factory is an async context manager that handles all resources.
async with sqlite_stream_factory(db_path) as open_stream:
stream_id = "my_first_stream"
# Write an event
async with open_stream(stream_id) as stream:
event = CandidateEvent(type="UserRegistered", data=b'{"user": "Alice"}')
await stream.write([event])
print(f"Event written. Stream version is now {stream.version}.")
# Read the event back
async with open_stream(stream_id) as stream:
all_events = [e async for e in stream.read()]
print(f"Read {len(all_events)} event(s) from the stream.")
print(f" -> Event type: {all_events[0].type}, Data: {all_events[0].data.decode()}, Version: {all_events[0].version}")
if __name__ == "__main__":
asyncio.run(main())
Basic Usage
For a detailed, fully-commented example covering all major features (writing, reading, snapshots, and watching), please see basic_usage.py.
To run the example:
uv run python3 basic_usage.py
Sample Output:
--- Example 1: Writing and Reading Events ---
Stream 'counter_stream_1' is now at version 3.
Reading all events from the stream:
- Event version 1: Increment
- Event version 2: Increment
- Event version 3: Decrement
--- Example 2: Reconstructing State from Events (Read Model & Projector) ---
Reconstructed state: Counter is 1.
--- Example 3: Watching for New Events ---
Watching for events (historical and new)...
- Watched event: Increment (version 1)
- Watched event: Increment (version 2)
- Watched event: Decrement (version 3)
Writing new events to trigger the watcher...
- Watched event: Increment (version 4)
- Watched event: Increment (version 5)
--- Example 4: Using Snapshots for Efficiency ---
Snapshot for 'counter' projection saved at version 5 with state: count = 3
State for 'counter' projection restored from snapshot at version 5. Count is 3.
Replaying events since snapshot...
- Applying event version 6: Increment
Final reconstructed state: Counter is 4.
Benchmarks
A benchmark script is included to measure write/read throughput and demonstrate the performance benefits of using snapshots. You can find the script at benchmark.py.
To run the benchmark:
uv run python3 benchmark.py
Sample Output:
--- Running benchmark with 1,000,000 events ---
Writing 1,000,000 events...
Finished writing 1,000,000 events in 9.41 seconds.
Write throughput: 106,320.88 events/sec.
--- Benchmark 1: Reconstructing state from all events ---
Reconstructed state (all events): Counter is 1,000,000.
Time to reconstruct from all events: 2.66 seconds.
Read throughput: 375,273.94 events/sec.
--- Benchmark 2: Reconstructing state using snapshot ---
Creating snapshot...
Snapshot created.
Writing 100 additional events...
Finished writing 100 additional events.
State restored from snapshot at version 1,000,000.
Reconstructed state (with snapshot): Counter is 1,000,100.
Time to reconstruct with snapshot: 0.00 seconds.
--- Benchmark Summary ---
Time to reconstruct from all 1,000,100 events: 2.66 seconds.
Time to reconstruct with snapshot (after 1,000,000 events): 0.0007 seconds.
Database file size: 148.13 MB
Testing
The test suite uses pytest. To run all tests, use the following command:
uv run pytest
Contributing
We welcome contributions! Please see our Contributing Guide for details on:
- Setting up a development environment
- Code style guidelines
- Testing requirements
- Pull request process
License
This project is licensed under the MIT License - see the LICENSE file for details.
Changelog
See CHANGELOG.md for a list of changes and version history.
Project details
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file py_event_sourcing-1.0.0.tar.gz.
File metadata
- Download URL: py_event_sourcing-1.0.0.tar.gz
- Upload date:
- Size: 85.1 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.11.9
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
0dc3208cf746f21a05f202e9168aa7d3efefa5482cdc83178db06418ac0e7680
|
|
| MD5 |
751bfa98d3b3994e7ff50f3e2da782cd
|
|
| BLAKE2b-256 |
02528a1330777cbc097e866abb07a18e424f8e99acb743cd9667988e03930b0d
|
File details
Details for the file py_event_sourcing-1.0.0-py3-none-any.whl.
File metadata
- Download URL: py_event_sourcing-1.0.0-py3-none-any.whl
- Upload date:
- Size: 20.4 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.11.9
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
9302349fa1b5af63c500a1c660e2b749160c07646d13f5b27042adeb10c617e7
|
|
| MD5 |
08a56678a0c4f9be30c017cc787f5532
|
|
| BLAKE2b-256 |
ce4af2799166716140217aaf155ae2d05dc7dcb019e768e1ba64b9d747917ed4
|