Skip to main content

async-lambda

Async-lambda is a Python framework for building scalable, event-driven AWS Lambda applications with first-class support for asynchronous invocation via SQS queues. It provides a high-level abstraction for orchestrating Lambda functions, managing event triggers, and handling complex workflows with minimal boilerplate.

async-lambda converts your application into a Serverless Application Model (SAM) template which can be deployed with the SAM cli tool or via Cloudformation.

Features

  • Async Task Abstraction: Define tasks as Lambda functions triggered by SQS, API Gateway, DynamoDB Streams, or scheduled events.
  • Automatic SAM Template Generation: Converts your application into a Serverless Application Model (SAM) template for deployment.
  • Lane-based Parallelism: Scale async tasks horizontally using multiple SQS queues (lanes) and control concurrency per lane.
  • Middleware Support: Register middleware to wrap task execution, enabling cross-cutting concerns (logging, auth, etc.).
  • Large Payload Handling: Offload large payloads to S3 automatically when exceeding SQS size limits.
  • Dead Letter Queues (DLQ): Robust error handling and message redrive for failed async invocations.
  • Configurable via Code and JSON: Set configuration at the app, stage, and task level using Python or config files.
  • Task Initialization: Run custom logic during Lambda INIT phase for caching, setup, or resource allocation.

Getting Started

Installation

Add async-lambda to your project (see requirements in pyproject.toml).

Basic Usage

from async_lambda import AsyncLambdaController, config_set_name, ScheduledEvent, ManagedSQSEvent

app = AsyncLambdaController()
config_set_name("project-name")
lambda_handler = app.async_lambda_handler  # Required Lambda export

@app.scheduled_task('ScheduledTask1', schedule_expression="rate(15 minutes)")
def scheduled_task_1(event: ScheduledEvent):
    app.async_invoke("AsyncTask1", payload={"foo": "bar"})

@app.async_task('AsyncTask1')
def async_task_2(event: ManagedSQSEvent):
    print(event.payload)  # {"foo": "bar"}

Packaging & Deployment

Use the async-lambda CLI to build and package your app:

async-lambda build app --stage <stage-name>

This generates a SAM template (template.json) and deployment bundle (deployment.zip). Deploy using AWS SAM CLI or CloudFormation.

Core Concepts

Controller & Tasks

  • AsyncLambdaController: Central orchestrator for registering tasks, managing middleware, and handling invocations.
  • Task: Each task is a Lambda function with a unique task_id and a trigger type (SQS, API, DynamoDB, schedule).

Task Decorators

All task decorators accept common configuration arguments:

  • memory: Memory allocation (MB)
  • timeout: Timeout (seconds)
  • ephemeral_storage: Ephemeral storage (MB)
  • maximum_concurrency: Max concurrency for SQS triggers (int or list per lane)
  • lane_count: Number of parallel lanes (for async tasks)
  • init_tasks: Functions to run during Lambda INIT phase

Example: Async Task

@app.async_task("TaskID")
def async_task(event: ManagedSQSEvent):
    print(event.payload)

It is quite easy to get into infinite looping situations when utilizing async-lambda and care should be taken.

INFINITE LOOP EXAMPLE

# If task_1 where to ever get invoked, then it would start an infinite loop with
# task 1 invoking task 2, task 2 invoking task 1, and repeat...

@app.async_task("Task1")
def task_1(event: ManagedSQSEvent):
    app.async_invoke("Task2", {})

@app.async_task("Task2")
def task_1(event: ManagedSQSEvent):
    app.async_invoke("Task1", {})

Example: Scheduled Task

@app.scheduled_task("TaskID", schedule_expression='rate(15 minutes)')
def scheduled_task(event: ScheduledEvent):
    ...

Example: API Task

@app.api_task("TaskID", path='/test', method='get')
def api_task(event: APIEvent):
    print(event.headers)
    print(event.querystring_params)
    print(event.body)

Example: Unmanaged SQS Task

@app.sqs_task("TaskID", queue_arn='queue-arn')
def sqs_task(event: UnmanagedSQSEvent):
    print(event.body)

Lanes & Concurrency

Async tasks can be scaled horizontally using lanes. Each lane is a separate SQS queue and can have its own concurrency limit. Lane assignment can be controlled at the controller, sub-controller, or task level.

