Skip to main content

Streamtask

Streamtask is a lightweight python parallel framework for parallelizing the computationally intensive pipelines. It is similar to Map/Reduce, while it is more lightweight. It parallelizes each module in the pipeline with a given processing number to make it possible to leverage the different speeds in different modules. It improves the performance especially there are some heavy I/O operations in the pipeline.

Example

Suppose we want to process the data in a pipline with 3 blocks, f1, f2 and f3. We can use the following code to parallelize the processing. We can also directly add a data list or give a file name by using add_data.

def f1(total):
    import time
    for i in range(total):
        time.sleep(0.002)
        yield i * 2

def f2(n, add, third = 0.01):
    time.sleep(0.02)
    return n + add + third

def f3_the_final(n):
    time.sleep(0.03)
    return n + 1

if __name__ == "__main__":
    total = 100000
    stk = StreamTask(total = total)
    stk.add_data(data = range(total), batch_size=10) # Also support directly stream reading file
    #stk.add_module(f1, 1, total = total)
    stk.add_module(f2, 2, args = [0.5], third = 0.02)
    stk.add_module(f3_the_final, 2)
    stk.run(parallel = True)
    stk.join()
    res = stk.get_results()
    print(stk.get_results())
stream_reader (1/1):  79%|██████████████████████████████████████████████▋            | 7923/10000 [00:05<00:01, 1567.68it/s]
f2 (2/2):  29%|████████████████████▉                                                   | 1927/6635 [00:04<00:09, 476.72it/s]
f3_the_final (2/2): 100%|█████████████████████████████████████████████████████████████▉| 1927/1928 [00:04<00:00, 476.55it/s]

Release files for streamtask 0.0.9

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for streamtask 0.0.9
File Size Uploaded
streamtask-0.0.9.tar.gz 4.8 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for streamtask 0.0.9
File Interpreter ABI Platform
streamtask-0.0.9-py3-none-any.whl Python 3 none any Details

Total release size: 10.0 kB

Release files / streamtask-0.0.9.tar.gz

Download URL streamtask-0.0.9.tar.gz
Size 4.8 kB
Tags Source
SHA-256 checksum
How to use checksums
490c1226e76f6c088dc4a2f07b88c9bf3838ea6f358cd3a44593b40b3fbd4cfd
BLAKE2b-256 checksum
How to use checksums
f5022d7baa4bbbd6bceab38e3fd00e552a8df2e332470ad25e11b11e9fc4e79b
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/4.0.2 CPython/3.8.13

Release files / streamtask-0.0.9-py3-none-any.whl

Download URL streamtask-0.0.9-py3-none-any.whl
Size 5.2 kB
Tags Python 3
SHA-256 checksum
How to use checksums
1eec316ecf80f7be4498fa2f28c5f49d53c73277adea12c7b23276d593e5849c
BLAKE2b-256 checksum
How to use checksums
150ef07689023dd51ef1bd0f8f16dd55137f8ced7a1bd76ef43e47d9cbbfce38
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/4.0.2 CPython/3.8.13

Release history Release notifications | RSS feed

This release

0.0.9 This release

2 release files

0.0.8

2 release files

0.0.7

2 release files

0.0.6

2 release files

0.0.5

2 release files

0.0.4

2 release files

0.0.3

2 release files

0.0.2

2 release files

0.0.1

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page