Skip to main content

Framework for building modular, event-driven data pipelines

Project description

FinDrum

FinDrum is a lightweight Python framework for building and orchestrating data pipelines with extensible architecture via operators, datasources, schedulers, and triggers.

This repository (FinDrum-Platform) is the core package and is meant to be used as a library. Custom logic (pipelines and extensions) should be defined in external projects.


Installation

pip install findrum-platform

Overview

Findrum pipelines are defined in YAML files and can include:

  • A sequence of operators
  • A datasource (provides data from external source)
  • A scheduler (to run periodically)
  • An event trigger (to respond to real-time events)

Example structures:

Example with scheduler and batch datasource:

scheduler:
  type: MyCustomScheduler

pipeline:
  - id: batch_ingest
    datasource: MyDataSource
    params:
      key: value

  - id: transform
    operator: MyOperator
    depends_on: batch_ingest
    params:
      key: value

Example with event and pipeline

event:
  type: MyTrigger
  config:
    key: value

pipeline:
  - id: step1
    operator: MyOperator
    depends_on: MyTrigger
    params:
      key: value

  - id: step2
    operator: DownstreamOperator
    depends_on: step1

Interfaces

Findrum provides a minimal interface for each pipeline component. These are abstract base classes that must be subclassed by your custom logic.

Operator – Core processing unit

from findrum.interfaces import Operator

class MyOperator(Operator):
    def run(self, input_data):
        ...

Use when defining a step in a pipeline. Must implement run(input_data). We recommend that it returns a pandas.DataFrame.


DataSource – Step that starts a pipeline

from findrum.interfaces import DataSource

class MySource(DataSource):
    def fetch(self, **kwargs):
        ...

We recommend that it returns a pandas.DataFrame. It feeds the pipeline with data.


Scheduler – Periodic trigger for pipelines

from findrum.interfaces import Scheduler

class MyScheduler(Scheduler):
    def register(self, scheduler):
        # e.g., add job to APScheduler instance
        ...

Implements logic to execute the pipeline on a time interval or schedule.


EventTrigger – React to system/file/bucket events

from findrum.interfaces import EventTrigger

class MyTrigger(EventTrigger):
    def start(self):
        # Starts a file watcher, webhook listener, etc.
        ...

Runs the pipeline when an external event occurs (e.g., new Kafka message, file in MinIO). The trigger should call self.emit(data) to push input into the pipeline.


Core Classes

You can import and use the main classes provided by Findrum:

from findrum import Platform
  • Platform: Main entrypoint to manage pipelines, register them, and run based on schedule or events.

CLI Usage: findrum-run

After installing findrum-platform, a CLI tool is available:

Run a pipeline immediately

findrum-run pipelines/my_pipeline.yaml

Use a custom config file for extensions

findrum-run pipelines/my_pipeline.yaml --config config/config.yaml

Enable logging (INFO level)

findrum-run pipelines/my_pipeline.yaml --verbose

Extension Discovery

Findrum requires a config.yaml file with registered class paths:

operators:
  - my_project.operators.MyCustomOperator

datasources:
  - my_project.datasources.MyDataSource

schedulers:
  - my_project.schedulers.MyScheduler

triggers:
  - my_project.triggers.MyTrigger

This lets Findrum dynamically import your components.


Minimal Example For a Non-CLI runner

from findrum import Platform

platform = Platform("config.yaml")
platform.register_pipeline("pipelines/my_pipeline.yaml")
platform.start()

You can also run your pipelines from a python file (like main.py for example) following the example above.


Clean Project Structure

A typical project using Findrum should look like:

your-project/
├── operators/
│   └── my_operator.py
├── schedulers/
│   └── my_scheduler.py
├── triggers/
│   └── my_trigger.py
├── datasources/
│   └── my_datasource.py
├── pipelines/
│   └── my_pipeline.yaml
├── config.yaml
└── main.py (optional)

Getting Started With Examples

To get started quickly, FinDrum includes runnable examples in the examples/ folder.

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

findrum_platform-1.0.1.tar.gz (9.5 kB view details)

Uploaded Source

Built Distribution

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

findrum_platform-1.0.1-py3-none-any.whl (10.7 kB view details)

Uploaded Python 3

File details

Details for the file findrum_platform-1.0.1.tar.gz.

File metadata

  • Download URL: findrum_platform-1.0.1.tar.gz
  • Upload date:
  • Size: 9.5 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.1.0 CPython/3.12.11

File hashes

Hashes for findrum_platform-1.0.1.tar.gz
Algorithm Hash digest
SHA256 352bb0fc8af31c458e6a313acaaf5777a61d8a4b4cb45da4c9013f6b2a6777ca
MD5 045a28c5b6113ba47ab0dc135ff00549
BLAKE2b-256 d10fed60747e9cc9b64061f0fe6bc30fa78836b44f3dc7d152931991ad41f9ff

See more details on using hashes here.

File details

Details for the file findrum_platform-1.0.1-py3-none-any.whl.

File metadata

File hashes

Hashes for findrum_platform-1.0.1-py3-none-any.whl
Algorithm Hash digest
SHA256 69ce6d80f720bead04f2af432976a443931d02109f744592d2c436580c5f4f6a
MD5 0872077a998391f24ccae150c09d559c
BLAKE2b-256 9c23f0b7a50e76b3751de4d5c9afb7aa503f3d9fde192d6b8aa80d70c5491c76

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