app = AsyncLambdaController(lane_count=2)

@app.async_task("SwitchBoard")
def switch_board(event: ManagedSQSEvent):
    lane = 1 if event.payload['value'] > 50000 else 0
    app.async_invoke("ProcessingTask", event.payload, lane=lane)

@app.async_task("ProcessingTask", maximum_concurrency=[10, 2])
def processing_task(event: ManagedSQSEvent):
    ...

Middleware

Middleware functions wrap task execution and can be used for logging, authentication, or modifying events/responses.

def async_lambda_middleware(event, call_next):
    print(f"Invocation Payload: {event}")
    result = call_next(event)
    print(f"Invocation Result: {result}")
    return result

controller = AsyncLambdaController(middleware=[([BaseEvent], async_lambda_middleware)])

If there are multiple middleware functions then call_next will actually be calling the next middleware function in the stack.

For example if there is middleware functions A and B registered in that order. Then the execution order would go:

A(Pre) -> B(Pre) -> Task -> B(Post) -> A(Post)

Large Payloads

If a payload exceeds SQS size limits, async-lambda automatically stores it in S3 and passes a reference key to the Lambda function.

Dead Letter Queues (DLQ)

All async tasks share a DLQ for failed messages. You can configure custom DLQ tasks for advanced error handling.

async-lambda config

Configuration options can be set with the .async_lambda/config.json file. The configuration options can be set at the app, stage, and task level. A configuration option set will apply unless overridden at a more specific level (app -> stage -> task -> stage). The override logic attempts to be non-destructive so if you have a layers of ['layer_1'] at the app level, and [layer_2] at the stage level, then the value will be ['layer_1', 'layer_2'].

Config file levels schema

{
    # APP LEVEL
    "stages": {
        "stage_name": {
            # STAGE LEVEL
        }
    },
    "tasks": {
        "task_id": {
            # TASK LEVEL
            "stages": {
                "stage_name": {
                    # TASK STAGE LEVEL
                }
            }
        }
    }
}

At any of these levels any of the configuration options can be set: With the exception of domain_name, tls_version, and certificate_arn which can not be set at the task level.

environment_variables

{
    "ENV_VAR_NAME": "ENV_VAR_VALUE"
}

This config value will set environment variables for the function execution. These environment variables will also be available during build time.

The value is passed to the Environment property on SAM::Serverless::Function

policies

[
    'IAM_POLICY_ARN' | STATEMENT
]

Use this config option to attach any arbitrary policies to the lambda functions execution role.

The value is passed to the Policies property on SAM::Serverless::Function, in addition to the async-lambda created inline policies.

layers

[
    "LAYER_ARN"
]

Use this config option to add any arbitrary lambda layers to the lambda functions. Ordering matters, and merging is done thru concatenation.

The value is passed to the Layers property on SAM::Serverless::Function

subnet_ids

[
    "SUBNET_ID"
]

Use this config option to put the lambda function into a vpc/subnet.

The value is passed into the SubnetIds field of the VpcConfig property on SAM::Serverless::Function

security_group_ids

[
    "SECURITY_GROUP_ID"
]

Use this config option to attach a security group to the lambda function.

The value is passed into the SecurityGroupIds field of the VpcConfig property on SAM::Serverless::Function

managed_queue_extras

[
    {
        # Cloudformation resource
    }
]

Use this config option to add extra resources for managed SQS queues (async_task tasks.)

For example this might be used to attach alarms to these queues.

Each item in the list should be a complete cloudformation resource. async-lambda provides a few custom substitutions so that you can reference the extras and the associated managed sqs resource by LogicalId.

  • $QUEUEID will be replaced with the LogicalId of the associated Managed SQS queue.
  • $EXTRA<index> will be replaced with the LogicalId of the extra at the specified index.

method_settings

This config value can only be set at the app or stage level.

[
    {...}
]

If your async-lambda app contains any api_task tasks, then a AWS::Serverless::Api resource is created.

The value is passed into the MethodSettings property of the AWS::Serverless::Api. The spec for MethodSetting can be found here.

domain_name

This config value can only be set at the app or stage level.

"domain_name"

