Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 7 additions & 8 deletions .github/workflows/lint.yml
Original file line number Diff line number Diff line change
Expand Up @@ -15,16 +15,15 @@ jobs:
name: lint
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v3
- uses: actions/setup-go@v4
- uses: actions/checkout@v5
- uses: actions/setup-go@v5
with:
go-version-file: 'go.mod'
cache: false
- name: golangci-lint
uses: golangci/golangci-lint-action@v3
uses: golangci/golangci-lint-action@v8
with:
version: v1.55
skip-pkg-cache: true
version: v2.13

govulncheck:
runs-on: ubuntu-latest
Expand All @@ -38,14 +37,14 @@ jobs:
commitlint:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v3
- uses: actions/checkout@v5
with:
fetch-depth: 0

- name: Install NodeJS
uses: actions/setup-node@v4
uses: actions/setup-node@v5
with:
node-version: 18
node-version: 24

- name: Install commitlint and convention
run: |
Expand Down
4 changes: 2 additions & 2 deletions .github/workflows/release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,11 @@ jobs:
pull-requests: write # to be able to comment on released pull requests
steps:
- name: Checkout
uses: actions/checkout@v3
uses: actions/checkout@v5
with:
fetch-depth: 0
- name: Setup Node.js
uses: actions/setup-node@v3
uses: actions/setup-node@v5
with:
node-version: "lts/*"
- name: Install dependencies
Expand Down
4 changes: 2 additions & 2 deletions .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,8 @@ jobs:
tests:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v3
- uses: actions/setup-go@v4
- uses: actions/checkout@v5
- uses: actions/setup-go@v5
with:
go-version-file: 'go.mod'
cache: false
Expand Down
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,9 @@
# Output of the go coverage tool, specifically when used with LiteIDE
*.out

# IDEs and editors
/.idea

# Dependency directories (remove the comment below to include it)
# vendor/

Expand Down
230 changes: 121 additions & 109 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,145 +33,156 @@ Let's create some simple examples to demonstrate how to use this library to proc
### Basic example

```go
package main

import (
"context"
"fmt"
"log"

formigo "github.com/Pod-Point/go-queue-worker"
workerSqs "github.com/Pod-Point/go-queue-worker/clients/sqs"

"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
"context"
"log"

formigo "github.com/Pod-Point/go-queue-worker"

"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
)

func main() {
ctx := context.Background()

awsCfg, err := config.LoadDefaultConfig(ctx)
if err != nil {
log.Fatalln("Unable to create AWS config", err)
}

sqsSvc := sqs.NewFromConfig(awsCfg)
sqsClient, err := formigo.NewSqsClient(ctx, formigo.SqsClientConfiguration{
Svc: sqsSvc,
ReceiveMessageInput: &sqs.ReceiveMessageInput{
QueueUrl: &queueUrl,
MaxNumberOfMessages: 1,
VisibilityTimeout: 30,
WaitTimeSeconds: 20,
},
})
if err != nil {
return fmt.Errorf("unable to create sqs client: %w", err)
}

wkr := formigo.NewWorker(formigo.Configuration{
Client: sqsClient,
Concurrency: 100,
Consumer: formigo.NewMessageConsumer(formigo.MessageConsumerConfiguration{
Handler: func(ctx context.Context, msg formigo.Message) error {
log.Println("Got Message", msg.Content())

// Assert the type of message to get the body or any other attributes
log.Println("Message body", *msg.Content().(types.Message).Body)

return nil
},
}),
})

err = wkr.Run(ctx)
if err != nil {
log.Fataln("Worker stopped with error", err)
}
ctx := context.Background()

queueUrl := "https://sqs.eu-west-1.amazonaws.com/123456789012/my-queue"

awsCfg, err := config.LoadDefaultConfig(ctx)
if err != nil {
log.Fatalln("Unable to create AWS config", err)
}

sqsSvc := sqs.NewFromConfig(awsCfg)
sqsClient, err := formigo.NewSqsClient(ctx, formigo.SqsClientConfiguration{
Svc: sqsSvc,
ReceiveMessageInput: &sqs.ReceiveMessageInput{
QueueUrl: &queueUrl,
MaxNumberOfMessages: 10,
VisibilityTimeout: 30,
WaitTimeSeconds: 20,
},
})
if err != nil {
log.Fatalln("Unable to create sqs client", err)
}

wkr := formigo.NewWorker(formigo.Configuration{
Client: sqsClient,
Concurrency: 100,
Consumer: formigo.NewMessageConsumer(formigo.MessageConsumerConfiguration{
Handler: func(ctx context.Context, msg formigo.Message) error {
// Assert the type of message to get the body or any other attributes
log.Println("Message body", *msg.Content().(types.Message).Body)

return nil
},
}),
})

if err := wkr.Run(ctx); err != nil {
log.Fatalln("Worker stopped with error", err)
}
}
```

In this example, we have created a worker that consumes messages one at a time from an AWS SQS queue. The polling phase retrieves 10 messages from the queue, but the handler processes them individually.

By default, the worker's concurrency is set to 100, meaning it can process up to 100 messages concurrently, optimizing throughput and efficiency.

If any errors occur during message handling, the worker will log them using log.PrintLn by default. Additionally, the worker is configured to stop if it encounters more than 3 errors within any 120-second interval.
If any errors occur during message handling, the worker will log them using log.Println by default. Additionally, the worker is configured to stop if it encounters more than 3 errors within any 120-second interval.

Please note that these are the default settings, and you can customize the concurrency level, error handling, and other parameters to suit your specific requirements.


### Batching

```go
package main

import (
"context"
"fmt"
"log"
"context"
"log"
"time"

formigo "github.com/Pod-Point/go-queue-worker"
workerSqs "github.com/Pod-Point/go-queue-worker/clients/sqs"
formigo "github.com/Pod-Point/go-queue-worker"

"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
)

func main() {
ctx := context.Background()

awsCfg, err := config.LoadDefaultConfig(ctx)
if err != nil {
log.Fatalln("Unable to create AWS config", err)
}

sqsSvc := sqs.NewFromConfig(awsCfg)
sqsClient, err := formigo.NewSqsClient(ctx, formigo.SqsClientConfiguration{
Svc: sqsSvc,
ReceiveMessageInput: &sqs.ReceiveMessageInput{
QueueUrl: &queueUrl,
MaxNumberOfMessages: 1,
VisibilityTimeout: 30,
WaitTimeSeconds: 20,
},
})
if err != nil {
return fmt.Errorf("unable to create sqs client: %w", err)
}

wkr := formigo.NewWorker(formigo.Configuration{
Client: sqsClient,
Concurrency: 100,
Consumer: formigo.BatchConsumer(formigo.BatchConsumerConfiguration{
BufferConfig: formigo.BatchBufferConfiguration{
Size: 100,
Timeout: time.Second * 5,
},
Handler: func(ctx context.Context, msgs []formigo.Message) error {
log.Printf("Got %d messages to process\n", len(msgs)

// Assert the type of message to get the body or any other attributes

for i, msg := range msgs {
log.Printf("Message %d body: %s", i, *msg.Content().(types.Message).Body)
}

return nil
},
}),
})

err = wkr.Run(ctx)
if err != nil {
log.Fataln("Worker stopped with error", err)
}
ctx := context.Background()

queueUrl := "https://sqs.eu-west-1.amazonaws.com/123456789012/my-queue"

awsCfg, err := config.LoadDefaultConfig(ctx)
if err != nil {
log.Fatalln("Unable to create AWS config", err)
}

sqsSvc := sqs.NewFromConfig(awsCfg)
sqsClient, err := formigo.NewSqsClient(ctx, formigo.SqsClientConfiguration{
Svc: sqsSvc,
ReceiveMessageInput: &sqs.ReceiveMessageInput{
QueueUrl: &queueUrl,
MaxNumberOfMessages: 10,
VisibilityTimeout: 30,
WaitTimeSeconds: 20,
},
})
if err != nil {
log.Fatalln("Unable to create sqs client", err)
}

wkr := formigo.NewWorker(formigo.Configuration{
Client: sqsClient,
Concurrency: 100,
Consumer: formigo.NewBatchConsumer(formigo.BatchConsumerConfiguration{
BufferConfig: formigo.BatchConsumerBufferConfiguration{
Size: 100,
Timeout: time.Second * 5,
},
Handler: func(ctx context.Context, msgs []formigo.Message) (formigo.BatchResponse, error) {
log.Printf("Got %d messages to process\n", len(msgs))

// Collect the ids of any messages that failed. They won't be deleted,
// so they return to the queue once their visibility timeout expires.
var failed []interface{}

for _, msg := range msgs {
// Assert the type of message to get the body or any other attributes
body := *msg.Content().(types.Message).Body

if err := process(ctx, body); err != nil {
log.Println("Unable to process message", err)
failed = append(failed, msg.Id())
}
}

return formigo.BatchResponse{FailedMessagesId: failed}, nil
},
}),
})

if err := wkr.Run(ctx); err != nil {
log.Fatalln("Worker stopped with error", err)
}
}

func process(ctx context.Context, body string) error { return nil }
```

In this example, we have created a worker that efficiently consumes batches of messages from an AWS SQS queue. The handler will be invoked either when the buffer is full or when a specified timeout expires.

It's essential to note that the timer starts as soon as the first message is added to the buffer.

Returning message ids in `BatchResponse.FailedMessagesId` marks those messages as failed: they are not deleted from the queue, so they become visible again once their visibility timeout expires. Every other message in the batch is deleted. Returning an error from the handler instead fails the whole batch and reports it to the `ErrorConfig.ReportFunc`.

By processing messages in batches, the worker can significantly enhance throughput for specific use cases or reduce resource consumption. For instance, it can be leveraged for batch insertions or deletions.

## Configuration
Expand All @@ -182,7 +193,8 @@ By processing messages in batches, the worker can significantly enhance throughp
| Concurrency | Number of Go routines that process the messages from the Queue. Higher values are useful for slow I/O operations in the consumer's handler. | 100 |
| Retrievers | Number of Go routines that retrieve messages from the Queue. Higher values are helpful for slow networks or when consumers are quicker. | 1 |
| ErrorConfig | Defines the error threshold and interval for worker termination and error reporting function. | None |
| Consumer | The message consumer, either MessageConsumer or BatchConsumer. | None |
| Consumer | The message consumer, built with either `NewMessageConsumer` or `NewBatchConsumer`. This is a required configuration. | None |
| DeleterConfig | Size and timeout of the buffer that batches messages before they are deleted from the queue. | 10, 500ms |

## License

Expand Down
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
module github.com/Pod-Point/go-queue-worker

go 1.21
go 1.26

require (
github.com/aws/aws-sdk-go-v2/service/sqs v1.34.5
Expand Down
Loading