fstream

package module
v1.0.1 Latest Latest
Warning

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

Go to latest
Published: Nov 5, 2019 License: MIT Imports: 11 Imported by: 0

README

fstream

fstream is the simple solution for collect/read input messages in binary files queue. Each file contains all messages that were received within 1 minute.

File format:

<uint16> <msg1> <uint16> <msg2> ...

Message size is stored in uint16 value.

All files are saved in the selected directory and have names similar to the following 000000, 000001, 000002, ... The name 000000 follows name 999999. The current file name is stored in idx file.

Writer saves the last 10000 files and and deletes older.

fstream is useful when there is a large flow of events and there is no way to process them on the fly or use solutions such as RabbitMQ. It's effective for statistics aggregation and subsequent processing.

Install

go get "github.com/belfinor/fstream"

Writer example

package main

import (
  "bufio"
  "os"
  "strings"

  "github.com/belfinor/fstream"
)

func main() {

  os.Mkdir(".data", 0777)

  w := fstream.NewWriter(".data", ".data/writer.idx")
  defer w.Close()

  br := bufio.NewReader(os.Stdin)

  for {

    str, err := br.ReadString('\n')
    if err != nil && str == "" {
      break
    }

    str = strings.TrimSpace(str)

    if str == "" {
      continue
    }

    w.Write([]byte(str))

  }
}

Reader example

package main

import (
  "fmt"
  "os"
  "time"

  "github.com/belfinor/fstream"
)

func handler(data []byte) {
  fmt.Println(string(data)
}


func main() {

  os.Mkdir(".data", 0777)

  w := fstream.NewReader(".data", ".data/reader.idx", handler)
  defer w.Close()

  wait := make(chan int)

  <-wait
}

Reader hooks

AfterReadFile

Called after Reader has read the next binary file. Example:

r.AfterReadFile = func(filename string) {
  fmt.Printf("%s processed", filename)
}

BeforeReadFile

Сalled before Reader starts processing the next binary file. Example:

r.BeforeReadFile = func(filename string) {
  fmt.Printf("process %s", filename)
}

Documentation

Index

Constants

View Source
const (
	MessageLimit  int   = 2048
	FileNumberMod int64 = 1000000
	SaveFiles     int64 = 10000
	SavePeriod    int64 = 60
)

Variables

This section is empty.

Functions

This section is empty.

Types

type Reader

type Reader struct {
	BeforeReadFile func(string)
	AfterReadFile  func(string)
	// contains filtered or unexported fields
}

func NewReader

func NewReader(path string, idx string, handler func([]byte)) *Reader

func (*Reader) Close

func (r *Reader) Close()

func (*Reader) ReadFile

func (r *Reader) ReadFile(filename string)

type Writer

type Writer struct {
	Input chan []byte
	File  *os.File
	Path  string
	Cnt   *fcounter.Counter
	// contains filtered or unexported fields
}

func NewWriter

func NewWriter(path string, idx string) *Writer

func (*Writer) Close

func (w *Writer) Close()

func (*Writer) Write

func (w *Writer) Write(data []byte) (int, error)

Jump to

Keyboard shortcuts

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