Skip to main content

Zero-copy shared memory IPC library for building complex streaming data pipelines capable of processing large datasets

Project description

MomentumX


MomentumX is a zero-copy shared memory IPC library for building complex streaming data pipelines capable of processing large datasets using Python.


Key Features:

  • High-Throughput, Low Latency
  • Supports streaming and synchronous modes for use within a wide variety of use cases.
  • Bring your own encoding, or use raw binary data.
  • Sane data protections to ensure reliability of data in a cooperative computing environment.
  • Pairs with other high-performance libraries, such as numpy and scipy, to support parallel processing of memory-intensive scientific data.
  • Works on most modern versions of Linux using shared memory (via /dev/shm).
  • Seamlessly integrates into a Docker environment with minimal configuration, and readily enables lightweight container-to-container data sharing.

Examples:

Below are some simplified use cases for common MomentumX workflows. Consult the examples in the examples/ directory for additional details and implementation guidance.

Stream Mode

# Producer Process
import momentumx as mx

# Create a stream with a total capacity of 10MB (1MB x 10)
stream = mx.Producer('my_stream', buffer_size=int(1e6), buffer_count=10, sync=False)

# Obtain the next available buffer for writing
buffer = stream.next_to_send()
buffer.write(b'1') 

buffer.send()
# NOTE: buffer.send() can also be passed an explicit number of bytes as well. 
# Otherwise an internally managed cursor will be used.
# Consumer Process(es)
import momentumx as mx

stream = mx.Consumer('my_stream')

# Receive from my_stream as long as the stream has not ended OR there are unread buffers 
while stream.has_next:

    # Block while waiting to receive buffer 
    # NOTE: Non-blocking receive is possible using blocking=False keyword argument
    buffer = stream.receive()
    
    # If we are here, either the stream ended OR we have a buffer, so check...
    if buffer is not None:

        # We have buffer containing data, so print the entire contents
        print(buffer.read(buffer.data_size))
    
        # See also "Implicit versus Explicit Buffer Release" section below.

Sync Mode

# Producer Process
import momentumx as mx
import threading
import signal

cancel_event = threading.Event()
signal.signal(signal.SIGINT, (lambda _sig, _frm: cancel_event.set()))

# Create a stream with a total capacity of 10MB
stream = mx.Producer(
    'my_stream', 
    buffer_size=int(1e6), 
    buffer_count=10, 
    sync=True
) # NOTE: sync set to True

min_subscribers = 1

while stream.subscriber_count < min_subscribers:
    print("waiting for subscriber(s)")
    if cancel_event.wait(0.5):
        break

print("All expected subscribers are ready")

# Write the series 0-999 to a consumer 
for n in range(0, 1000):
    if stream.subscriber_count == 0:
        cancel_event.wait(0.5)

    # Note: sending strings directly is possible via the send_string call
    elif stream.send_string(str(n)):
        print(f"Sent: {n}")
# Consumer Process(es)
import momentumx as mx

stream = mx.Consumer('my_stream')

while stream.has_next:
    data = stream.receive_string() 

    if data is not None: 
        # Note: receiving strings is possible as well via the receive_string call
        print(f"Received: {data}")

Iterator Syntax

Working with buffers is even easier using iter() builtin:

import momentumx as mx

stream = mx.Consumer(STREAM)

# Iterate over buffers in the stream until stream.receive() returns None
for buffer in iter(stream.receive, None):     
    # Now, buffer is guaranteed to be valid, so no check required -  
    # go ahead and print all the contents again, this time using 
    # the index and slice operators!
    print(buffer[0])                    # print first byte
    print(buffer[1:buffer.data_size])   # print remaining bytes

Numpy Integration

import momentumx as mx
import numpy as np

# Create a stream
stream = mx.Consumer('numpy_stream')

# Receive the next buffer (or if a producer, obtain the next_to_send buffer)
buffer = stream.receive()

# Create a numpy array directly from the memory without any copying
np_buff = np.frombuffer(buffer, dtype=uint8)

Implicit versus Explicit Buffer Release

MomentumX Consumers will, by default, automatically release a buffer under the covers once all references are destroyed. This promotes both usability and data integrity. However, there may be cases where the developer wants to utilize a different strategy and explicity control when buffers are released to the pool of available buffers.

