a8m / a8m/kinesis-producer

Use In AWS Lambda

Ouverte
#28 2 commentaires 0 réactions 0 personnes assignées Voir sur GitHub
Langage dominant
Go
Étoiles
150
Forks
47
Métriques de merge des PR
Aucune PR mergée en 30 j

Description

I want to use this library with AWS Lambda but producer cannot be reused after stop.
Below example code occuers `Unable to Put record. Producer is already stopped` when runs at second times.

```go
package main

import (
"fmt"
"github.com/a8m/kinesis-producer"
"github.com/aws/aws-lambda-go/events"
"github.com/aws/aws-lambda-go/lambda"
"github.com/aws/aws-sdk-go/aws/session"
"github.com/aws/aws-sdk-go/service/kinesis"
"golang.org/x/sync/errgroup"
"os"
)

var pr = producer.New(&producer.Config{
StreamName: os.Getenv("KINESIS_STREAM"),
Client: kinesis.New(session.Must(session.NewSession())),
})

func handle(e events.KinesisEvent) error {
eg := errgroup.Group{}

pr.Start()
eg.Go(func() error {
for r := range pr.NotifyFailures() {
return r
}
return nil
})

for _, r := range e.Records {
// Any logic for each records
if err := pr.Put(r.Kinesis.Data, r.Kinesis.PartitionKey); err != nil {
return err
}
}
pr.Stop()
return eg.Wait()
}

func main() {
lambda.Start(handle)
}
```

Of course, if I generate Producer every time, it works well but I want to reuse the producer as much as possible.

```go
var kc = kinesis.New(session.Must(session.NewSession()))

func handle(e events.KinesisEvent) error {
var pr = producer.New(&producer.Config{
StreamName: os.Getenv("KINESIS_STREAM"),
Client: kc,
})

eg := errgroup.Group{}

pr.Start()
eg.Go(func() error {
for r := range pr.NotifyFailures() {
return r
}
return nil
})

for _, r := range e.Records {
if err := pr.Put(r.Kinesis.Data, r.Kinesis.PartitionKey); err != nil {
return err
}
}
pr.Stop()
return eg.Wait()
}
```

Is it possible to make the Producer discretion by making the `flush()` method of the Producer public or by `Stop()` and then `Start()` again?

Guide de contribution

Aucun guide de contribution indexé pour ce dépôt

Évaluation

Cette issue n'a pas encore été évaluée.

Recevez les nouvelles issues par e-mail

Un résumé court des issues GitHub adaptées aux débutants.