Skip to main content

流式数据简易处理包

Project description

flowdata

流式数据简易处理工具

本项目支持对流式数据的处理进行多进程加速

安装 flowdata

pip install flowdata

Quick Start

数据读取支持

  • 简单封装了txt, json, jsonl, excel文件的读写接口。请参考FileTool, JsonTool, JsonlTool, ExcelTool类。

1、单任务

代码中通过add_task将任务加载到任务流中,且可以根据需要指定不同进程数量

import time
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。任务执行按照task的添加顺序执行。前一个任务的输出是下一个任务的输入。

import time
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.5.tar.gz (9.4 kB view details)

Uploaded Source

File details

Details for the file flowdata-0.1.5.tar.gz.

File metadata

  • Download URL: flowdata-0.1.5.tar.gz
  • Upload date:
  • Size: 9.4 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.5.tar.gz
Algorithm Hash digest
SHA256 a1242a05f7dda59a9efc8e7affa4c538062ceabb6b1ebb910aaecb3c654459fc
MD5 67aa6168a816ec88b3ac388a95d541e9
BLAKE2b-256 9ba8ba1e231e7fd00aef4161c9d9a95d18e18cdd2cadeb6d5dd64a029deb3595

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