stream = mx.Consumer('my_stream')

buffer = stream.receive()

# Access to buffer is safe!
buffer.read(10)

# Buffer is being returned back to available buffer pool. 
# Be sure you are truly done with your data!
buffer.release() 

# DANGER: DO NOT DO THIS! 
# All operations on a buffer after calling `release` are considered unsafe! 
# All safeguards have been removed and the memory is volatile!
buffer.read(10) 

Isolated Contexts

MomentumX allows for the usage of streams outside of /dev/shm (the default location). Pass the context kwarg pointing to a directory on the filesystem for both the Producer and all Consumer instances to create isolated contexts.

This option is useful if access to /dev/shm is unsuitable.

import momentumx as mx

# Create a producer attached to the context path /my/path
stream = mx.Producer('my_stream', ..., context='/my/path/')
...

# Create Consumer elsewhere attached to the same context of /my/path
stream = mx.Consumer('my_stream', context='/my/path/')

License

Captivation Software, LLC offers MomentumX under an Unlimited Use License to the United States Government, with all other parties subject to the GPL-3.0 License.

Inquiries / Requests

All inquiries and requests may be sent to opensource@captivation.us.

Copyright © 2022-2023 - Captivation Software, LLC.

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

momentumx-2.9.1.tar.gz (60.3 kB view details)

Uploaded Source

Built Distributions

If you're not sure about the file name format, learn more about wheel file names.

momentumx-2.9.1-cp315-cp315-manylinux2014_x86_64.manylinux_2_17_x86_64.whl (219.8 kB view details)

Uploaded CPython 3.15manylinux: glibc 2.17+ x86-64

momentumx-2.9.1-cp314-cp314-manylinux2014_x86_64.manylinux_2_17_x86_64.whl (219.8 kB view details)

Uploaded CPython 3.14manylinux: glibc 2.17+ x86-64

momentumx-2.9.1-cp313-cp313-manylinux2014_x86_64.manylinux_2_17_x86_64.whl (219.8 kB view details)

Uploaded CPython 3.13manylinux: glibc 2.17+ x86-64

momentumx-2.9.1-cp312-cp312-manylinux2014_x86_64.manylinux_2_17_x86_64.whl (219.8 kB view details)

Uploaded CPython 3.12manylinux: glibc 2.17+ x86-64

momentumx-2.9.1-cp311-cp311-manylinux2014_x86_64.manylinux_2_17_x86_64.whl (220.3 kB view details)

Uploaded CPython 3.11manylinux: glibc 2.17+ x86-64

momentumx-2.9.1-cp310-cp310-manylinux2014_x86_64.manylinux_2_17_x86_64.whl (220.2 kB view details)

Uploaded CPython 3.10manylinux: glibc 2.17+ x86-64

momentumx-2.9.1-cp39-cp39-manylinux2014_x86_64.manylinux_2_17_x86_64.whl (220.6 kB view details)

Uploaded CPython 3.9manylinux: glibc 2.17+ x86-64

File details

Details for the file momentumx-2.9.1.tar.gz.

File metadata

  • Download URL: momentumx-2.9.1.tar.gz
  • Upload date:
  • Size: 60.3 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.10.9

File hashes

Hashes for momentumx-2.9.1.tar.gz
Algorithm Hash digest
SHA256 13e2fa6f0e43a5996e9f688604900118818ffbff1bda2662b365aa80143ace6a
MD5 f0ef429416c0630db10ae94aeae22da4
BLAKE2b-256 b92c9afa3f3eb05a525a1701c9b6c985085136e4c2700e9b8b63579ef12862de

See more details on using hashes here.

File details

Details for the file momentumx-2.9.1-cp315-cp315-manylinux2014_x86_64.manylinux_2_17_x86_64.whl.

File metadata

File hashes

Hashes for momentumx-2.9.1-cp315-cp315-manylinux2014_x86_64.manylinux_2_17_x86_64.whl
Algorithm Hash digest
SHA256 13f762d80f9cfa128cc717f4a3f99a4bd33f4fc1db883942c71ac16cee1c0f46
MD5 8c31296fd59f2506c8aeb01c51d3cec4
BLAKE2b-256 6fa8ea74dd96117fcd6495a645f750e424df94c5abb250b31c6b1ad1d5bd1915

