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
13 changes: 8 additions & 5 deletions internal/pkg/pipeline/task/sqs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,12 +18,13 @@ The task automatically determines its mode based on the presence of input/output
| `name` | string | - | Task name for identification |
| `type` | string | `sqs` | Must be "sqs" |
| `queue_url` | string | - | SQS queue URL (required) |
| `concurrency` | int | `10` | Number of concurrent workers that acknowledge (delete) fully-processed messages |
| `concurrency` | int | `10` | Number of concurrent `DeleteMessage` calls |
| `max_messages` | int | `10` | Maximum number of messages to receive per batch |
| `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` |
| `message_group_id` | string | - | Message group ID for FIFO queues |
| `delivery` | string | `at-least-once` | Read mode: `at-least-once` deletes after downstream finishes; `at-most-once` deletes on receive. |
| `task_concurrency` | int | `1` | Number of competing-consumer workers for this task |
| `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 |
Expand Down Expand Up @@ -73,10 +74,12 @@ tasks:
## Message Acknowledgment

When reading from a queue, a message's receipt is deleted only once every downstream task
has finished with the record produced from it. A task that returns an error while holding a
record leaves the receipt alone, so SQS redelivers the message after the visibility timeout
rather than losing it. Delivery is therefore at-least-once: a pipeline may see a message
more than once.
has finished with the record produced from it (`delivery: at-least-once`, the default). A task
that returns an error while holding a record leaves the receipt alone, so SQS redelivers the
message after the visibility timeout rather than losing it. Delivery is therefore at-least-once:
a pipeline may see a message more than once. `delivery: at-most-once` forwards the record
without an ack and deletes on the same `concurrency` pool, so FIFO groups are not blocked by
downstream work; a crash can lose the message.

That covers failures, not drops. A task configured to skip a bad record counts it as
finished, so the receipt is deleted and the message does not come back — and `jq`'s
Expand Down
59 changes: 43 additions & 16 deletions internal/pkg/pipeline/task/sqs/sqs.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,13 @@ const (
defaultRegion = "us-west-2"
)

type deliveryMode string

const (
deliveryAtMostOnce deliveryMode = "at-most-once"
deliveryAtLeastOnce deliveryMode = "at-least-once"
)

var (
awsRegionRegex = regexp.MustCompile(`^[a-z]{2}-[a-z]+-\d+$`)
ctx = context.Background()
Expand All @@ -41,12 +48,13 @@ var (

type sqs struct {
task.ServerBase `yaml:",inline" json:",inline"`
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"`
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
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"`
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
Delivery deliveryMode `yaml:"delivery,omitempty" json:"delivery,omitempty"`

client *qs.Client
tracker *ack.Tracker
Expand All @@ -70,14 +78,23 @@ func (s *sqs) Init() error {
return fmt.Errorf("queue_url is required")
}

switch s.Delivery {
case "", deliveryAtLeastOnce:
s.Delivery = deliveryAtLeastOnce
case deliveryAtMostOnce:
default:
return fmt.Errorf("invalid delivery mode %q: must be %q or %q", s.Delivery, deliveryAtMostOnce, deliveryAtLeastOnce)
}

s.tracker = ack.NewTracker(s.Concurrency)

region := s.extractRegionFromQueueURL()
awsConfig, err := config.LoadDefaultConfig(ctx, config.WithRegion(region))
if err != nil {
return fmt.Errorf("failed to load AWS config: %w", err)
}

s.client = qs.NewFromConfig(awsConfig)
s.tracker = ack.NewTracker(s.Concurrency)

return nil
}
Expand Down Expand Up @@ -170,28 +187,38 @@ func (s *sqs) getMessages(ctx context.Context, output chan<- *record.Record) err

for _, m := range receiveMessageOutput.Messages {

// nothing to forward to, so there's no downstream ack to wait
// for: delete the receipt right away.
if output == nil {
s.deleteMessage(m.MessageId, m.ReceiptHandle)
continue
}

if s.Delivery == deliveryAtMostOnce {
s.SendData(ctx, []byte(*m.Body), output)
del := ack.New()
del.AddBranch(1)
del.Done()
s.enqueueDelete(m, del)
continue
}

msgAck := ack.New()
s.SendData(ack.WithContext(ctx, msgAck), []byte(*m.Body), output)

s.outstanding.Add(1)
s.tracker.Track(msgAck, &messageAck{
sqs: s,
messageId: m.MessageId,
receiptHandle: m.ReceiptHandle,
})
s.enqueueDelete(m, msgAck)
}
}
}

}

func (s *sqs) enqueueDelete(m types.Message, a *ack.Ack) {
s.outstanding.Add(1)
s.tracker.Track(a, &messageAck{
sqs: s,
messageId: m.MessageId,
receiptHandle: m.ReceiptHandle,
})
}

// messageAck acknowledges one received message on behalf of ack.Tracker.
type messageAck struct {
sqs *sqs
Expand Down
Loading