If your async-lambda app contains any api_task tasks, then a AWS::Serverless::Api resource is created.

This config value will set the DomainName field of the Domain property

tls_version

This config value can only be set at the app or stage level.

"tls_version"

If your async-lambda app contains any api_task tasks, then a AWS::Serverless::Api resource is created.

This config value will set the SecurityPolicy field of the Domain property

Possible values are TLS_1_0 and TLS_1_2

certificate_arn

This config value can only be set at the app or stage level.

"certificate_arn"

If your async-lambda app contains any api_task tasks, then a AWS::Serverless::Api resource is created.

This config value will set the CertificateArn field of the Domain property

hosted_zone_id

This config value can only be set at the app or stage level.

"hosted_zone_id"

If your async-lambda app contains any api_task tasks, then a AWS::Serverless::Api resource is created.

This config value will set the HostedZoneId field of the Route53 property of the Domain property.

This will create a DNS record on the given hosted zone for the api gateway endpoint created by the SAM deployment.

tags

{
    "TAG_NAME": "TAG_VALUE"
}

This config value will set the Tags field of all resources created by async-lambda. This will not set the field on managed_queue_extras resources.

The keys framework and framework-version will always be set and the system values will override any values set by the user.

For managed queues the tags async-lambda-queue-type will be set to dlq, dlq-task, or managed depending on the queue type.

For async_task queues (non dlq-task) the async-lambda-lane will be set.

logging_config

{
    "ApplicationLogLevel": "TRACE" | "DEBUG" | "INFO" | "WARN" | "ERROR" | "FATAL",
    "LogFormat": "Text" | "JSON",
    "LogGroup": "",
    "SystemLogLevel": "DEBUG" | "INFO" | "WARN"
}

The value is passed directly to the LoggingConfig Cloudformation parameter for lambda function/s.

See above for full details on configuration schema and merging behavior.

Advanced Usage

Task Initialization

Use the init_tasks argument to run setup logic during Lambda INIT phase. This is useful for caching, resource allocation, or one-time setup.

def setup_cache(task_id):
    ...

@app.async_task("TaskID", init_tasks=[setup_cache])
def async_task(event: ManagedSQSEvent):
    ...

Defer Utility

The Defer class allows you to cache values during INIT and only execute a function when its value is requested.

cache = Defer(get_a_value, 10, 100)

@app.async_task("Task", init_tasks=[cache.execute])
def task(event: ManagedSQSEvent):
    for i in range(cache.value):
        ...

Shared Function Role

By default each task declares Policies, so SAM generates one AWS::IAM::Role per task.

config_set_shared_function_role() emits a single AsyncLambdaSharedFunctionRole and points every task's Role at it instead:

from async_lambda import config_set_shared_function_role

config_set_shared_function_role()

The shared role carries everything the per-task roles carried, with one difference: SQS consume access is granted on the app's queue-name prefix ({name}-*) rather than each task's own queues. Unmanaged SQS queues are still listed individually, since they're named by their producer.

Two caveats:

  • A task with task-specific policies keeps its own generated role, since the shared role is built from app-wide and stage-wide config only.
  • Every app-wide policy must be either a managed policy ARN or a {"Statement": [...]} document. SAM policy templates can't be expressed in a plain IAM role and will fail the build.

Consolidated Queue-Age Alarms

managed_queue_extras attaches a copy of each extra resource to every managed queue, so using it for a CloudWatch alarm costs one stack resource per queue lane.

Setting queue_age_alarm in the build config emits a small number of grouped alarms instead, each watching ApproximateAgeOfOldestMessage across up to ten queues:

{
  "queue_age_alarm": {
    "threshold": 5000,
    "period": 86400,
    "alarm_actions": ["arn:aws:sns:us-east-1:123456789012:alarms"],
    "ok_actions": ["arn:aws:sns:us-east-1:123456789012:alarms"]
  }
}

Only threshold is required and period defaults to one day. Coverage is every managed queue lane plus the app DLQ. An app with 104 queues gets 11 alarms rather than 104.

CloudWatch rejects an alarm carrying more than 10 metrics.

Known Limitations

  • Not all Lambda configuration options are supported (see code for extension points)
  • Payloads must be JSON serializable
  • Infinite loops are possible if tasks invoke each other recursively

