a8m / a8m/kinesis-producer

Use In AWS Lambda

未关闭
#28 2 条评论 0 个 reaction 已指派 0 人 在 GitHub 查看
主要语言
Go
星标
150
派生
47
PR 合并指标
30 天内没有已合并 PR

描述

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?

贡献指南

这个仓库没有索引到贡献指南

评估

这个 Issue 还没有评估数据。

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。