Skip to main content

流式数据简易处理包

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)

Uploaded Source

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

Hashes for flowdata-0.1.1.tar.gz
Algorithm Hash digest
SHA256 29aeb78f9ae0ed20196df380281f483f0237e501e9ddf033acecf8b984dc56dd
MD5 bf17db6eed4784f3d9d93e238c28481f
BLAKE2b-256 b143033a1c4c4b67e58190b1ac78ff328843085583b41167a053f2d038726ef6

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