stream

package module
v0.0.0-...-2ed5266 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Feb 3, 2023 License: BSD-3-Clause Imports: 13 Imported by: 0

README

Streams is a pipeline based, stream processing action for the Project Flogo Ecosystem

Flogo Stream

Edge devices have the potential for producing millions or even billions of events at rapid intervals, often times the events on their own are meaningless, hence the need to provide basic streaming operations against the slew of events.

A native streaming action as part of the Project Flogo Ecosystem accomplishes the following primary objectives:

  • Enables apps to implement basic streaming constructs in a simple pipeline fashion
  • Provides non-persistent state for streaming operations
    • Streams are persisted in memory until the end of the pipeline
  • Serves as a pre-process pipeline for raw data to perform basic mathematical and logical operations. Ideal for feeding ML models

Some of the key highlights include:

😀 Simple pipeline construct enables a clean, easy way of dealing with streams of data
Stream aggregation across streams using time or event tumbling & sliding windows
🙌 Join streams from multiple event sources
🌪 Filter out the noise with stream filtering capabilities

Getting Started

We’ve made building powerful streaming pipelines as easy as possible. Develop your pipelines using:

  • A simple, clean JSON-based DSL
  • Golang API

See the sample below of an aggregation pipeline (for brevity, the triggers and metadata of the resource has been omitted). Also don’t forget to check out the examples in the repo.

  "stages": [
    {
      "ref": "github.com/project-flogo/stream/activity/aggregate",
      "settings": {
        "function": "sum",
        "windowType": "timeTumbling",
        "windowSize": "5000"
      },
      "input": {
        "value": "=$.input"
      }
    },
    {
      "ref": "github.com/project-flogo/contrib/activity/log",
      "input": {
        "message": "=$.result"
      }
    }
  ]

Try out the example

Firstly you should install the install the Flogo CLI.

Next you should download our aggregation example agg-flogo.json.

We'll create a our application using the example file, we'll call it myApp

$ flogo create -f agg-flogo.json myApp

Now, build it...

$ cd myApp/
$ flogo build

Activities

Flogo Stream also provides some activities to assist in stream processing.

  • Aggregate : This activity allows you to aggregate data and calculate an average or sliding average.
  • Filter : This activity allows you to filter out data in a streaming pipeline.

License

Flogo source code in this repository is under a BSD-style license, refer to LICENSE

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type ActionFactory

type ActionFactory struct {
	// contains filtered or unexported fields
}

func (*ActionFactory) Initialize

func (f *ActionFactory) Initialize(ctx action.InitContext) error

func (*ActionFactory) New

func (f *ActionFactory) New(config *action.Config) (action.Action, error)

type Settings

type Settings struct {
	StreamURI     string `md:"streamURI"`
	PipelineURI   string `md:"pipelineURI"`
	GroupBy       string `md:"groupBy"`
	OutputChannel string `md:"outputChannel"`
}

type StreamAction

type StreamAction struct {
	// contains filtered or unexported fields
}

func (*StreamAction) IOMetadata

func (s *StreamAction) IOMetadata() *metadata.IOMetadata

func (*StreamAction) Info

func (s *StreamAction) Info() *action.Info

func (*StreamAction) Metadata

func (s *StreamAction) Metadata() *action.Metadata

func (*StreamAction) Run

func (s *StreamAction) Run(context context.Context, inputs map[string]interface{}, handler action.ResultHandler) error

Directories

Path Synopsis
activity
aggregate Module
filter Module
service
telemetry Module
trigger
streamtester Module

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL