diff --git a/internal/pkg/pipeline/task/sqs/README.md b/internal/pkg/pipeline/task/sqs/README.md index 37f291c..fab5705 100644 --- a/internal/pkg/pipeline/task/sqs/README.md +++ b/internal/pkg/pipeline/task/sqs/README.md @@ -19,7 +19,8 @@ The task automatically determines its mode based on the presence of input/output | `type` | string | `sqs` | Must be "sqs" | | `queue_url` | string | - | SQS queue URL (required) | | `concurrency` | int | `10` | Number of concurrent `DeleteMessage` calls | -| `max_messages` | int | `10` | Maximum number of messages to receive per batch | +| `max_messages` | int | `10` | Per poll: max messages per `ReceiveMessage` (AWS cap 10) | +| `max_records` | int | `0` (unlimited) | Read-mode cap on records forwarded downstream (per worker). `0` = unlimited. | | `wait_time_seconds` | int | `10` | Long polling wait time in seconds | | `exit_on_empty` | bool | `false` | Exit when a receive returns no messages. FIFO: empty poll is not drain while this task holds receipts or the queue still has visible, in-flight, or delayed messages. | | `end_after` | duration | - | Stop polling after this much time (read mode); e.g. `5m` | @@ -29,8 +30,8 @@ The task automatically determines its mode based on the presence of input/output | `context` | map | - | JQ expressions whose results are stored on each record for downstream tasks | | `fail_on_error` | bool | `false` | Whether to stop the pipeline if this task encounters an error | -In read mode the task polls until the queue drains (`exit_on_empty`) or `end_after` elapses; -with neither set it polls indefinitely. +In read mode the task polls until the queue drains (`exit_on_empty`), `end_after` elapses, or `max_records` is reached; +with none of those set it polls indefinitely. ## Example Configurations diff --git a/internal/pkg/pipeline/task/sqs/sqs.go b/internal/pkg/pipeline/task/sqs/sqs.go index 97b0f08..ca711ec 100644 --- a/internal/pkg/pipeline/task/sqs/sqs.go +++ b/internal/pkg/pipeline/task/sqs/sqs.go @@ -51,6 +51,7 @@ type sqs struct { QueueURL string `yaml:"queue_url" json:"queue_url" validate:"required"` Concurrency int `yaml:"concurrency,omitempty" json:"concurrency,omitempty"` MaxMessages int32 `yaml:"max_messages,omitempty" json:"max_messages,omitempty"` + MaxRecords int `yaml:"max_records,omitempty" json:"max_records,omitempty" validate:"omitempty,gte=0"` WaitTimeSeconds int `yaml:"wait_time_seconds,omitempty" json:"wait_time_seconds,omitempty"` ExitOnEmpty bool `yaml:"exit_on_empty,omitempty" json:"exit_on_empty,omitempty"` MessageGroupId string `yaml:"message_group_id,omitempty" json:"message_group_id,omitempty"` // used for FIFO queues @@ -153,6 +154,7 @@ func (s *sqs) getMessages(ctx context.Context, output chan<- *record.Record) err defer cancel() } + recordsRead := 0 for { select { case <-ctx.Done(): @@ -189,21 +191,23 @@ func (s *sqs) getMessages(ctx context.Context, output chan<- *record.Record) err if output == nil { s.deleteMessage(m.MessageId, m.ReceiptHandle) - continue - } - - if s.Delivery == deliveryAtMostOnce { + } else if s.Delivery == deliveryAtMostOnce { s.SendData(ctx, []byte(*m.Body), output) del := ack.New() del.AddBranch(1) del.Done() s.enqueueDelete(m, del) - continue + } else { + msgAck := ack.New() + s.SendData(ack.WithContext(ctx, msgAck), []byte(*m.Body), output) + s.enqueueDelete(m, msgAck) } - msgAck := ack.New() - s.SendData(ack.WithContext(ctx, msgAck), []byte(*m.Body), output) - s.enqueueDelete(m, msgAck) + recordsRead++ + if s.MaxRecords > 0 && recordsRead >= s.MaxRecords { + fmt.Printf("SQS max_records (%d) reached for queue %s, stopping reader\n", s.MaxRecords, s.QueueURL) + return nil + } } } }