See more details on using hashes here.

File details

Details for the file momentumx-2.9.1-cp314-cp314-manylinux2014_x86_64.manylinux_2_17_x86_64.whl.

File metadata

File hashes

Hashes for momentumx-2.9.1-cp314-cp314-manylinux2014_x86_64.manylinux_2_17_x86_64.whl
Algorithm Hash digest
SHA256 ea48f283d4480ea2f08992a005b27a462ae2a2a830f5c0d3e4176cfd5314146d
MD5 2ce8af61661d4825f130420e2c86b9ff
BLAKE2b-256 209fd7506ca5f6c8200fa8720132cd799999e5832207867a10db026eefe535f7

See more details on using hashes here.

File details

Details for the file momentumx-2.9.1-cp313-cp313-manylinux2014_x86_64.manylinux_2_17_x86_64.whl.

File metadata

File hashes

Hashes for momentumx-2.9.1-cp313-cp313-manylinux2014_x86_64.manylinux_2_17_x86_64.whl
Algorithm Hash digest
SHA256 dc456f30a1e6f6102caa4187c390cadb0e1b419930a28499709799611d3870e2
MD5 8cd89969b58ca75e051815f5d41c2c96
BLAKE2b-256 ab4e7ddabda2abda1debef103c080aa3e8eb464757ce8c43e349f9564314d7af

See more details on using hashes here.

File details

Details for the file momentumx-2.9.1-cp312-cp312-manylinux2014_x86_64.manylinux_2_17_x86_64.whl.

File metadata

File hashes

Hashes for momentumx-2.9.1-cp312-cp312-manylinux2014_x86_64.manylinux_2_17_x86_64.whl
Algorithm Hash digest
SHA256 30567afb3c464a7f67abca2f5bbb124f7f10958e5d4437911f7c8227bf76b625
MD5 0e7f2faccd53296ed93fd5c7919d681f
BLAKE2b-256 b7d9420a9b0d6033139aa7c936ac0addaf3c0b60725516f3c319ffda2fa28ac9

See more details on using hashes here.

File details

Details for the file momentumx-2.9.1-cp311-cp311-manylinux2014_x86_64.manylinux_2_17_x86_64.whl.

File metadata

File hashes

Hashes for momentumx-2.9.1-cp311-cp311-manylinux2014_x86_64.manylinux_2_17_x86_64.whl
Algorithm Hash digest
SHA256 097c8f656e2f96d5c4fa83483314383202b10d46cbac9f2b93a68c1b809333d1
MD5 2b0953879a6ef857128460e8b189d439
BLAKE2b-256 d6876051520952b0079f453ab91a027b287883348b96b8e6f18cfcbd90725cb8

See more details on using hashes here.

File details

Details for the file momentumx-2.9.1-cp310-cp310-manylinux2014_x86_64.manylinux_2_17_x86_64.whl.

File metadata

File hashes

Hashes for momentumx-2.9.1-cp310-cp310-manylinux2014_x86_64.manylinux_2_17_x86_64.whl
Algorithm Hash digest
SHA256 dd839d4dd20daa182dbf2ec244d86be821d65bcd46e4b0cbea47602cfdeb6a86
MD5 5db2c2238b5e37bf179b17d742c95e05
BLAKE2b-256 dce904b935881b0bc3d9de87a91cc43885281df288dfe67a7fa48ff85156421c

See more details on using hashes here.

File details

Details for the file momentumx-2.9.1-cp39-cp39-manylinux2014_x86_64.manylinux_2_17_x86_64.whl.

File metadata

File hashes

Hashes for momentumx-2.9.1-cp39-cp39-manylinux2014_x86_64.manylinux_2_17_x86_64.whl
Algorithm Hash digest
SHA256 68b163a88e5c8c34acfd8d9c2f8a14295a622a43c6af6ad1f99ab1c9dd7f1812
MD5 317a2ca26647230ceb9f68541b10d4d0
BLAKE2b-256 6237a2f75ca6a2cec54cd5d290c95d30596f4a81fa204cd89a2333c9531d2a2e

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page