流式数据简易处理包
Project description
flowdata
流式数据简易处理工具
本项目支持对流式数据的处理进行多进程加速
安装 flowdata
pip install flowdata
Quick Start
1、单任务
代码中通过add_task将任务加载到任务流中,且可以根据需要指定不同进程数量
from flowdata import FlowBase
from flowdata import add_task
class TaskFlow(FlowBase):
@add_task(work_num=1)
def task(self, item: dict, *args, **kwargs) -> dict:
time.sleep(.2)
item['id'] += 1
return item
def get_data(self):
for i in range(20):
yield {"id": i}
def save_data(self, item_iter):
list(item_iter)
TaskFlow().main()
2、多任务
假设一个处理数据的任务可以细分为多个子任务,例如,task_a, task_b。
from flowdata import FlowBase
from flowdata import add_task
# 多个任务
class TaskFlow(FlowBase):
@add_task(work_num=2)
def task_a(self, item: dict, *args, **kwargs) -> dict:
time.sleep(.2)
item['id'] += 1
return item
@add_task(work_num=2)
def task_b(self, item: dict, *args, **kwargs) -> dict:
time.sleep(.2)
item['id'] += 1
return item
def get_data(self):
for i in range(20):
yield {"id": i}
def save_data(self, item_iter):
list(item_iter)
TaskFlow().main()
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
flowdata-0.1.1.tar.gz
(7.6 kB
view details)
File details
Details for the file flowdata-0.1.1.tar.gz.
File metadata
- Download URL: flowdata-0.1.1.tar.gz
- Upload date:
- Size: 7.6 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/4.0.2 CPython/3.7.10
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
29aeb78f9ae0ed20196df380281f483f0237e501e9ddf033acecf8b984dc56dd
|
|
| MD5 |
bf17db6eed4784f3d9d93e238c28481f
|
|
| BLAKE2b-256 |
b143033a1c4c4b67e58190b1ac78ff328843085583b41167a053f2d038726ef6
|