Project Structure

  • async_lambda/: Core framework code
    • controller.py: Main controller and orchestration logic
    • models/: Event, response, and task models
    • middleware.py: Middleware registration and execution
    • client.py, env.py, config.py: AWS clients and configuration
    • defer.py: Defer utility for caching
    • payload_encoder.py, util.py: Utilities
  • example/: Example usage and sample app
  • scripts/: Linting and testing scripts
  • test/: Unit tests

License

See LICENSE file for details.

Metadata

Release files for async-lambda-unstable 0.6.14

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for async-lambda-unstable 0.6.14
File Size Uploaded
async_lambda_unstable-0.6.14.tar.gz 46.1 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for async-lambda-unstable 0.6.14
File Interpreter ABI Platform
async_lambda_unstable-0.6.14-py2.py3-none-any.whl Python 3, Python 2 none any Details

Total release size: 100.9 kB

Release files / async_lambda_unstable-0.6.14.tar.gz

Download URL async_lambda_unstable-0.6.14.tar.gz
Size 46.1 kB
Tags Source
SHA-256 checksum
How to use checksums
f55d06ab09d4b25614e38577ae49b2959cc0f797c8ae3c55c823db5167e12723
BLAKE2b-256 checksum
How to use checksums
f009989ef3b32b130e7a94bde20f43f9f317fc2fdc1694acab45de5823b4cb3e
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via uv/0.11.33 {"installer":{"name":"uv","version":"0.11.33","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"macOS","version":null,"id":null,"libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}

Release files / async_lambda_unstable-0.6.14-py2.py3-none-any.whl

Download URL async_lambda_unstable-0.6.14-py2.py3-none-any.whl
Size 54.8 kB
Tags Python 2 Python 3
SHA-256 checksum
How to use checksums
d10081a6bde360b2bb4060a25ded4b7d6f18128cfa6facab215a26e53cdf2d2b
BLAKE2b-256 checksum
How to use checksums
ef304d69165fffc2eb328bbf825c2cc3fc80b2fe619d4fd598672983db2e6f6c
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via uv/0.11.33 {"installer":{"name":"uv","version":"0.11.33","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"macOS","version":null,"id":null,"libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}

Release history Release notifications | RSS feed

This release

0.6.14 This release

2 release files

0.6.13

2 release files

0.6.12

2 release files

0.6.11

2 release files

0.6.10

2 release files

0.6.9

2 release files

0.6.8

2 release files

0.6.7

2 release files

0.6.6

2 release files

0.6.5

2 release files

0.6.4

2 release files

0.6.3

2 release files

0.6.2

2 release files

0.6.1

2 release files

0.6.0

2 release files

0.5.9

2 release files

0.5.8

1 release file

0.5.7

2 release files

0.5.6

2 release files

0.5.5

2 release files

0.5.4

2 release files

0.5.3

2 release files

0.5.2

2 release files

0.5.1

2 release files

0.4.18

2 release files

0.4.16

2 release files

0.4.14

2 release files

0.4.13

2 release files

0.4.12

2 release files

0.4.11

2 release files

0.4.10

2 release files

0.4.9

2 release files

0.4.8

2 release files

0.4.7

2 release files

0.4.6

2 release files

0.4.5

2 release files

0.4.4

2 release files

0.4.3

2 release files

0.4.1

2 release files

0.3.12

2 release files

0.3.11

2 release files

0.3.10

2 release files

0.3.9

2 release files

0.3.8

2 release files

0.3.7

2 release files

0.3.6

2 release files

0.3.5

2 release files

0.3.4

2 release files

0.3.2

2 release files

0.3.1

2 release files

0.3.0

2 release files

0.2.22

2 release files

0.2.20

2 release files

0.2.19

2 release files

0.2.18

2 release files

0.2.17

2 release files

0.2.16

2 release files

0.2.14

2 release files

0.2.13

2 release files

0.2.12

2 release files

0.2.9

2 release files

0.2.8

2 release files

0.2.7

2 release files

0.2.6

2 release files

0.2.5

2 release files

0.2.4

2 release files

0.2.3

2 release files

0.2.2

2 release files

0.2.1

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page