Kanject.Core.Queue.Provider.AwsSqs
3.8.2
Prefix Reserved
See the version list below for details.
dotnet add package Kanject.Core.Queue.Provider.AwsSqs --version 3.8.2
NuGet\Install-Package Kanject.Core.Queue.Provider.AwsSqs -Version 3.8.2
<PackageReference Include="Kanject.Core.Queue.Provider.AwsSqs" Version="3.8.2" />
<PackageVersion Include="Kanject.Core.Queue.Provider.AwsSqs" Version="3.8.2" />
<PackageReference Include="Kanject.Core.Queue.Provider.AwsSqs" />
paket add Kanject.Core.Queue.Provider.AwsSqs --version 3.8.2
#r "nuget: Kanject.Core.Queue.Provider.AwsSqs, 3.8.2"
#:package Kanject.Core.Queue.Provider.AwsSqs@3.8.2
#addin nuget:?package=Kanject.Core.Queue.Provider.AwsSqs&version=3.8.2
#tool nuget:?package=Kanject.Core.Queue.Provider.AwsSqs&version=3.8.2
Kanject.Core.Queue.Provider.AwsSqs
The Amazon SQS implementation of Kanject.Core.Queue.Abstractions. It covers publishing (single, batch and routed messages), consuming (AWS Lambda SQS events or long-polling hosted consumers in any .NET host), partial-batch failure reporting, and optional queue and dead-letter-queue creation.
Use it in services that exchange work through SQS. Declare messages and consumers with attributes; the Kanject.Core.Queue.Provider.AwsSqs.Annotations generator writes the registration, send and Lambda-dispatch code for them.
Installation
<ItemGroup>
<PackageReference Include="Kanject.Core.Queue.Provider.AwsSqs" />
<PackageReference Include="Kanject.Core.Queue.Provider.AwsSqs.Annotations" PrivateAssets="all" />
</ItemGroup>
or
dotnet add package Kanject.Core.Queue.Provider.AwsSqs
dotnet add package Kanject.Core.Queue.Provider.AwsSqs.Annotations
Targets .NET 8, .NET 9 and .NET 10. This package depends on the Annotations package, so the generators normally reach your project through it. Add the direct Annotations reference whenever you use the attributes anyway, so the generator version is explicit.
Quick start
1. Declare the message and the consumer
using Amazon.Lambda.SQSEvents;
using Kanject.Core.Queue.Abstractions.Models;
using Kanject.Core.Queue.Provider.AwsSqs.Annotations.Attributes;
namespace Shop.Orders;
[QueueMessage("orders")] // SQS queue "orders_queue"
public sealed class OrderPlaced
{
public string OrderId { get; set; } = string.Empty;
public decimal Total { get; set; }
}
[QueueConsumer(Message = typeof(OrderPlaced))] // queue name comes from [QueueMessage]
[QueueConsumerDependency(typeof(IFulfilmentService))]
public partial class OrderPlacedConsumer
{
protected override async Task ConsumeAsync(
List<MessageContext<OrderPlaced>> messages, CancellationToken cancellationToken)
{
foreach (var context in messages)
{
try
{
if (context.Message is null) // body could not be deserialized
throw new InvalidOperationException("Unreadable message.");
await FulfilmentService.ShipAsync(context.Message, cancellationToken);
await AcknowledgeAsync(context);
}
catch (Exception)
{
Response.BatchItemFailures.Add(new SQSBatchResponse.BatchItemFailure
{
ItemIdentifier = context.MessageId
});
}
}
}
}
The generator makes OrderPlacedConsumer derive from AbstractQueueConsumer<OrderPlaced>. It also adds a constructor that takes IQueueRegistry plus each declared dependency, and exposes each dependency as a property (FulfilmentService here). Finally, it emits AddOrderPlacedConsumer(), ProcessOrderPlacedEventAsync(...), AddOrdersQueue() and EnqueueOrderPlacedAsync(...).
2. Consume in AWS Lambda
using Amazon.Lambda.Core;
using Amazon.Lambda.SQSEvents;
using Kanject.Core.Queue.Provider.AwsSqs.Extensions;
using Microsoft.Extensions.DependencyInjection;
using Shop.Orders;
public sealed class Function
{
private static readonly IServiceProvider Services = BuildServices();
private static IServiceProvider BuildServices()
{
var services = new ServiceCollection();
services.AddSingleton<IFulfilmentService, FulfilmentService>();
services
.AddAwsSqsGlobalQueueConfiguration(options =>
{
// No keys: the SQS client uses the default AWS credential chain, which in Lambda
// is the function's execution role. Set AWSAccessKey/AWSSecretKey only for static keys.
options.AWSRegion = "eu-west-1";
})
.AddOrderPlacedConsumer();
return services.BuildServiceProvider();
}
// SQS event source mapping with ReportBatchItemFailures enabled.
public Task<SQSBatchResponse> HandleAsync(SQSEvent sqsEvent, ILambdaContext context)
=> Services.ProcessOrderPlacedEventAsync(sqsEvent);
}
ProcessOrderPlacedEventAsync deserializes each record into a MessageContext<OrderPlaced>, passes the batch to ConsumeAsync, and returns the consumer's Response. Only the ids you added to BatchItemFailures are retried. Every invocation starts with an empty Response, so a failure recorded in one invocation is never returned again by the next one on a warm container. It also takes an optional CancellationToken, which it forwards to ConsumeAsync.
3. Send
using Kanject.Core.Queue.Abstractions.Interfaces;
using Shop.Orders;
// Registration: services.AddAwsSqsGlobalQueueConfiguration(...).AddOrdersQueue();
public sealed class CheckoutService(IQueueManagerService queues)
{
public async Task CompleteAsync(OrderPlaced order)
{
if (!await queues.EnqueueOrderPlacedAsync(order))
throw new InvalidOperationException($"Order {order.OrderId} was not queued.");
}
}
The message is serialized with System.Text.Json. Other generated overloads take metadata (sent as SQS message attributes), an IList<OrderPlaced> batch, or (message, metadata) tuples.
Consuming in a long-running host
Every registered consumer is also an IHostedService. Polling is opt-in: set WatchQueue and each consumer long-polls its queue when the host starts. With WatchQueue off (the default) the hosted service does nothing, so a host that starts its hosted services inside Lambda never polls alongside the event source mapping.
using Kanject.Core.Queue.Provider.AwsSqs.Abstractions.Extensions;
using Kanject.Core.Queue.Provider.AwsSqs.Extensions;
using Shop.Orders;
var builder = WebApplication.CreateBuilder(args);
builder.Services
.AddAwsSqsGlobalQueueConfiguration(options =>
{
// credentials and region as above
options.WatchQueue = true; // hosted consumers poll their queue
options.WatcherSleepTime = 5; // seconds between polls
options.UseDeadLetterQueue = true;
})
.AddOrdersQueue()
.AddOrderPlacedConsumer();
var app = builder.Build();
// Resolves every registered queue in the background, creating missing ones
// when CreateQueueIfNotExist is true (the default).
app.UseAwsSqsQueueProvider();
app.Run();
For a generic host, call host.Services.UseAwsSqsQueueProvider(). Each poll receives up to MaximumReceiveMessageCount messages and passes them to ConsumeAsync with a token that's cancelled when the host stops the consumer. Polling failures back off exponentially, up to 30 seconds. In this mode, AcknowledgeAsync is what deletes a processed message from the queue.
Routed queues
One queue can carry several message kinds, dispatched by a route tag:
[RouteQueueConsumer(QueueName = "order-events", Message = typeof(OrderEvent))]
[QueueConsumerRoute("order.created")]
[QueueConsumerRoute("order.updated")]
public partial class OrderChangedHandler
{
protected override Task ConsumeAsync(
List<MessageContext<OrderEvent>> messages, CancellationToken cancellationToken)
{
/* handle, acknowledge, or add to Response.BatchItemFailures */
return Task.CompletedTask;
}
}
// Registration: maps every [QueueConsumerRoute] declared for "order-events"
services.AddOrderEventsQueueQueueRouter();
// Lambda handler
Task<SQSBatchResponse> HandleAsync(SQSEvent sqsEvent, ILambdaContext context)
=> Services.RouteSqsEventToOrderEventsQueueAsync(sqsEvent);
// Publisher: the route travels as the "QueueRoute" message attribute
await queues["order-events"].EnqueueAsync("order.created", orderEvent);
Route consumers are registered as scoped services and resolved in a fresh scope per event. With WatchQueue = true, a hosted SqsQueueConsumerRouter polls and dispatches the queue instead. Without the generator, the same wiring is services.AddQueueConsumerRouter(q => q.UseQueue("order-events").MapConsumerRoute<OrderChangedHandler>("order.created")). MapDefaultConsumerRoute<T>() catches routes that have no mapped consumer, and records that carry no QueueRoute attribute at all. A record that reaches no consumer — no mapped route and no default route, or a mapped consumer that isn't registered — is reported in BatchItemFailures, so SQS retries it and eventually moves it to the dead-letter queue instead of deleting it.
Configuration
AddAwsSqsGlobalQueueConfiguration sets the defaults for every queue and consumer registered after it. Each registration copies the global settings at the moment it's made, and you can override them per queue through the generated Add…(options => …) overloads.
AwsSqsQueueConfiguration property |
Default | Effect |
|---|---|---|
AWSAccessKey, AWSSecretKey |
empty | Static credentials for the SQS client. Leave either empty to use the default AWS credential chain (see below) |
AWSRegion |
empty | Region for the SQS client. Empty uses the SDK's region resolution (AWS_REGION in Lambda) |
EnforceQueueOrdering |
false |
Use a FIFO queue (.fifo name). Every send carries a MessageGroupId (see below) |
CreateQueueIfNotExist |
true |
Create the queue during UseAwsSqsQueueProvider() if it doesn't exist |
UseDeadLetterQueue |
false |
Also create <queue>_dlq and attach a redrive policy |
MaximumMessageRetry |
10 |
Redrive maxReceiveCount before a message moves to the DLQ |
WatchQueue |
false |
Hosted consumers and the hosted router poll the queue only when this is true |
WatcherSleepTime |
10 |
Seconds between polls |
MaximumReceiveMessageCount |
10 |
Messages per receive call |
ReceiveMessageWaitTimeInSeconds |
20 |
Long-poll wait time |
VisibilityWaitTimeInSeconds |
20 |
Visibility timeout |
DelaySeconds |
0 |
Delivery delay for every message |
MessageRetentionPeriod / DlqMessageRetentionPeriod |
4 / 5 |
Retention, in days |
Queue attributes such as delay, retention, wait time and visibility timeout are applied only when the provider creates a queue. Existing queues are left unchanged.
Things to know
- Credentials: when both
AWSAccessKeyandAWSSecretKeyare set, the queue's SQS client signs with them. Otherwise it uses the default AWS credential chain: the Lambda or ECS task role, an EC2 instance profile, theAWS_*environment variables, or the shared profile. The provider callssqs:GetQueueUrl,sqs:GetQueueAttributes,sqs:SendMessage,sqs:ReceiveMessageandsqs:DeleteMessage. It also callssqs:CreateQueuewhen it's allowed to create queues. - Queue names:
ordersbecomes the SQS queueorders_queue(seeQueueNameHelperinKanject.Core.Queue.Abstractions). The queue must exist, either provisioned by your infrastructure or created at startup withUseAwsSqsQueueProvider(). - Send results: enqueue calls return
falseinstead of throwing when SQS rejects the request or the queue can't be resolved; the exception is written to the console. Always check the result. Using a queue key that isn't registered throwsKeyNotFoundExceptionfrom theIQueueManagerServiceindexer. - Batches: SQS accepts at most 10 entries per
SendMessageBatchrequest, so the batch overloads send larger batches in order, 10 messages per request, one request at a time. A batch returnstrueonly if SQS accepted every message. It returnsfalseas soon as a request is rejected or lists any entry underFailed(SQS reports those with HTTP 200), and the rest of the batch isn't sent. The failed entries are written to the console. Earlier requests may already have been delivered, so resending the whole batch can deliver those messages twice. Metadata and routes passed to a batch overload are sent with every message in it. - FIFO queues: with
EnforceQueueOrdering, every send carries aMessageGroupId, which SQS requires on a FIFO queue. By default all messages share one group named after the queue, which gives one strict order. Put aMessageGroupIdentry in the metadata (AwsQueueConstants.MessageGroupIdMetadataKey) to order per group instead. Queues the provider creates have content-based deduplication enabled; a FIFO queue created elsewhere needs it too, because the provider doesn't send a deduplication id. - Unreadable messages: in the Lambda path, a record whose body can't be deserialized still reaches
ConsumeAsync, withMessageleft at its default value. Guard for it and report the failure. The polling path logs such messages and skips them. - Lifetimes: typed consumers (
[QueueConsumer],AddQueueConsumer<T>) are singletons. Any service they depend on is captured for the lifetime of the app, so useIServiceScopeFactoryfor scoped work.Responseis per invocation: insideConsumeAsyncit is always the current batch's response, even if invocations overlap; read after the call returns, it is the most recent invocation's. - Missing consumer: if a consumer isn't registered,
ProcessSqsEventWithQueueConsumerAsync<T>,ProcessSqsEventWithRouteQueueConsumerAsync<T>and the generatedProcess…EventAsyncmethods throwInvalidOperationException. The invocation fails, so SQS retries the whole batch and eventually moves it to the dead-letter queue, whether or not the event source mapping hasReportBatchItemFailuresenabled. - Unrouted records: a record without a
QueueRouteattribute goes to the default route. A record that no consumer receives (no mapped route and no default route, or a mapped consumer that isn't registered) is reported as a failure, not dropped.
Public surface at a glance
| Type or member | Purpose |
|---|---|
AddAwsSqsGlobalQueueConfiguration(Action<AwsSqsQueueConfiguration>) |
Global defaults |
CreateQueue(...), AddSqsQueue(name, ...) |
Register a queue for publishing |
AddQueueConsumer<T>(...) |
Register a consumer as a singleton and hosted service |
AddQueueConsumerRouter(...) with UseQueue, MapConsumerRoute<T>, MapDefaultConsumerRoute<T> |
Register a routed queue |
RouteIncomingSqsQueueEventAsync(sqsEvent, queue) |
Dispatch a Lambda SQSEvent to route consumers without the generator. Return its SQSBatchResponse from the handler. It replaces the obsolete SqsQueueConsumerRouter.ProcessSqsEventAsync, which can't return a response and throws when any record fails |
SQSEvent.CategorizeRecordsByFilter(attributeName, allowedValues) |
Split an event into buckets by a JSON body property |
AbstractQueueConsumer<T>, AbstractRawQueueConsumer |
Base classes for typed and raw (Amazon.SQS.Model.Message) consumers |
AbstractRouteQueueConsumer<T>, AbstractRouteRawQueueConsumer |
Base classes for route consumers |
SqsQueueManager, SqsQueueManagerService, QueueRegistry, SqsQueueConsumerRouter |
The implementations behind the contracts |
The Lambda helpers (ProcessSqsEventWithQueueConsumerAsync<T>, ProcessSqsEventWithRouteQueueConsumerAsync<T>), UseAwsSqsQueueProvider(), and AcknowledgeOrFailAsync for [Parallel] outcomes ship in Kanject.Core.Queue.Provider.AwsSqs.Abstractions.
Related packages
| Package | Role | Availability |
|---|---|---|
Kanject.Core.Queue.Abstractions |
Provider-neutral queue contracts | nuget.org |
Kanject.Core.Queue.Provider.AwsSqs.Abstractions |
SQS contracts, AwsSqsQueueConfiguration, Lambda handler and partial-batch helpers |
nuget.org |
Kanject.Core.Queue.Provider.AwsSqs.Annotations |
Source generator and analyzers | nuget.org |
Kanject.Core.Queue.Provider.AwsSqs.Annotations.Attributes |
The queue attributes | nuget.org |
Kanject.Core |
Base library, including the [Parallel] feature |
nuget.org |
Kanject.Core.Queue.Provider.AwsSqs.Extensions.EventBridge |
Scheduled and recurring delivery through EventBridge Scheduler | Commercial license (not on nuget.org) |
Kanject.Core.Queue.Provider.AwsSqs.Extensions.EventBridge.Annotations |
Generator for typed schedule helpers | Commercial license (not on nuget.org) |
License
Licensed under the Kanject Code Libraries License Agreement (KCLLA); the full text ships in this package as LICENSE.md. Organizations whose trailing-twelve-month gross revenue and total funding raised are each below US$250,000 may use it at no cost under the Free Tier. At or above either threshold a commercial license is required — contact commercial@kanjectbusiness.com.
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net8.0 is compatible. net8.0-android was computed. net8.0-browser was computed. net8.0-ios was computed. net8.0-maccatalyst was computed. net8.0-macos was computed. net8.0-tvos was computed. net8.0-windows was computed. net9.0 is compatible. net9.0-android was computed. net9.0-browser was computed. net9.0-ios was computed. net9.0-maccatalyst was computed. net9.0-macos was computed. net9.0-tvos was computed. net9.0-windows was computed. net10.0 is compatible. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
-
net10.0
- Amazon.Lambda.SQSEvents (>= 3.0.1)
- Kanject.Core.Queue.Provider.AwsSqs.Abstractions (>= 3.8.2)
- Kanject.Core.Queue.Provider.AwsSqs.Annotations (>= 3.8.1)
-
net8.0
- Amazon.Lambda.SQSEvents (>= 3.0.1)
- Kanject.Core.Queue.Provider.AwsSqs.Abstractions (>= 3.8.2)
- Kanject.Core.Queue.Provider.AwsSqs.Annotations (>= 3.8.1)
-
net9.0
- Amazon.Lambda.SQSEvents (>= 3.0.1)
- Kanject.Core.Queue.Provider.AwsSqs.Abstractions (>= 3.8.2)
- Kanject.Core.Queue.Provider.AwsSqs.Annotations (>= 3.8.1)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|---|---|
| 3.9.0 | 36 | 10/2/2026 |
| 3.8.2 | 37 | 10/1/2026 |
| 3.8.1 | 106 | 9/27/2026 |
| 3.8.0 | 82 | 9/27/2026 |
| 3.7.7 | 90 | 9/26/2026 |
| 3.7.6 | 118 | 9/7/2026 |
| 3.7.5 | 102 | 8/27/2026 |
| 3.7.4 | 106 | 8/22/2026 |
| 3.7.3 | 117 | 8/10/2026 |
| 3.7.2 | 111 | 8/9/2026 |
| 3.7.1 | 113 | 8/5/2026 |
| 3.7.0 | 111 | 8/5/2026 |
| 3.6.0 | 114 | 8/3/2026 |
| 3.5.9 | 128 | 7/30/2026 |
| 3.5.8 | 121 | 7/18/2026 |
| 3.5.7 | 140 | 7/13/2026 |
| 3.5.6 | 120 | 7/11/2026 |
| 3.5.5 | 128 | 7/11/2026 |
| 3.5.4 | 128 | 7/9/2026 |
| 3.5.3 | 124 | 7/9/2026 |