Tag Archives: Advanced (300)

How AppFolio transformed its data streaming architecture with Amazon MSK Express brokers

Post Syndicated from Brandon Stanley original https://aws.amazon.com/blogs/big-data/how-appfolio-transformed-its-data-streaming-architecture-with-amazon-msk-express-brokers/

Real-time data streaming and event processing are critical components of modern distributed systems architectures. Apache Kafka has emerged as a leading platform for building real-time data pipelines and enabling asynchronous communication between microservices and applications. However, running and managing Kafka clusters at scale can be challenging, requiring specialized expertise and significant operational overhead.

Amazon Managed Streaming for Apache Kafka (Amazon MSK) is a fully managed service that you can use to build and run production Kafka applications. With Amazon MSK, you can rely on AWS to handle the heavy lifting of provisioning and managing Kafka clusters, while you focus on building innovative applications and real-time data processing pipelines.

In this post, you learn how AppFolio adopted Amazon MSK Express brokers to replace hours-long rebalances and manual storage planning with a streaming platform that scales automatically.

About AppFolio and its data streaming platform

AppFolio is a leading Real Estate Performance Management platform, serving thousands of property management companies across the United States. AppFolio’s platform processes millions of transactions daily, from rent collection and maintenance requests to lease management and financial reporting. In this data-intensive environment, reliable streaming infrastructure isn’t only important. It’s mission-critical.

At AppFolio, real-time data is the foundation of the company’s ability to deliver powerful, intelligent solutions that power the real estate industry. To achieve this level of performance, AppFolio engineered a modern streaming data architecture built on Amazon MSK with Express brokers. This infrastructure enables high-throughput, real-time applications at scale. With Amazon MSK Express brokers, AppFolio reliably ingests massive volumes of diverse data, including Change Data Capture (CDC) and server-side events, and makes it available to downstream consumers, such as real-time fraud detection, financial reporting, and automated property management workflows, within seconds of origin.

Previous architecture and AppFolio’s evolving requirements

Until early 2025, AppFolio ran their streaming platform on a single Amazon MSK cluster with Standard brokers, supporting both customer-facing and internal workloads. The architecture served them well through earlier growth phases. As AppFolio’s data platform evolved to support increasingly complex use cases and higher throughput, two characteristics of their workload led them to look for a more elastic streaming foundation.

AppFolio’s previous architecture: a single Amazon MSK cluster with Standard brokers serving both customer-facing and internal workloads

Figure 1: AppFolio’s previous architecture with Amazon MSK Standard brokers

First, AppFolio makes extensive use of log-compacted topics for their CDC streams. Compacted topics retain the latest value for each key indefinitely, which is exactly what they want for streams that mirror the state of operational tables. As their footprint grew, AppFolio wanted an infrastructure model that could scale storage automatically alongside data growth, without ongoing capacity planning that took multiple hours every month.

Second, AppFolio’s throughput continued to grow as they onboarded new use cases and added more event sources. They wanted the ability to scale the cluster quickly in response to traffic shifts, with minimal lead time for partition reassignments.

Third, as AppFolio’s platform matured, they needed workload isolation between customer-facing and internal data flows. Running customer-facing and internal workloads on a single cluster made it harder to size and tune each independently. As both grew, AppFolio wanted dedicated resources so each could be sized and tuned independently.

Based on these needs, AppFolio identified the following key requirements for their next-generation streaming platform:

  1. Elastic, automatically managed storage that scales with AppFolio compaction-heavy CDC workloads, removing the need for upfront broker capacity planning.
  2. Faster horizontal scaling and partition reassignment so AppFolio can adjust cluster shape in response to actual traffic in minutes rather than hours.
  3. Workload isolation between customer-facing and internal data flows, so each workload can be sized and tuned for its own traffic pattern.

Why AppFolio chose Amazon MSK Express brokers

After evaluating their options, AppFolio chose Amazon MSK Express brokers as the foundation for their next-generation streaming platform. Express brokers are a broker type offered under MSK Provisioned. They include pay-as-you-go elastic storage that scales automatically, intelligent partition rebalancing, and Kafka configuration defaults tuned for production workloads. Express brokers mapped directly to the requirements AppFolio identified:

  1. Elastic storage that scales with their data. Express brokers remove broker disk sizing and provisioning, with storage scaling automatically alongside data growth. AppFolio pays only for the storage actually used.
  2. AWS benchmarks showed up to 20 times faster scaling. Horizontal scaling and partition reassignment that previously took hours now complete in minutes, letting AppFolio react to traffic shifts on a much shorter cycle.
  3. Production-tuned defaults. Express brokers come pre-configured with Kafka best-practice defaults and built-in client throughput quotas, simplifying AppFolio’s operational model.
  4. Full Kafka API compatibility. AppFolio was able to migrate without changes to its producer and consumer applications.

As part of the migration, AppFolio also took the opportunity to rethink how the cluster was being used. Rather than recreating a single shared cluster on Express brokers, they segmented their MSK clusters by workload type. This gives customer-facing and internal workloads dedicated resources, providing better isolation and more predictable performance for each workload class.

Current architecture

AppFolio’s current architecture consists of multiple Amazon MSK clusters with Express brokers, segmented by workload type. Each cluster is sized and tuned for its specific traffic pattern, providing improved isolation and more predictable performance. The following diagram shows the deployment.

Current architecture: multiple Amazon MSK clusters with Express brokers, segmented by workload type into customer-facing and internal clusters

Figure 2: Current architecture with workload-segmented Amazon MSK clusters using Express brokers

Benefits achieved

By migrating to Amazon MSK Express brokers and adopting a workload-segmented cluster design, AppFolio has realized several key benefits:

Elastic, hands-off storage

The pay-as-you-go storage of Express brokers scales automatically with AppFolio’s data growth. Storage capacity is no longer something the platform team plans, provisions, or monitors, and AppFolio pays only for what they use. For a workload that runs heavily on compacted topics, this is the single largest operational improvement they have seen.

Faster scaling

Partition reassignment and broker scaling that previously took hours now complete in minutes, enabling AppFolio to adjust cluster shape in response to actual traffic instead of running ahead of forecasts.

Improved workload isolation

Splitting their streaming traffic into workload-segmented clusters has given AppFolio more predictable performance. Customer-facing and internal workloads now run on dedicated infrastructure, and each cluster can be sized and tuned for its own traffic pattern.

Stable environment as data volumes grow

Since the migration, AppFolio has maintained a stable environment with no significant downtime, even as data volumes continue to grow.

Reduced operational overhead

Hands-off storage management and intelligent rebalancing have removed several recurring tasks from the AppFolio platform team’s queue, including the constant monitoring and manual intervention that storage planning required under their previous architecture.

Conclusion

By using Amazon MSK Express brokers and adopting a workload-segmented cluster design, AppFolio has built a streaming foundation that scales elastically with their data growth and adapts quickly to changes in traffic. The pay-as-you-go storage and faster scaling of Express brokers let AppFolio’s platform team focus engineering effort on building new capabilities for customers, rather than on Kafka capacity planning. As AppFolio continues to expand its platform for the real estate industry, the Amazon MSK Express brokers infrastructure provides a scalable foundation for future growth.

To learn more about Express brokers for Amazon MSK, see the Express brokers for Amazon MSK documentation and the AWS announcement post Introducing Express brokers for Amazon MSK.


About the authors

Brandon Stanley

Brandon Stanley

Brandon is a Staff Data Engineer at AppFolio, responsible for architecting, building, and evolving AppFolio’s near real-time data platform, which captures, ingests, and serves database change logs, custom server-side events, and clickstream events from customer databases across product domains to targets including data warehouses, OLTP databases, and data lakehouses.

Devarsh Patel

Devarsh Patel

Devarsh is a Data Engineer at AppFolio, where he builds and operates large-scale, production-grade streaming data infrastructure that powers real-time analytics across the organization. His areas of focus include change data capture (CDC) pipelines, Apache Flink, Snowflake, and AWS infrastructure automation using Terraform and Kubernetes.

Ryan D’Souza

Ryan D’Souza

Ryan is a Staff Data Engineer at AppFolio. He architects, builds, and scales the data platform powering AppFolio’s AI solutions, customer-facing applications, and product analytics. He specializes in streaming data pipelines and data lakehouse architectures on AWS.

Aarjvi Desai

Aarjvi Desai

Aarjvi is a Sr Technical Account Manager and container specialist at AWS, based in the San Francisco Bay Area. She helps customers solve cloud challenges and build scalable, reliable solutions for generative AI workloads. Her expertise spans Kubernetes architecture, GPU accelerated workloads, and helping enterprises navigate AI infrastructure at scale.

Kalyan Janaki

Kalyan Janaki

Kalyan is Senior Big Data & Analytics Specialist at AWS. He helps customers architect and build highly scalable, performant, and secure cloud-based solutions on AWS.

Shilpa Bondale

Shilpa Bondale

Shilpa is a Senior Solutions Architect at AWS, based in the San Francisco Bay Area. She partners with companies to solve complex engineering challenges across databases, analytics, machine learning, and AI. She helps customers architect scalable, production-grade solutions, from real-time data pipelines to large-scale ML inference – using the breadth of AWS services.

Recovery strategies to meet data residency requirements

Post Syndicated from Jamie Pasterick original https://aws.amazon.com/blogs/architecture/recovery-strategies-to-meet-data-residency-requirements/

Data residency requirements can affect how government agencies, regulated industries such as financial services, healthcare, and power and utilities, and businesses that make residency commitments plan for the recovery of their critical workloads on AWS. These requirements must be considered and balanced against applicable workload recovery objectives. This post assumes familiarity with AWS Regions, disaster recovery concepts, and AWS encryption services.

Where residency requirements are scoped at the national level, AWS provides multiple Regions within the same country in the United States, Canada, Australia, India, Japan, Germany, and China (operated by Sinnet and NWCD). In some cases, residency requirements can span national borders. For example, AWS offers multiple Regions within the European Union (EU), including the AWS European Sovereign Cloud, giving EU customers options for hosting and recovering workloads across member states where pan-national regulations treat the EU as a unified jurisdiction for data protection. This allows customers to use multi-Region recovery architectures and maintain data residency.

Where residency requirements are scoped to a country or countries served by a single AWS Region, or you need to address failure scenarios not fully mitigated by multiple Regions in the same country, alternate strategies can help you achieve your recovery requirements. In this post, we present three strategies that customers can use in close collaboration with their regulators to achieve bounded recovery while addressing data residency requirements. These strategies range from encryption-based compensating controls on multi-Region replication to fully in-country architectures.

Recovery strategies

We present three strategies for backing up critical business data (source code and data you cannot reproduce from other sources) and launching recovery infrastructure from those backups at a location distinct from your primary Region. Each strategy represents a different set of constraints on where data and administrative operations can reside. You should select the strategy that best matches your requirements and risk appetite, then evaluate the options within that strategy in collaboration with your regulator. The most important factor of success for whichever strategy you choose is your ability to test it end-to-end, continuously, to build and maintain confidence it will work when required.

Strategy 1: Cryptographic boundary

This approach replicates data into another AWS Region in a geopolitically aligned country with compatible data protection frameworks using encryption as a compensating technical control. Customers use AWS Key Management Service (AWS KMS) keys to encrypt the data, so that no one, including AWS operators, can access the data without the customer-controlled data encryption keys. This approach supports replication using features like Amazon Simple Storage Service (Amazon S3) Cross-Region Replication (CRR) or AWS Backup cross-Region copy.

With AWS KMS, you can also use key policies to explicitly deny decryption operations in the recovery Region. This provides you with strong assurance that your data in the recovery Region cannot be decrypted under any circumstances until you modify the key policy. You can work with your regulators to determine when to update these key policies as part of your recovery process.

Server-side encryption architecture using AWS KMS with cross-Region replication

Figure 1 – Server-side encryption architecture using AWS KMS with cross-Region replication.

You can also use this approach with client-side encryption for backups you manage and store in S3. Customers manage their own backup processes and use the
AWS Encryption SDK or their own clients with
multi-Region AWS KMS keys to encrypt the data. Multi-Region keys allow you to encrypt and decrypt replicated S3 data in multiple Regions using the same key material.

Client-side encryption approach using multi-Region AWS KMS keys

Figure 2 – Client-side encryption approach using multi-Region AWS KMS keys.

You can explore additional strategies for enhancing controls on the key policies to meet your requirements. For example, you can require Multi-Factor Authentication (MFA) to update your key policies. This allows the MFA holder and credential holder to be two distinct parties. They could be two different teams within an organization, or you could consider greater separation by providing the MFA device to a trusted third party such as a regulator. Another option is to implement controls in your identity provider to issue specific IAM session tags that provide conditional access to update the key policy. You should also scope key policy update permissions to a set of named IAM principals using condition keys, so that only explicitly authorized identities can modify decryption access.

Choose this strategy when storing encrypted data in a partner country is acceptable. Key policies prevent unauthorized access, and this approach offers the simplest operational model with the lowest recovery time.

Strategy 2: Data boundary

In this strategy, you store backups or operate a pilot light recovery environment on AWS Outposts in an on-premises site within the source country or other approved location. You replicate data from your primary Region using tools such as AWS DataSync for S3 data or other replication tools like MySQL binlog or Postgres logical replication for Amazon Relational Database Service (Amazon RDS) instances. You maintain full control of where your business data physically resides at all times. You can also use third-party backup solutions to replicate data from your primary Region to on-premises or in-country storage.

AWS Outposts recovery architecture with data replicated to on-premises infrastructure

Figure 3 – AWS Outposts recovery architecture with data replicated to on-premises infrastructure.

Data access occurs directly over the local network in your on-premises facility through the
data planes of the resources hosted on the Outposts infrastructure. These are resources like
Amazon Elastic Compute Cloud (Amazon EC2) instances, Amazon RDS database instances, and S3 buckets. You provision and configure those resources through each service’s
control plane, which is hosted in the parent Region you select when you order your Outposts racks.

Choose a parent Region that is different from your primary Region. This prevents simultaneous impact to your primary workloads and your ability to use control plane operations for recovery, such as launching new instances on your Outposts. Note that Outposts are not designed for disconnected operations or environments with limited to no connectivity. Maintain highly available networking connections from your on-premises site back to the parent AWS Region.

During recovery, you can restore your environment directly on the Outposts infrastructure. AWS Outposts support a subset of the available AWS services in a Region, so you need to design your workloads to use the services available. Alternatively, you may obtain regulator agreement to restore your environment from backups stored on your Outposts to an AWS Region in a different country during extreme circumstances.

Select this strategy when you must physically maintain your backups and data in specific locations, want a consistent experience using AWS services and hardware in the cloud and on-premises, and using control planes for your Outposts infrastructure from outside your primary Region is acceptable.

Strategy 3: Strict local autonomy boundary

Some data residency frameworks, such as those in financial services or national security contexts, may require customer data and the control plane systems used to manage that data and recovery environments remain within national borders. Two options achieve this outcome.

Option 1: On-premises infrastructure

In this option, you operate hardware and software in an on-premises location to store backup data. Like the Outposts option, this provides flexibility on where backed-up data is restored: it could be restored to on-premises physical hardware or to an AWS Region in a different country.

On-premises backup architecture with data copied from Amazon S3 to local storage

Figure 4 – On-premises backup architecture with data copied from Amazon S3 to local storage.

This solution requires you to self-manage backups in S3, then copy them to on-premises storage using tools like AWS DataSync. You need to decide how to manage encryption of your data on-premises. Encrypting and decrypting data on-premises should not depend on the availability of the primary Region.

Option 2: Multi-cloud

In a multi-cloud solution, you can replicate backups from your primary cloud to an environment on another cloud provider and recover your workloads there. You can also use a lifeboat strategy.

A lifeboat strategy involves running a separate set of systems in another cloud provider which meets your residency requirements. These systems are not replicas of the primary platform. They are built and developed natively to provide a subset of critical functionality. This approach avoids architecting to the least common denominator of services across providers. You can take full advantage of all services available on AWS for your primary workload while the stand-in system uses a separate, purpose-built architecture.

Multi-cloud lifeboat architecture with a purpose-built stand-in system

Figure 5 – Multi-cloud lifeboat architecture with a purpose-built stand-in system.


Monzo Bank’s stand-in system is a well-documented example of this pattern. Monzo operates its primary banking system on AWS with thousands of microservices. Rather than replicating that entire stack, they built a small set of purpose-built services on a separate cloud provider that supports only the operations most critical to their customer experience: card payments, bank transfers, and balance information. According to Monzo, the stand-in shares no code and no infrastructure with the primary system.

The lifeboat pattern follows several key design principles. The stand-in supports only key functionality, which helps minimize the cost of the solution. Using different software reduces the probability that the same defect or failure mode affects both systems simultaneously. The stand-in accepts eventually consistent data, which avoids strong coupling between the two environments and preserves availability independence. The recovery does not appear transparent to end users. The experience is intentionally degraded to a subset of services, which is an explicit trade-off for maintaining availability.

A multi-cloud lifeboat provides protection against service disruptions of an AWS Region in a single country, but using multi-cloud to keep data in-country may not fully mitigate scenarios where all major cloud providers in a geography face simultaneous disruption.

Testing

You must continuously test your recovery strategy to build and maintain confidence it will succeed during a real event. Testing must be performed end-to-end: validating the integrity and consistency of backups, launching compute and database resources, and running synthetic test traffic through the recovered system. The more frequently you test, the more recovery becomes a standard operational process rather than a one-off monthly or quarterly activity. Your testing must keep up with the rate of change in your environment. If a change breaks your recovery process, you want to know about it and fix the procedures as quickly as possible.

The approach to testing is generally the same as traditional multi-Region recovery testing, but the additional complexities and operational processes of these solutions require an increased level of rigor:

  • Strategy 1 (Cryptographic boundary): You need to decide if decrypting data and launching recovery environments with that data is acceptable. Ideally, you run full end-to-end recovery tests during approved test windows. If you cannot, you need to test your recovery procedures using synthetic data that is not subject to data residency requirements. This helps validate your recovery procedures, but it does not prove you can recover your critical business data.
  • Strategy 2 (Data boundary): Validate replication lag meets your recovery requirements. Test launching recovery workloads on Outposts infrastructure and confirm that the parent Region control plane can orchestrate recovery while the primary Region is simulated as unavailable.
  • Strategy 3 (Strict local autonomy boundary): For on-premises recovery, validate that backups can be restored to your target environment and that workloads function correctly outside of AWS. For multi-cloud lifeboats, test failover activation and verify that the lifeboat has no dependencies on your primary environment.

Summary

In this post, we presented three strategies for implementing disaster recovery solutions that support data residency requirements:

  • Strategy 1 (Cryptographic boundary) uses encryption to meet the intent of data residency while using multi-Region AWS infrastructure for recovery. This offers the lowest operational complexity and best recovery performance, if permitted by applicable regulations.
  • Strategy 2 (Data boundary) uses AWS Outposts to maintain data in-country in customer-controlled facilities while using an out-of-country control plane for management operations, if permitted by applicable regulations.
  • Strategy 3 (Strict local autonomy boundary) keeps both data and control in-country during a recovery event using on-premises infrastructure or a multi-cloud lifeboat strategy. This offers the highest degree of control but with the greatest operational complexity.

The option that works best will be a joint decision between your business, regulators, and your customers. You should consider potential failures proactively and build recovery plans before an event occurs. This framework provides a structured way to evaluate the trade-offs and determine where to invest based on your regulatory environment, risk appetite, and operational capabilities.

Next steps

To learn more about the services and approaches discussed in this post, see the following resources:

If you have questions about applying these strategies to your specific workloads and regulatory environment, reach out to your AWS account team to discuss your recovery requirements in detail.


About the authors

Burst to Region: Overflow AWS Outposts workloads to Amazon EC2

Post Syndicated from Diya . original https://aws.amazon.com/blogs/compute/burst-to-region-overflow-aws-outposts-workloads-to-amazon-ec2/

AWS Outposts brings AWS infrastructure into your data center, giving on-premises workloads the low latency and data locality they need. But unlike the AWS Region, an Outposts rack has a fixed amount of compute. When your workload needs more instances than the rack can provide, you have two options: drop requests, or overflow them somewhere with room to grow. This post shows you how to automate the second option. You build a Burst to Region pattern that detects capacity constraints on your Outpost, launches Amazon Elastic Compute Cloud (Amazon EC2) instances in the parent Region, gradually shifts traffic to them, and returns traffic to local instances once capacity recovers.

To implement this pattern you configure Amazon CloudWatch, Amazon Simple Notification Service (Amazon SNS), AWS Lambda, Amazon EC2 Auto Scaling, Elastic Load Balancing (Application Load Balancer), and Amazon EventBridge. You trade a moderate latency increase for continued availability during capacity events.

When to use this pattern

This pattern assumes your Outposts workload scales out through Amazon EC2 Auto Scaling. Burst to Region reacts to instance-capacity exhaustion on the rack. It engages when your workload tries to launch more instances than the available Outpost capacity supports. If your fleet is fixed size and degrades under load without scaling out, the capacity alarm never fires and overflow never triggers. For those workloads, monitor per-instance saturation (CPU, latency) separately.

Good candidates prefer local capacity but can tolerate Region latency under pressure. If your application runs on Outposts for proximity yet degrades gracefully when some traffic takes the longer path to the Region, it fits this pattern. Examples include:

  • Internal enterprise applications.
  • Stateless web frontends and API layers.
  • Pre-processing tiers where single-digit to tens-of-milliseconds additional round-trip latency during peaks is acceptable.

Poor candidates cannot absorb any added latency or must stay on the Outpost. Avoid this pattern for:

  • Applications with sub-millisecond requirements.
  • Workloads with strict data residency or sovereignty mandates that prevent traffic from leaving the on-premises environment.
  • Real-time control systems with hard timing constraints.
  • Applications tightly coupled to on-premises data stores with no Region replica.

The core tradeoff is explicit. During capacity events, you accept moderately higher latency to maintain availability. If your workload cannot tolerate any latency increase, keep it pinned to Outposts and reserve capacity through other means, such as Capacity Reservations.

Solution overview

Burst to Region works in three moves: detect capacity pressure on the Outpost, launch overflow compute in the parent Region, and shift traffic gradually until local capacity recovers. Six AWS services coordinate to make this automatic. The following diagram shows the reference architecture for the Burst to Region pattern, illustrating how the six AWS services interact during capacity detection, overflow scaling, traffic distribution, and recovery.

Reference architecture for Burst to Region on AWS Outposts showing capacity detection, overflow scaling, traffic distribution, and recovery

Figure 1: Reference architecture for Burst to Region on AWS Outposts

The pattern uses six AWS services working together:

  • Amazon CloudWatch monitors Outposts capacity utilization and raises alarms.
  • Amazon SNS provides event fan-out from alarm to orchestrator.
  • AWS Lambda orchestrates the burst logic (scale-out, weight adjustment, recovery)
  • Amazon EC2 Auto Scaling manages the overflow fleet lifecycle.
  • Application Load Balancer distributes traffic across both locations using weighted target groups.
  • Amazon EventBridge handles periodic recovery evaluation.

You must configure five phases for this pattern:

  1. Monitor. CloudWatch tracks Outposts capacity utilization metrics in the AWS/Outposts namespace.
  2. Detect. A CloudWatch alarm fires when utilization exceeds a threshold (for example, 80%).
  3. Overflow. The alarm triggers a Lambda function through Amazon SNS. Lambda scales out a Region-based Amazon EC2 Auto Scaling group and adjusts ALB target group weights.
  4. Distribute. The ALB splits traffic between Outposts instances and Region instances using weighted forwarding.
  5. Recover. An Amazon EventBridge scheduled rule periodically evaluates capacity. When Outposts recovers, Lambda scales down the overflow fleet and returns all traffic to local instances.

Design decisions

We chose Application Load Balancer with weighted forwarding over Amazon Route 53 weighted routing for traffic distribution. ALB provides health-aware routing to only healthy overflow instances and target group stickiness for session consistency. Weight changes take effect for new connections after calling the ModifyRule API. DNS-based shifting through Route 53 provides too coarse control for rapid weight adjustments, and TTL propagation delays make recovery slower.

The burst orchestrator runs as a Lambda function rather than a long-running service. It executes only during state transitions, so there is no steady-state compute cost. Lambda integrates natively with Amazon SNS and Amazon EventBridge for event-driven invocation without additional infrastructure.

You implement recovery with an Amazon EventBridge scheduled rule (every 5 minutes) rather than relying solely on the CloudWatch alarm to return to OK state. The alarm confirms capacity is available, but does not confirm that overflow instances have drained active connections. The scheduled rule provides gradual, safe scale-down.

Implementation

This section walks through the key components of the Burst to Region pattern. For the complete deployable AWS SAM template, see the GitHub repository.

Prerequisites

To deploy this pattern, you need:

  • An AWS account with a configured AWS Outposts rack.
  • An Amazon Virtual Private Cloud (Amazon VPC) with subnets associated with your Outposts and subnets in the parent AWS Region.
  • IAM permissions to create CloudWatch alarms, Lambda functions, Auto Scaling groups, and ALB resources.
  • AWS Serverless Application Model (AWS SAM) CLI installed and configured.
  • Existing Amazon EC2 Auto Scaling group running on your Outpost (these become your baseline fleet)
  • A custom domain name with a DNS record (Route 53 alias or CNAME) pointing to your Application Load Balancer, and an AWS Certificate Manager (ACM) certificate for that domain to enable HTTPS.

Capacity monitoring and alarm

The CloudWatch alarm monitors instance utilization on the Outpost and triggers the burst workflow when capacity is constrained.

The InstanceTypeCapacityUtilization metric reports the percentage of a given instance type’s capacity in use. Note that this metric includes capacity consumed by managed services such as Amazon Relational Database Service (Amazon RDS) or Application Load Balancer running on the Outpost — not only your application’s EC2 instances. Factor this into your threshold planning.

OutpostsCapacityAlarm:
  Type: AWS::CloudWatch::Alarm
  Properties:
    AlarmName: outposts-capacity-high
    Namespace: AWS/Outposts
    MetricName: InstanceTypeCapacityUtilization
    Dimensions:
      - Name: OutpostId
        Value: !Ref OutpostId
      - Name: InstanceType
        Value: !Ref OutpostInstanceType
    Statistic: Average
    Period: 300
    EvaluationPeriods: 2
    Threshold: !Ref CapacityThreshold
    ComparisonOperator: GreaterThanOrEqualToThreshold
    AlarmActions:
      - !Ref BurstSNSTopic
    TreatMissingData: notBreaching

Why these values matter:

  • Period: 300 and EvaluationPeriods: 2 require 10 minutes of sustained high utilization before triggering. This avoids false alarms from transient spikes.
  • Threshold: 80 (recommended starting point) leaves a 20% buffer. A threshold set too high (95%) risks launch failures before the overflow fleet is ready. A threshold set too low (50%) causes unnecessary bursts.
  • TreatMissingData: notBreaching prevents false alarms when data points are missing. Since this alarm is scoped to a single instance type, treating missing data as breaching could trigger unnecessary bursts when the instance type is simply not in use.
  • Separate scale-out from scale-in: This alarm triggers burst scale-out at 80%. Recovery is handled separately by the Amazon EventBridge scheduled rule, which uses a lower threshold (for example, 60%) before scaling in. This hysteresis gap prevents flapping where scaling down immediately pushes utilization back above the alarm threshold.

Burst orchestrator (Lambda)

The Lambda function handles two event paths: alarm-triggered scale-out and scheduled recovery evaluation. The following pseudocode shows the orchestration flow:

def handler(event, context):
    # Route based on event source
    if is_scheduled_recovery(event):
        return handle_recovery_check()

    alarm_state = parse_sns_alarm_state(event)

    if alarm_state == 'ALARM':
        # Scale out the overflow Auto Scaling group
        scale_out_overflow(desired=OVERFLOW_CAPACITY)
        # Don't shift traffic yet --- wait for healthy instances
        publish_burst_metric(active=True)


def handle_recovery_check():
    """Called every 5 minutes by EventBridge."""
    # Check if burst is active
    if not is_burst_active():
        return

    # If overflow instances are healthy and registered, shift traffic
    if overflow_targets_healthy():
        current_weights = get_current_alb_weights()
        if current_weights['region'] == 0:
            # First shift --- instances are now warm
            set_alb_weights(outposts=90, region=10)
        elif needs_more_overflow():
            step_up_region_weight()

    # If Outposts capacity has recovered, begin scale-down
    if outposts_capacity_recovered():
        step_down_region_weight()
        if get_current_alb_weights()['region'] == 0:
            # All traffic back to Outposts, drain and terminate overflow
            wait_for_connection_draining()
            scale_down_overflow(desired=0)
            publish_burst_metric(active=False)

The key actions the function performs:

  • scale_out_overflow — Sets the overflow Auto Scaling group desired capacity from 0 to your configured burst size.
  • set_alb_weights — Calls the ModifyListener API to adjust weighted forwarding between the Outposts and Region target groups.
  • publish_burst_metric — Writes a custom CloudWatch metric (BurstActive) for dashboard visibility.
  • handle_recovery_check — Called every 5 minutes by Amazon EventBridge. Confirms Outposts capacity has recovered, steps weights back gradually, waits for connection draining, then scales down the overflow fleet.

Important: The orchestrator does not shift ALB weights immediately upon scale-out. It waits for the next Amazon EventBridge invocation (up to 5 minutes) to confirm that overflow instances have passed health checks and are registered as healthy in the target group. This helps prevent routing traffic to instances that have not finished launching.

For the production-ready implementation with error handling, gradual weight stepping, and connection draining verification, see the GitHub repository.

Overflow Auto Scaling group

The overflow fleet starts at zero and scales only when the Lambda function sets desired capacity during a burst event:

OverflowASG:
  Type: AWS::AutoScaling::AutoScalingGroup
  Properties:
    AutoScalingGroupName: burst-overflow-fleet
    LaunchTemplate:
      LaunchTemplateId: !Ref OverflowLaunchTemplate
      Version: !GetAtt OverflowLaunchTemplate.LatestVersionNumber
    MinSize: 0
    MaxSize: !Ref MaxOverflowCapacity
    DesiredCapacity: 0
    VPCZoneIdentifier:
      - !Ref RegionSubnet1
      - !Ref RegionSubnet2
    TargetGroupARNs:
      - !Ref RegionTargetGroup
    HealthCheckType: ELB
    HealthCheckGracePeriod: 120
    MetricsCollection:
      - Granularity: 1Minute

The overflow fleet starts at zero capacity and incurs no cost at rest. During a burst event, the Lambda function calls the SetDesiredCapacity API to launch overflow instances. During recovery, it sets desired capacity back to zero.

The launch template mirrors your Outposts instance type to maintain consistent performance characteristics across both locations.

ALB weighted forwarding

The ALB listener uses weighted forwarding across two target groups. In steady state, all traffic goes to Outposts (weight 100/0). During burst, the Lambda function adjusts these weights dynamically using the ModifyListener API. Clients reach the ALB through a DNS record — either a Route 53 alias or a CNAME pointing to the ALB’s DNS name.

ALBListener:
  Type: AWS::ElasticLoadBalancingV2::Listener
  Properties:
    LoadBalancerArn: !Ref ApplicationLoadBalancer
    Port: 443
    Protocol: HTTPS
    SslPolicy: ELBSecurityPolicy-TLS13-1-2-2021-06
    Certificates:
      - CertificateArn: !Ref CertificateArn
    DefaultAction:
      Type: forward
      ForwardConfig:
        TargetGroups:
          - TargetGroupArn: !Ref OutpostsTargetGroup
            Weight: 100
          - TargetGroupArn: !Ref RegionTargetGroup
            Weight: 0
        TargetGroupStickinessConfig:
          Enabled: true
          DurationSeconds: 300

RegionTargetGroup:
  Type: AWS::ElasticLoadBalancingV2::TargetGroup
  Properties:
    Name: burst-region-targets
    Protocol: HTTP
    Port: 80
    VpcId: !Ref VpcId
    HealthCheckEnabled: true
    HealthCheckIntervalSeconds: 30
    HealthCheckPath: /health
    HealthyThresholdCount: 2
    UnhealthyThresholdCount: 3
    TargetGroupAttributes:
      - Key: deregistration_delay.timeout_seconds
        Value: "300"
      - Key: slow_start.duration_seconds
        Value: "120"

Note on stickiness: Target group stickiness keeps a client pinned to whichever target group served its first request for DurationSeconds. We set this to 300 seconds (5 minutes) to match the Amazon EventBridge evaluation interval. This balances session consistency for stateful workloads against the need for weight changes to take effect within a reasonable window. For purely stateless workloads, you can disable stickiness entirely to allow immediate weight convergence. For workloads requiring longer session affinity, increase the duration but understand that weight transitions will converge more slowly — existing sticky sessions continue going to the original target group until they expire.

Traffic weight progression

Use stepped transitions rather than abrupt weight changes. The following table shows the recommended progression:

Phase Outposts weight Region weight Condition to advance
Normal 100 0 Steady state
Burst step 1 90 10 Region target group has at least 1 healthy host
Burst step 2 70 30 Region target group healthy for 2 consecutive checks
Burst step 3 50 50 Only if Outposts capacity exceeds 95% used
Recovery step 1 80 20 Outposts capacity below 70%
Recovery step 2 100 0 Outposts capacity below 60% for 2 checks

Avoid jumping directly from 0% to 50% Region traffic. Cold overflow instances need time to warm caches and stabilize before absorbing significant load.

Best practices

Apply these best practices to get the most from this pattern while avoiding common pitfalls.

Traffic tiering

Classify your workloads into two tiers at the ALB listener level. Latency-critical paths use routing rules with the Outposts target group only. These never overflow regardless of capacity state. Overflow-eligible paths use the weighted forwarding rule. This separation helps make sure that your most latency-sensitive flows are not impacted by the burst mechanism.

Managing data gravity

For stateless workloads, Burst to Region requires no special data handling. For workloads with session state or shared data:

Anti-pattern: Do not burst workloads that write to Outposts-local storage and expect synchronous consistency. The latency and complexity of cross-location writes defeats the purpose of the pattern.

Cost optimization

The overflow fleet consumes On-Demand pricing by default since it starts at zero and scales only during peaks.

Burst profile Recommended pricing Rationale
Unpredictable spikes (minutes) On-Demand Maximum flexibility, no commitment waste
Predictable daily peaks (hours) Savings Plans (Compute) Covers overflow hours at discount
Frequent, long bursts Reserved capacity plus On-Demand Baseline discount plus burst flexibility

Monitor your BurstActive custom metric over time. If overflow is active more than 30% of the time, you likely need additional Outposts capacity rather than relying on Region overflow.

Security consistency

Maintain identical security posture across both environments:

  • Use the same security group rules for Outposts and Region instances.
  • Deploy with AWS CloudFormation StackSets to support consistency.
  • Share the same IAM instance profile. The overflow launch template references the same role as your Outposts instances.
  • Apply the same AWS Systems Manager patch baselines and compliance rules to both fleets.

Observability

Build a CloudWatch dashboard that provides visibility into burst state and performance. The SAM template in the repository deploys a pre-configured dashboard tracking:

  • Burst status: Custom BurstActive metric (1 = active, 0 = normal)
  • Capacity headroom: UsedInstanceType_Count compared to AvailableInstanceType_Count. Note that UsedInstanceType_Count includes instances consumed by managed services (Amazon RDS, ALB), so your available application capacity may be lower than the raw availability count suggests.
  • Overflow fleet size: Auto Scaling group GroupInServiceInstances.
  • Latency comparison: TargetResponseTime per target group (Outposts compared to Region)
  • Traffic distribution: RequestCount per target group.

Set a CloudWatch alarm on Region target group TargetResponseTime exceeding your acceptable threshold. This provides early warning if overflow latency degrades beyond your tolerance.

Because the ALB resides in the Region, all traffic to Outposts targets traverses the service link. Keep the following in mind:

Bandwidth planning: Steady-state traffic to Outposts targets flows over the service link. Verify that your connection meets the minimum 500 Mbps per compute rack recommended by AWS, with sufficient headroom for both application traffic and Outposts control plane communication. Monitor service link VIF throughput using IfTrafficIn and IfTrafficOut metrics (on service link VIFs) to detect saturation before it impacts performance.

Latency impact: The service link adds latency compared to a locally deployed load balancer. The exact impact depends on your service link connection type and distance to the parent Region (AWS specifies a maximum of 175 ms round-trip for service link). For internet-facing workloads, this is typically negligible relative to the client-to-Region round trip. For workloads serving on-premises users through the Local Gateway, consider Route 53 weighted routing between an ALB on Outposts and a separate ALB in the Region instead.

Connection draining: When scaling down the overflow fleet, allow sufficient time for in-flight requests to complete. The deregistration delay configured on the target group (default 300 seconds) and the Auto Scaling scale-in cool-down period work together to help provide graceful termination and minimize the risk of dropping active connections.

Failure modes: If the service link goes down, the ALB cannot reach Outposts targets. Health checks fail, and all traffic automatically shifts to Region targets. This provides an unintentional but useful failover behavior. However, note that the overflow fleet is sized for burst capacity, not for sustaining 100% of production traffic. Monitor the ConnectedStatus metric (under the AWS/Outposts namespace, dimension OutpostId) and alert on degradation. If you need full failover capability, architect a separate disaster recovery solution with appropriately sized Region capacity.

Limitations

Be aware of these constraints when implementing this pattern:

  • ALB requirement: The pattern requires an Application Load Balancer in the Region. Workloads that rely on direct IP access through the Local Gateway (without an ALB) cannot use this pattern without an architecture change.
  • Stateful workloads: Applications with local disk state or in-memory sessions require external session stores (ElastiCache, DynamoDB) before they can burst. Without this, overflow instances serve requests without session context.
  • Database coupling: If your application writes to a database running exclusively on the Outpost, overflow instances in the Region cannot reach it without a cross-location replica or proxy. Read-heavy workloads with a Region read replica are ideal candidates.
  • Service link as single path: All ALB-to-Outpost traffic shares the service link with AWS control plane operations. Under extreme load, bandwidth contention can degrade both application traffic and management operations.
  • ALB on Outposts: As of this writing, ALB on Outposts does not support weighted target groups spanning both locations. The ALB must reside in the Region for this pattern to work.

Testing the pattern

Validate the burst mechanism before relying on it in production:

Simulate capacity pressure:

aws cloudwatch set-alarm-state \
  --alarm-name outposts-capacity-high \
  --state-value ALARM \
  --state-reason "Testing burst mechanism"

Verify overflow fleet launched:

aws autoscaling describe-auto-scaling-groups \
  --auto-scaling-group-names burst-overflow-fleet \
  --query "AutoScalingGroups[0].DesiredCapacity"

Verify ALB weights shifted (after recovery check runs):

aws elbv2 describe-listeners \
  --listener-arns <your-listener-arn> \
  --query "Listeners[0].DefaultActions[0].ForwardConfig.TargetGroups[*].[TargetGroupArn,Weight]"

Trigger recovery:

aws cloudwatch set-alarm-state \
  --alarm-name outposts-capacity-high \
  --state-value OK \
  --state-reason "Testing recovery"

Confirm overflow fleet scales back to zero and all traffic returns to Outposts targets. Recovery is gradual — the Amazon EventBridge rule evaluates every 5 minutes and steps weights back before scaling down, so full recovery may take 10–15 minutes depending on your weight progression configuration.

Clean up

To avoid ongoing charges, verify that the overflow Auto Scaling group has scaled to zero, then delete the stack:

sam delete --stack-name burst-to-region-stack

This removes all resources created by the template, including the Lambda function, CloudWatch alarm, SNS topic, Amazon EventBridge rule, and the overflow Auto Scaling group.

Conclusion

This Burst to Region pattern extends AWS Outposts capacity into the parent Region during peak demand. You trade a moderate latency increase for continued availability when local capacity is exhausted.

The pattern works best when you clearly classify which workloads can overflow, implement gradual traffic transitions, and maintain security and observability parity across both environments.

For the complete deployable AWS SAM template including the Lambda orchestrator, CloudWatch dashboard, and all IAM roles, see the GitHub repository. To learn more about capacity planning for Outposts, see Managing your AWS Outposts capacity using Amazon CloudWatch and AWS Lambda and AWS Outposts monitoring and reporting: A comprehensive Amazon EventBridge solution.

For more information, see the AWS Outposts User Guide and the Amazon EC2 Auto Scaling User Guide.

Build an AI email pipeline with Amazon Bedrock and SES Mail Manager

Post Syndicated from Zip Zieper original https://aws.amazon.com/blogs/messaging-and-targeting/build-an-ai-email-pipeline-with-amazon-bedrock-and-ses-mail-manager/

Processing inbound email attachments at scale involves extracting files, routing them by recipient, scanning for malware, and classifying content. This traditionally requires stitching together polling loops, event rules, and multiple integration points. Amazon Simple Email Service (Amazon SES) Mail Manager now provides two new rule actions that simplify this pattern. The Lambda action invokes AWS Lambda functions directly from rule sets, and the Bounce action returns rejection responses. Together, they let you build multi-step email processing pipelines with declarative configuration.

In this post, you learn how to build an attachment processing pipeline that automatically extracts email attachments and classifies them with Amazon Bedrock. The pipeline also rejects infected files with RFC-compliant bounce responses. The complete implementation is available as an AWS Cloud Development Kit (AWS CDK) deployment in the companion GitHub repository sample-amazon-ses-mail-manager-attachment-pipeline. You can deploy it manually using the steps in this post, or hand it off to an AI coding agent such as Kiro or Claude Code. The repository includes a machine-readable agentic deployment guide that walks an agent through every deployment step, from prerequisite checks to post-deploy verification.

Architecture of the inbound email pipeline: SES Mail Manager routes messages through a traffic policy and rule set to AWS Lambda, Amazon Simple Storage Service (Amazon S3), Amazon DynamoDB, and Amazon Bedrock

The problem: scaling document intake for a multi-tenant platform

Consider a fictitious SaaS platform from AnyCompany that lets customers submit documents by email. Each customer sends invoices, contracts, and supporting files to a dedicated address (for example, [email protected] or [email protected]). They expect those attachments to land in their isolated storage, classified and ready for downstream processing.

Without a purpose-built pipeline, the typical approach looks like this: an Amazon S3 event notification triggers a Lambda function that polls for new MIME objects, parses them, looks up the recipient in a routing table, and fans out extraction to another function. Worse, it relies on a separate virus-scanning step having run first. Orchestration lives in AWS Step Functions or Amazon EventBridge rules. Adding a new customer means updating routing configuration in multiple places. Adding classification means bolting on yet another Lambda in the chain.

The result is fragile. When volume spikes during month-end invoice runs or onboarding waves, the polling loop backs up and retries cascade. Infected files occasionally slip past the scanner because the scan and extraction steps are not transactionally linked.

This pipeline solves the problem declaratively. Mail Manager’s traffic policy rejects unauthorized senders and enforces size limits at the SMTP connection level. This filtering happens before any processing resources are consumed. The rule set handles virus scanning, bouncing, archiving, classification, and extraction in a single ordered sequence. Each step completes before the next begins. If an attachment is infected, the sender gets an immediate SMTP bounce. There are no silent failures and no orphaned files in downstream storage.

The result is a pipeline where:

  • Adding a customer means adding an email address to the Mail Manager address list and a row in Amazon DynamoDB. No changes to code.
  • Adding a classification category means editing a prompt string. No schema migration.
  • Infected files never reach storage because the bounce fires during the SMTP transaction, before any Lambda is invoked.

Pipeline architecture overview

Table 1: Architecture components and their roles in the email processing pipeline

Component Role
Amazon SES Mail Manager open Ingress Endpoint Email arrives via public internet at a Mail Manager open ingress point over SMTP.
Mail Manager traffic policy Filters spam using the Abusix (or Spamhaus) email add-on, then enforces a recipient allowlist at the connection level.
Mail Manager rule set Messages for allowed recipients are passed to the rule set, which sequentially evaluates each message against two rules.
Rule 1 Uses the Trend Micro email add-on to scan for infected attachments, then bounces any unsafe messages back to sender (using Amazon SES outbound).
Rule 2 Clean messages passed from Rule 1 are copied to a Mail Manager archive and written as raw Multipurpose Internet Mail Extensions (MIME) objects to a “landing-zone” Amazon S3 bucket.
AWS Lambda (AttachmentProcessor) Triggered by the arrival of objects in the S3 bucket, this function parses MIME email, extracts attachments, and routes them to per-recipient S3 buckets.
AWS Lambda (EmailCategorizer) Triggered by the arrival of objects in the landing-zone S3 bucket, this function classifies each email using Amazon Nova Micro via Amazon Bedrock and writes results to Amazon DynamoDB.
Amazon S3 (landing zone + per-recipient buckets) Stores raw MIME objects in a shared landing-zone bucket; stores extracted attachments in isolated per-recipient buckets keyed by local part (for example, invoices/ for [email protected]).
Amazon DynamoDB (RecipientBucketLookup) Maps recipient email addresses to their designated S3 bucket and key prefix.
Amazon DynamoDB (EmailCategories) Stores Amazon Bedrock classification results: category, urgency, and summary.
Amazon Bedrock (Amazon Nova Micro) Classifies each email into a category (invoice, contract, HR, unknown) and urgency level.
AWS IAM roles Mail Manager and Lambda execution permissions following the principle of least privilege.

How the Mail Manager traffic policy filters connections

The traffic policy (Receive-attachments) makes connection-level decisions before any message content is processed. It evaluates two statements in order:

  1. Deny spam — Connections from senders flagged by Abusix as spam sources are denied immediately.
  2. Allow approved recipients — Connections where the recipient is in the approved-recipients address list pass through to the rule set.

The policy uses a default action of DENY, so any connection that does not match an explicit ALLOW statement is rejected. The policy also enforces a 35 MB maximum message size. You can add additional statements to enforce SPF, DKIM, or DMARC authentication results. This is useful in regulated industries where sender verification is required before any processing occurs.

The PolicyStatements array defines the evaluation order (deny first, then allow):

PolicyStatements=[
    {   # Statement 1: Deny connections from known spam sources
        "Action": "DENY",
        "Conditions": [{"BooleanExpression": {
            "Evaluate": {"Analysis": {"Analyzer": "ABUSIX_ADDON_ARN", "ResultField": "isListed"}},
            "Operator": "IS_TRUE",
        }}],
    },
    {   # Statement 2: Allow only recipients in the approved list
        "Action": "ALLOW",
        "Conditions": [{"BooleanExpression": {
            "Evaluate": {"IsInAddressList": {"Attribute": "RECIPIENT", "AddressLists": ["ADDRESS_LIST_ARN"]}},
            "Operator": "IS_TRUE",
        }}],
    },
]

For the complete create_traffic_policy call with all parameters, see the companion repository.

API reference: CreateTrafficPolicy

Rule set: the processing pipeline

Messages that pass the traffic policy enter the rule set (attachment-pipeline-rules), which evaluates two rules in order.

Rule 1 — Virus scan and bounce

This rule checks the Trend Micro add-on result. If Trend Micro reports isPassed = FALSE (infected attachment detected) — note that Mail Manager has already accepted the message by this point — the rule fires a Bounce action, which generates a non-delivery report (NDR) back to the sender with SMTP 550 (permanent failure) and status 5.7.1 (security/policy reason). It then Drops the message. No further rules run.

This after-the-fact NDR prevents infected messages from entering your processing pipeline while still providing clear guidance to legitimate senders.

Rule 2 — Process clean email

This rule has no conditions, so it applies to every message that passed the virus scan. It runs four actions in sequence:

  1. Archive — Mail Manager stores a copy in the archive for compliance and electronic discovery (eDiscovery).
  2. WriteToS3 — Mail Manager writes the raw MIME object to the amzn-s3-demo-bucket-general-receiving S3 bucket, keyed by message ID.
  3. InvokeLambda (EmailCategorizer, REQUEST_RESPONSE) — Mail Manager invokes the categorizer, which classifies the email with Amazon Bedrock and writes results to Amazon DynamoDB.
  4. InvokeLambda (AttachmentProcessor, REQUEST_RESPONSE) — Mail Manager invokes the processor, which extracts attachments and routes them to per-recipient S3 locations.

The categorizer fires before the attachment processor by design: the attachment processor deletes the original MIME from Amazon S3 after successfully extracting attachments. By running first, the categorizer is guaranteed to find the MIME in Amazon S3.

Because the Bounce and Drop actions fire in Rule 1, the Lambda functions in Rule 2 are never invoked for infected messages. There is no risk of malicious content reaching your Amazon S3 buckets or Amazon Bedrock.

API reference: CreateRuleSet

How Amazon Bedrock classifies inbound email

The MailManager-EmailCategorizer function uses Amazon Nova Micro (amazon.nova-micro-v1:0) to classify each email. Amazon Nova Micro is a fast, lightweight text-only model optimized for classification and structured output tasks. Access to all Amazon Bedrock foundation models, including Amazon Nova Micro, is available by default in all commercial AWS Regions. No access request is needed.

The function performs the following steps:

  1. Parses the recipient, message ID, and subject from the Mail Manager event.
  2. Retrieves the raw MIME from the amzn-s3-demo-bucket-general-receiving S3 bucket.
  3. Extracts the plain-text or HTML body from the MIME structure.
  4. Sends the subject (capped at 500 characters) and body (capped at 4,000 characters) to Amazon Bedrock with a classification prompt.
  5. Writes the structured result to the EmailCategories DynamoDB table.

The classification prompt returns a structured JSON response:

{
    "category": "invoice | contract | hr | unknown",
    "urgency": "urgent | non-urgent",
    "summary": "<50-word summary>"
}

If Amazon Bedrock returns an error or malformed JSON, the function falls back to category: unknown, urgency: non-urgent and continues. It never blocks the attachment processor.

Choosing a classification model

To customize the classification categories for your use case, update the SYSTEM_PROMPT in the categorizer Lambda function. The prompt uses a structured instruction format that you can extend with additional categories, urgency levels, or routing rules. For example, an insurance carrier could add categories like claim_new, claim_status, document_submission, and complaint to automatically triage patient email. You can also update the COMPANY_NAME environment variable to inject your organization’s name into the classification prompt without modifying the function code.

To switch the model, update the BEDROCK_MODEL_ID environment variable. The following table compares supported options:

Model Model ID Best for Latency Relative cost
Amazon Nova Micro amazon.nova-micro-v1:0 Fast structured classification, low latency ~200ms Lowest
Amazon Nova Lite amazon.nova-lite-v1:0 Richer summaries, multi-label classification ~400ms Moderate
Anthropic Claude 3 Haiku anthropic.claude-3-haiku-20240307-v1:0 Complex reasoning, nuanced categorization ~600ms Higher

Attachment extraction and routing

The MailManager-AttachmentProcessor function handles MIME parsing, recipient-based routing, and cleanup. It performs the following steps:

  1. Parses the recipient email address and message ID from the Mail Manager event information.
  2. Retrieves the raw MIME message from the amzn-s3-demo-bucket-general-receiving S3 bucket using the message ID from the event as the S3 key.
  3. Looks up the recipient’s S3 destination in the RecipientBucketLookup DynamoDB table, or creates a new entry if this is the first email for that recipient.
  4. Extracts attachment parts from the MIME message, skipping plain-text and HTML body parts that have no file name.
  5. Copies each attachment to the recipient’s S3 bucket at the prefix {local_part}/ (for example, invoices/ for [email protected]).
  6. Deletes the original MIME object from the landing-zone bucket, but only if every attachment copy succeeded. If any copy failed, the MIME is retained for retry.
  7. Returns a response to Mail Manager indicating success or failure.

This synchronous invocation pattern allows the rule set to make routing decisions based on the Lambda function’s response. If attachment extraction fails, subsequent rules can bounce the message or route it to a quarantine location.

Attachment detection logic

The function detects attachments using three criteria:

  1. Content-Disposition containing attachment.
  2. Any MIME part with a file name (even if disposition is inline or missing).
  3. Non-text, non-multipart parts (such as application/pdf or image/*).

For parts without a file name, the function generates one from the content type (for example, attachment.pdf).

Input validation and security

The pipeline implements the following input validation to protect against malicious content and unexpected inputs:

  • messageId validation — the messageId from the Mail Manager event is validated against an alphanumeric-plus-hyphen pattern ([a-zA-Z0-9\-]+) before use as an S3 key. Unexpected formats raise a ValueError, which causes Mail Manager to apply the ActionFailurePolicy.
  • Attachment filename sanitization — filenames from MIME Content-Disposition headers are attacker-controlled. Before use as S3 key components, each filename is processed through os.path.basename() to strip directory components, leading-dot stripping to prevent hidden-file creation, and a character allowlist ([\w.\- ]). Filenames are also truncated to 255 characters.
  • Prompt size caps — the email body sent to Amazon Bedrock is capped at 4,000 characters. The subject line is capped at 500 characters, preventing oversized prompts and excessive token usage.

The following additional controls are recommended before adapting this pipeline for production:

  • Validate attachment file types against an approved allowlist (such as .pdf, .docx, .xlsx). Reject or quarantine messages with disallowed file types.
  • Implement per-attachment size limits in addition to the overall 35 MB message size limit.
  • Verify MIME structure integrity before parsing. Handle malformed MIME structures as error conditions.
  • Log validation failures to Amazon CloudWatch for security monitoring and audit purposes.

AWS CloudFormation and CDK support for Mail Manager rule actions

The InvokeLambda and Bounce rule actions are supported natively in AWS::SES::MailManagerRuleSet as of March 2026. The companion CDK stack uses CfnMailManagerRuleSet directly. No Custom Resource is required.

When using the Python CDK L1 bindings, note that typed property classes for Bounce and InvokeLambda are not yet exposed in the Python bindings. Pass these actions as plain dicts with camelCase keys matching the AWS CloudFormation property names. RuleActionProperty accepts Dict[str, Any] for each field:

ses.CfnMailManagerRuleSet.RuleActionProperty(
    bounce={
        "smtpReplyCode": "550",
        "statusCode": "5.7.1",
        "diagnosticMessage": "Your attachment was infected.",
        "sender": "[email protected]",
        "roleArn": role.role_arn,
        "actionFailurePolicy": "CONTINUE",
    }
)

API reference: AWS::SES::MailManagerRuleSet | AWS CDK API Reference

Prerequisites

This post and companion GitHub project assume familiarity with SMTP protocols, email infrastructure concepts, AWS Lambda, Amazon S3, Amazon DynamoDB, and AWS IAM.

Estimated time: 20–30 minutes to deploy and test.

Estimated cost: This pipeline uses a Mail Manager open ingress endpoint that costs $50/mo in addition to various AWS services that are charged based on actual usage. In a low-volume test environment (fewer than 1,000 email messages per day), costs should typically be under $60 USD per month driven primarily by Mail Manager archiving, S3 storage, Lambda invocations, and Amazon Bedrock token usage. Use the AWS Pricing Calculator to estimate costs for your expected volume.

AWS IAM permissions: The deploying user needs permissions to create and manage AWS CloudFormation stacks, Lambda functions, S3 buckets, DynamoDB tables, AWS IAM roles, and Amazon SES Mail Manager resources. For testing, AdministratorAccess is sufficient. For production, scope permissions to the specific actions required: cloudformation:CreateStack, lambda:CreateFunction, s3:CreateBucket, dynamodb:CreateTable, iam:CreateRole, iam:PassRole, ses:CreateTrafficPolicy, ses:CreateRuleSet, and ses:CreateAddressList. (Separately, the Lambda functions’ own execution roles, created by the stack, grant bedrock:InvokeModel at runtime; that permission is not needed by the person deploying the stack.)

To deploy this pipeline, you need the following:

  1. An active AWS account.
  2. AWS Command Line Interface (AWS CLI) version 2.x or later installed and configured with credentials and default region.
  3. AWS CDK version 2.x or later installed (npm install -g aws-cdk) and Python 3.12 or later.
  4. Amazon SES configured with production access in the target region with a verified Amazon SES identity for the bounce sender address.
  5. Ability to administer the DNS entries for the Amazon SES identity to add an MX record pointing to the Mail Manager ingress endpoint’s A record.

Deployment

Tip: Whichever path you choose, review the Prerequisites section first to make sure your AWS account has the necessary permissions and that you have a verified domain available in Amazon SES. The complete solution is available as an open-source reference implementation. To deploy it in your AWS account, clone the companion repository:

git clone https://github.com/aws-samples/sample-amazon-ses-mail-manager-attachment-pipeline.git
cd sample-amazon-ses-mail-manager-attachment-pipeline

From here, you have two paths to get up and running:

Option 1: Deploy manually

Follow the step-by-step instructions in the repository’s README.md. At a high level, you will:

  1. Install prerequisites (AWS CDK, Node.js, Python).
  2. Configure your environment variables (AWS account, region, verified domain).
  3. Bootstrap your CDK environment.
  4. Deploy the stack with cdk deploy.
  5. Complete post-deployment verification (confirm email receiving rules are active and test with a sample message).

Option 2: Deploy with a coding agent

If you use an AI-powered coding assistant (such as Amazon Q Developer CLI or Kiro), install the AWS MCP server and SES/Mail Manager skills to empower your AI assistants with deep context on Amazon SES and Mail Manager. These resources give your assistant live access to AWS APIs and CDK documentation, which significantly reduces trial-and-error during deployment. The repository’s AGENTS.md file contains machine-readable guidance, deployment failure recovery patterns, and region handling notes specifically for AI assistants. Simply point your AI assistant at the AGENTS.md file in the repository root. This file provides structured, machine-readable instructions that guide the agent through the full deployment, from prerequisite checks through stack deployment and validation, without manual intervention.

# Example: point your agent at the instructions
@agent follow AGENTS.md

Validating the deployment

Once your stack is deployed and the MX record is in place, send a test email with an attachment to one of your approved recipient addresses. Then confirm each stage of the pipeline executed successfully:

1. Check Lambda execution

Open Amazon CloudWatch Logs for both functions and confirm they completed without errors:

aws logs tail /aws/lambda/MailManager-EmailCategorizer --follow
aws logs tail /aws/lambda/MailManager-AttachmentProcessor --follow

You should see log entries showing the message ID being processed by each function in sequence: the categorizer first, then the attachment processor.

2. Confirm email classification

Query the EmailCategories DynamoDB table to verify Amazon Bedrock classified your test message:

aws dynamodb scan --table-name EmailCategories --max-items 1

A successful record includes category, urgency, and a short summary, all generated by Amazon Nova Micro from the email’s subject and body.

3. Verify attachment extraction

Look up your recipient’s S3 destination in the RecipientBucketLookup table, then list the bucket contents to confirm the attachment arrived:

aws dynamodb get-item --table-name RecipientBucketLookup \
  --key '{"recipient": {"S": "[email protected]"}}'

aws s3 ls s3://<bucket-name>/<prefix>/ --recursive

If all three checks pass, your pipeline is fully operational. Email messages are being scanned, classified, and routed to per-recipient storage without any external orchestration.

Troubleshooting

If your test email does not flow through the pipeline as expected, start with these common issues:

Symptom Likely cause Resolution
Bounce action fails silently — infected emails are dropped without notification The bounce_sender identity is not verified in the deployment region. Amazon SES identities are regional. Verify the domain in your target region: aws sesv2 create-email-identity --email-identity example.com --region <region>, add the DKIM CNAMEs to DNS, and wait for verification. No redeployment required.
Bounce action returns a validation error bounce_sender is set to a bare domain instead of an email address Use a full address like [email protected], not just example.com

For CDK deployment issues, stack rollback errors, and teardown conflicts, see the repository troubleshooting guide.

General debugging tip: Both Lambda functions log to /aws/lambda/MailManager-EmailCategorizer and /aws/lambda/MailManager-AttachmentProcessor in Amazon CloudWatch Logs. Start there for any runtime failures.

Clean up

To avoid ongoing charges, destroy the stack when you are done:

AWS_DEFAULT_REGION= cdk destroy

Note: If the destroy fails with a ConflictException, detach the ingress point from the traffic policy first. Amazon DynamoDB tables created with RETAIN policies may also need manual deletion. See the repository’s Common failure modes table for details.

Do not forget to remove the MX record from your domain’s DNS once the ingress point is deleted. After completing the clean up, verify on the AWS Management Console that the Mail Manager ingress endpoint, Amazon S3 buckets, Amazon DynamoDB tables, and Lambda functions no longer appear in your account.

Conclusion

The Lambda action and Bounce action in Amazon SES Mail Manager support multi-step inbound email processing without complex orchestration workarounds. This pipeline demonstrates how these capabilities work together in production: scanning attachments for malware, classifying email content with AI, extracting and routing files to per-recipient storage, and providing immediate RFC-compliant feedback to senders. The modular architecture supports extension: add new classification categories, integrate additional scanning engines, or chain Lambda functions for multi-stage processing. The synchronous invocation pattern means that every processing step completes before the next begins, giving you full control over the pipeline flow. Get started by cloning the sample-amazon-ses-mail-manager-attachment-pipeline repository and deploying to your account. For an overview of the four new Mail Manager capabilities used in this pipeline, see Four new Amazon SES Mail Manager capabilities, explained.

FAQ

Q: Can I use a different Amazon Bedrock model for email classification?

Yes. Update the BEDROCK_MODEL_ID environment variable on the MailManager-EmailCategorizer Lambda function. No changes to code are required. See the preceding model comparison table for supported options.

Q: Do I need to request access to Amazon Nova Micro?

No. In all commercial AWS Regions, access to Amazon Bedrock foundation models including Amazon Nova Micro is available by default. AWS GovCloud (US) regions require an explicit access request through the Amazon Bedrock console.

Q: What happens if the Lambda function times out or fails?

A REQUEST_RESPONSE invocation is time-bounded to approximately 30 seconds, or sooner if your function’s own configured timeout is shorter. In either case, Mail Manager applies the ActionFailurePolicy configured on the rule action. If set to CONTINUE, the pipeline moves to the next action. If set to DROP, the message is discarded. This pipeline uses CONTINUE, so a transient classification failure does not block attachment delivery.

Q: Can I add more classification categories?

Yes. Edit the SYSTEM_PROMPT in the categorizer Lambda function. The function writes whatever categories the model returns to Amazon DynamoDB. No schema changes are needed.

Q: How does the pipeline handle email messages with no attachments?

The AttachmentProcessor detects zero attachment parts, skips extraction, deletes the raw MIME from the landing-zone bucket, and returns success. The EmailCategorizer still classifies the message normally.

Q: What is the maximum attachment size supported?

The traffic policy enforces a 35 MB maximum message size (total MIME payload including all attachments and base64 encoding overhead). Individual attachments are not size-limited beyond this total cap.

Q: Can I deploy this with an AI coding agent?

Yes. The repository includes an AGENTS.md file with machine-readable deployment instructions. Point your AI assistant (Kiro, Claude Code, Amazon Q Developer CLI) at this file and it handles the full deployment without manual intervention.

Q: Is the Bounce action RFC-compliant?

Yes, with one clarification: it is not a live SMTP-transaction rejection. Mail Manager first accepts the message, then the rule set runs. If the Bounce action fires, it generates a non-delivery report (NDR) back to the sender with an RFC 5321-compliant SMTP reply code and an RFC 3463-compliant enhanced status code.


About the authors

Centralized CloudTrail monitoring across 100+ AWS accounts

Post Syndicated from Jagdish Komakula original https://aws.amazon.com/blogs/big-data/centralized-cloudtrail-monitoring-across-100-aws-accounts/

Organizations running workloads across dozens or hundreds of AWS accounts face a common challenge: centralized security monitoring at scale. Security teams need to search through hundreds of gigabytes of AWS CloudTrail logs daily to detect threats and satisfy compliance requirements for SOC 2, PCI DSS, and HIPAA audits. They also need to provide role-based access to multiple teams with different responsibilities.

Without purpose-built infrastructure, this often involves manual log searching that takes hours and compliance report generation that takes days. It also leads to fragmented code bases of custom AWS Lambda functions managing index lifecycles across environments. A single shared search domain without consistent access control compounds the problem further.

In this post, we show you how to build a centralized CloudTrail monitoring solution on Amazon OpenSearch Service. Terraform manages the full stack, from domain provisioning to access control and lifecycle policies. The solution handles 200 GB/day of CloudTrail logs, provides automated threat detection alerts, and gives 4 different teams isolated, role-appropriate access to the data.

Solution overview

The following diagram shows the architecture. CloudTrail logs flow from 100+ AWS accounts through an organization trail into a centralized S3 bucket. Amazon Simple Queue Service (Amazon SQS) notifications trigger the OpenSearch Ingestion pipeline. The pipeline auto-scales between 2 and 10 OpenSearch Compute Units (OCUs) to parse and index the logs into the OpenSearch domain. Four team-specific roles access the data through OpenSearch Dashboards with tenant isolation.

CloudTrail logs flow from 100+ accounts into S3, then through Amazon SQS and OpenSearch Ingestion into the OpenSearch domain used by 4 team roles

The key components are:

  • CloudTrail aggregation. An organization trail sends logs from 100+ accounts into a centralized Amazon Simple Storage Service (Amazon S3) bucket.
  • Ingestion. An Amazon OpenSearch Ingestion pipeline picks up new logs through Amazon SQS notifications on the S3 bucket. It automatically scales between 2 and 10 OCUs based on queue depth. Throttling on lower-environment queues prevents development and test spikes from starving production ingestion.
  • Amazon OpenSearch Service domain. 6 or1.4xlarge data nodes (OpenSearch Optimized instances) with 3 dedicated r8g.large master nodes, fine-grained access control, encryption at rest, and node-to-node encryption.
  • Infrastructure as code. Index templates, Index State Management (ISM) policies, roles, role mappings, tenants, alerting monitors, and dashboards are all declared in Terraform and applied consistently across environments.

Prerequisites

To implement this solution, you need the following:

  • An organization in AWS Organizations with CloudTrail enabled across member accounts.
  • Terraform v1.5+ with the AWS provider and the OpenSearch provider.
  • A virtual private cloud (VPC) with private subnets for the OpenSearch domain.
  • IAM roles for each team that will access the OpenSearch domain.
  • An Amazon Simple Notification Service (Amazon SNS) topic for security alert notifications.
  • Familiarity with Amazon OpenSearch Service, Terraform, and AWS CloudTrail.
  • Sample Terraform code is available in the GitHub repository

Implementation

This section walks through the Terraform code for each component of the solution, starting with the workload profile that informed our sizing decisions.

Workload profile

Before sizing the cluster, we defined the workload characteristics and SLAs for the centralized CloudTrail monitoring platform:

Metric Value
Index throughput 200 GB/day (~18,000 docs/sec)
Search queries ~2,000 queries/day (~0.023 QPS)
Average search latency < 100 ms (achieved: 76 ms)
Saved searches 600+
Dashboards and visualizations 100+
User teams 4 (Security Ops, Incident Response, Compliance, DevOps)
Retention 30 days (hot tier)
Availability target 99.9%

This is a write-heavy ingestion workload. The primary use case is automated alerting and periodic compliance queries rather than continuous interactive search. This workload profile informed the decision to use OR1 (storage-optimized) instances with zero replicas, prioritizing indexing throughput over search parallelism.

Domain provisioning

Start by provisioning the Amazon OpenSearch Service domain with encryption, fine-grained access control, and VPC placement:

resource "aws_opensearch_domain" "cloudtrail" {
  domain_name    = var.domain_name
  engine_version = "OpenSearch_3.3"

  cluster_config {
    instance_type          = "or1.4xlarge.search"
    instance_count         = 6
    zone_awareness_enabled = true
    zone_awareness_config {
      availability_zone_count = 3
    }
  }

  dedicated_master_config {
    dedicated_master_enabled = true
    dedicated_master_type    = "r6g.large.search"
    dedicated_master_count   = 3
  }

  ebs_options {
    ebs_enabled = true
    volume_type = "gp3"
    volume_size = 500
    iops        = 3000
    throughput  = 125
  }

  encrypt_at_rest { enabled = true }
  node_to_node_encryption { enabled = true }

  domain_endpoint_options {
    enforce_https       = true
    tls_security_policy = "Policy-Min-TLS-1-2-PFS-2023-10"
  }

  advanced_security_options {
    enabled                        = true
    internal_user_database_enabled = false
    master_user_options {
      master_user_arn = var.master_user_arn
    }
  }

  vpc_options {
    subnet_ids         = var.vpc_subnet_ids
    security_group_ids = var.vpc_security_group_ids
  }

  tags = {
    Environment = "production"
    Project     = "centralized-cloudtrail-monitoring"
    ManagedBy   = "terraform"
  }
}

This solution was built on OR1 instances, which are storage-optimized and use Amazon Elastic Block Store (Amazon EBS) (gp3 or io1) for local storage, with data copied synchronously to Amazon S3 as it arrives. This storage structure provides increased indexing throughput because indexing is performed exclusively on primary shards. Replicas are backed by Amazon S3 through segment replication, eliminating the CPU overhead of document replication on replica nodes. For new deployments, we recommend OR2 instances, which offer up to 26% higher indexing throughput compared to OR1 while maintaining the same storage-optimized architecture.

Ingestion pipeline

The Amazon OpenSearch Ingestion pipeline provides serverless, auto scaling ingestion from Amazon S3 into the OpenSearch domain. It picks up new CloudTrail logs through Amazon SQS notifications on the centralized S3 bucket and scales between 2 and 10 OpenSearch Compute Units (OCUs) based on queue depth:

resource "aws_iam_role" "osis_pipeline" {
  name = "cloudtrail-osis-pipeline-role"
  assume_role_policy = jsonencode({
    Version = "2012-10-17"
    Statement = [{
      Action    = "sts:AssumeRole"
      Effect    = "Allow"
      Principal = { Service = "osis-pipelines.amazonaws.com" }
    }]
  })
}

resource "aws_iam_policy" "osis_pipeline" {
  name = "cloudtrail-osis-pipeline-policy"
  policy = jsonencode({
    Version = "2012-10-17"
    Statement = [
      {
        Action   = ["s3:GetObject", "s3:ListBucket"]
        Effect   = "Allow"
        Resource = [var.cloudtrail_bucket_arn, "${var.cloudtrail_bucket_arn}/*"]
      },
      {
        Action   = ["sqs:ReceiveMessage", "sqs:DeleteMessage", "sqs:GetQueueAttributes"]
        Effect   = "Allow"
        Resource = var.cloudtrail_sqs_queue_arn
      },
      {
        Action   = ["es:DescribeDomain", "es:ESHttp*"]
        Effect   = "Allow"
        Resource = "${aws_opensearch_domain.cloudtrail.arn}/*"
      }
    ]
  })
}

resource "aws_iam_role_policy_attachment" "osis_pipeline" {
  role       = aws_iam_role.osis_pipeline.name
  policy_arn = aws_iam_policy.osis_pipeline.arn
}

resource "aws_cloudwatch_log_group" "osis_pipeline" {
  name              = "/aws/vendedlogs/OpenSearchIngestion/cloudtrail-pipeline"
  retention_in_days = 30
}

resource "aws_osis_pipeline" "cloudtrail" {
  pipeline_name = "cloudtrail-ingestion"
  pipeline_configuration_body = <<-EOT
    version: "2"
    cloudtrail-pipeline:
      source:
        s3:
          notification_type: "sqs"
          codec:
            json:
          compression: "gzip"
          sqs:
            queue_url: "${var.cloudtrail_sqs_queue_url}"
          aws:
            sts_role_arn: "${aws_iam_role.osis_pipeline.arn}"
            region: "${data.aws_region.current.name}"
      processor:
        - date:
            from_time_received: true
            destination: "@timestamp"
      sink:
        - opensearch:
            hosts: ["https://${aws_opensearch_domain.cloudtrail.endpoint}"]
            index: "cloudtrail-%{yyyy.MM.dd}"
            aws:
              sts_role_arn: "${aws_iam_role.osis_pipeline.arn}"
              region: "${data.aws_region.current.name}"
  EOT
  min_units = 2
  max_units = 10
  log_publishing_options {
    is_logging_enabled = true
    cloudwatch_log_destination {
      log_group = aws_cloudwatch_log_group.osis_pipeline.name
    }
  }
  tags = {
    Environment = "production"
    Project     = "centralized-cloudtrail-monitoring"
    ManagedBy   = "terraform"
  }
}

The pipeline uses the S3 source plugin with SQS-based notifications. When new CloudTrail log files land in S3, an SQS message triggers the pipeline to fetch and parse them. The min_units and max_units parameters control auto scaling. The pipeline starts at 2 OCUs and scales up to 10 based on queue depth, handling ingestion spikes without manual intervention. For lower environments (development and testing), you can apply throttling on the Amazon SQS queue to prevent non-production spikes from starving production ingestion capacity.

Index template

Define index templates up front to avoid painful reindexing later. The following template sets explicit mappings for CloudTrail fields, optimizes for write throughput with async translog durability, and integrates with ISM for automatic rollover:

resource "opensearch_index_template" "cloudtrail" {
  name = "cloudtrail-template"
  body = jsonencode({
    index_patterns = ["cloudtrail-*"]
    priority       = 100
    template = {
      settings = {
        number_of_shards                                  = 6
        number_of_replicas                                = 0
        "index.refresh_interval"                          = "10s"
        "index.translog.durability"                       = "async"
        "index.translog.sync_interval"                    = "30s"
        "plugins.index_state_management.rollover_alias"   = "cloudtrail"
      }
      mappings = {
        properties = {
          "@timestamp"        = { type = "date" }
          eventSource         = { type = "keyword" }
          eventName           = { type = "keyword" }
          awsRegion           = { type = "keyword" }
          sourceIPAddress     = { type = "ip" }
          errorCode           = { type = "keyword" }
          errorMessage        = { type = "text" }
          recipientAccountId  = { type = "keyword" }
          userIdentity = {
            properties = {
              type      = { type = "keyword" }
              arn       = { type = "keyword" }
              accountId = { type = "keyword" }
              userName  = { type = "keyword" }
              sessionContext = {
                properties = {
                  sessionIssuer = {
                    properties = {
                      type     = { type = "keyword" }
                      arn      = { type = "keyword" }
                      userName = { type = "keyword" }
                    }
                  }
                }
              }
            }
          }
          requestParameters = { type = "object", enabled = true }
          responseElements  = { type = "object", enabled = true }
        }
      }
    }
  })
}

Defining mappings before ingestion prevents mapping conflicts and avoids the need to reindex data after the fact.

Lifecycle management (ISM policy)

The following ISM policy replaces custom Lambda functions with a single declarative policy. The rollover action uses two OR conditions: min_index_age and min_primary_shard_size. Whichever threshold is reached first triggers the rollover. This keeps shard sizes bounded while ensuring timely rotation even during low-volume periods:

resource "opensearch_ism_policy" "cloudtrail_lifecycle" {
  policy_id = "cloudtrail-lifecycle"
  body = jsonencode({
    policy = {
      description   = "CloudTrail lifecycle - rollover, retain 30d, delete"
      default_state = "hot"
      ism_template  = [{ index_patterns = ["cloudtrail-*"], priority = 100 }]
      states = [
        {
          name    = "hot"
          actions = [{ rollover = { min_primary_shard_size = "30gb", min_index_age = "1d" } }]
          transitions = [{ state_name = "delete", conditions = { min_index_age = "30d" } }]
        },
        {
          name        = "delete"
          actions     = [{ delete = {} }]
          transitions = []
        }
      ]
    }
  })
}

This approach reduces lifecycle management code by approximately 60% compared to per-environment Lambda functions, and changes deploy in a single terraform apply.

Note: For indexes ingesting more than 100 GB/day (such as CloudTrail at 200 GB/day in this deployment), you can override min_index_age to 12h to roll over more frequently. The two conditions are OR-based in OpenSearch ISM. If a shard reaches 30 GB before 1 day, it rolls over on size. If 1 day passes before 30 GB, it rolls over on age.

Multi-team access control

When multiple teams need different access levels to the same data, define all roles declaratively and use for_each to create them consistently. The following example defines 4 team roles with varying permissions:

locals {
  team_roles = {
    security_ops = {
      description         = "Security Operations - full read, alert management"
      cluster_permissions = ["cluster_monitor", "cluster:admin/opendistro/alerting/*"]
      index_permissions = [
        { index_patterns = ["cloudtrail-*"], allowed_actions = ["read", "search", "get"] },
        { index_patterns = [".opendistro-alerting-*"], allowed_actions = ["read", "write", "search", "get", "delete"] }
      ]
    }
    incident_response = {
      description         = "Incident Response - full read for investigation"
      cluster_permissions = ["cluster_monitor"]
      index_permissions = [
        { index_patterns = ["cloudtrail-*"], allowed_actions = ["read", "search", "get"] }
      ]
    }
    compliance_auditors = {
      description         = "Compliance - read-only"
      cluster_permissions = []
      index_permissions = [
        { index_patterns = ["cloudtrail-*"], allowed_actions = ["read", "search"] }
      ]
    }
    devops = {
      description         = "DevOps - infra metrics and limited CloudTrail"
      cluster_permissions = ["cluster_monitor"]
      index_permissions = [
        { index_patterns = ["infra-metrics-*"], allowed_actions = ["read", "search", "get"] },
        { index_patterns = ["cloudtrail-*"], allowed_actions = ["read", "search"] }
      ]
    }
  }
}

resource "opensearch_role" "teams" {
  for_each            = local.team_roles
  role_name           = each.key
  description         = each.value.description
  cluster_permissions = each.value.cluster_permissions

  dynamic "index_permissions" {
    for_each = each.value.index_permissions
    content {
      index_patterns  = index_permissions.value.index_patterns
      allowed_actions = index_permissions.value.allowed_actions
    }
  }

  dynamic "tenant_permissions" {
    for_each = [each.key]
    content {
      tenant_patterns = [each.key]
      allowed_actions = ["kibana_all_write"]
    }
  }
}

resource "opensearch_roles_mapping" "teams" {
  for_each      = local.team_roles
  role_name     = opensearch_role.teams[each.key].role_name
  backend_roles = var.team_iam_roles[each.key]
}

resource "opensearch_tenant" "teams" {
  for_each    = local.team_roles
  tenant_name = each.key
  description = "Dashboard workspace for ${replace(each.key, "_", " ")}"
}

This approach maps IAM roles (not individual users) to OpenSearch roles. Adding a new team means adding one entry to the locals block and running terraform apply. Each team gets an isolated tenant in OpenSearch Dashboards, preventing cross-team interference with saved searches, visualizations, and dashboard configurations. We chose OpenSearch Dashboards because tenants, roles, visualizations, and saved objects can all be managed programmatically through the Terraform OpenSearch provider, keeping the entire stack under infrastructure-as-code governance. For teams building new visualizations outside of Terraform-managed workflows, we recommend OpenSearch UI. This next-generation analytics interface supports multiple data sources, provides workspaces for team isolation, and remains available during cluster upgrades.

Alerting

Define alerting monitors in Terraform to detect security-critical events automatically. The following monitor catches CloudTrail tampering attempts (StopLogging, DeleteTrail) and sends alerts through Amazon SNS:

resource "opensearch_monitor" "cloudtrail_tampering" {
  body = jsonencode({
    name     = "CloudTrail Tampering Detection"
    type     = "monitor"
    enabled  = true
    schedule = { period = { interval = 1, unit = "MINUTES" } }
    inputs = [{
      search = {
        indices = ["cloudtrail-*"]
        query = {
          size = 5
          query = {
            bool = {
              must = [{ terms = { eventName = ["StopLogging", "DeleteTrail",
                "UpdateTrail", "PutEventSelectors", "DeleteEventDataStore"] } }]
              filter = [{ range = { "@timestamp" = { gte = "now-1m" } } }]
            }
          }
        }
      }
    }]
    triggers = [{
      name     = "trail_tampering_detected"
      severity = "1"
      condition = { script = {
        source = "ctx.results[0].hits.total.value > 0"
        lang   = "painless"
      } }
      actions = [{
        name             = "notify_security"
        destination_id   = var.sns_destination_id
        message_template = { source = "CRITICAL: CloudTrail tampering detected." }
      }]
    }]
  })
}

Results and performance

After deploying the solution, we measured steady-state performance against the SLAs defined in the workload profile:

Metric Target Achieved
Index throughput 200 GB/day 200 GB/day sustained (~18,000 docs/sec)
Search latency (avg) < 100 ms 76 ms
Search availability 99.9% 99.95%+ (no unplanned downtime in 145 days)
Alert detection time < 2 minutes ~1 minute (monitor interval)
Compliance report generation < 5 minutes On-demand via saved searches

Key outcomes:

  • Threat detection dropped from hours to minutes. Automated alerting replaced manual log searching. The CloudTrail tampering monitor detects suspicious activity within one minute of the event.
  • Compliance reports generate on demand. With 600+ saved searches and 100+ dashboards, compliance teams produce SOC 2, PCI DSS, and HIPAA audit evidence in minutes rather than days.
  • Four teams operate independently. Each team has its own isolated tenant in OpenSearch Dashboards, preventing cross-team interference with saved searches and dashboard configurations.
  • Zero custom Lambda functions. ISM policies, index templates, and access control are all managed declaratively through Terraform, eliminating the previous fragmented code base.

Best practices

  • Define index templates before ingesting anything. Changing mappings on existing indices means reindexing. Get this right first.
  • Set rollover thresholds based on your actual ingestion rate. At 200 GB/day, rolling over at 30 GB keeps shard counts manageable while balancing query performance.
  • Test ISM transitions in a lower environment first. Warm and cold migrations on large indices take time.
  • Map IAM roles, not users. People change teams. Roles stay stable. This simplifies access management.
  • Put an Amazon SQS queue between S3 and the ingestion pipeline. This gives you per-environment throttling control without modifying pipeline configuration.
  • Use for_each aggressively. Roles, tenants, index patterns, and monitors all follow a pattern across teams or environments, so use for_each to eliminate copy-paste drift.
  • Consider OpenSearch UI for new visualization workflows. This solution uses OpenSearch Dashboards for Terraform-managed tenants and roles. OpenSearch UI is a next-generation interface that supports multiple data sources, stays available during cluster upgrades, and includes workspaces for team isolation. It is the recommended interface for creating new dashboards and visualizations going forward.

Optional: Extending with cold storage for longer retention

For organizations with compliance requirements mandating longer retention (for example, 7 years for PCI DSS or HIPAA), you can extend the ISM policy with warm and cold tiers. The following example adds tiered storage that moves data through hot, warm, cold, and delete states:

states = [
  {
    name    = "hot"
    actions = [{ rollover = { min_primary_shard_size = "30gb", min_index_age = "1d" } }]
    transitions = [{ state_name = "warm", conditions = { min_index_age = "30d" } }]
  },
  {
    name = "warm"
    actions = [
      { warm_migration = {} },
      { force_merge = { max_num_segments = 1 } }
    ]
    transitions = [{ state_name = "cold", conditions = { min_index_age = "365d" } }]
  },
  {
    name    = "cold"
    actions = [{ cold_migration = { timestamp_field = "@timestamp" } }]
    transitions = [{ state_name = "delete", conditions = { min_index_age = "2555d" } }]
  },
  {
    name        = "delete"
    actions     = [{ cold_delete = {} }]
    transitions = []
  }
]

Warm storage uses force-merge to reduce segment count (lowering query overhead), while cold storage moves data entirely to Amazon S3 for minimal cost. This tiered approach keeps hot-tier performance high while meeting long-term audit requirements.

Cleanup

To avoid incurring ongoing charges, remove the resources created in this post by running:

terraform destroy

This removes the OpenSearch domain, ingestion pipeline, IAM roles, SQS queues, and all associated configurations. Verify that you have exported any data or dashboards you want to retain before running destroy.

Conclusion

In this post, we showed you how to build a centralized CloudTrail monitoring solution on Amazon OpenSearch Service with Terraform managing the entire stack. The approach moves from fragile, manually configured systems with redundant Lambda code to a version-controlled, peer-reviewed, consistently deployed infrastructure.

Threat detection drops from hours to minutes with automated alerting. Compliance reports that took days now generate on demand. And your team spends time on security analysis instead of infrastructure maintenance.

To get started, use the AWS Terraform provider aws_opensearch_domain resource for the domain, then use the Terraform OpenSearch provider for index templates, ISM policies, roles, and monitors. Configure your ingestion pipeline to transform and enrich incoming CloudTrail logs before indexing, building a modern, scalable security foundation that grows with your organization.

The complete source code for this solution is available in the GitHub repository: GitHub repository

For more on the services used in this solution:


About the authors

Jagdish Komakula

Jagdish Komakula

Jagdish is a Senior Delivery Consultant at AWS Professional Services, focused on Amazon OpenSearch Service and Infrastructure automation. He has spent the last several years guiding financial services customers through building data platforms that scale.

Aditya Ambati

Aditya Ambati

Aditya is a Delivery Consultant at AWS Professional Services, focused on DevOps and infrastructure as code. He works with customers on automating cloud operations and implementing GitOps practices.

Securing your Amazon S3 buckets: Identifying and remediating over-permissioned access

Post Syndicated from Hetal Kolekar original https://aws.amazon.com/blogs/security/securing-your-amazon-s3-buckets-identifying-and-remediating-over-permissioned-access/

Misconfigured Amazon Simple Storage Service (Amazon S3) buckets can expose your data to unauthorized access. Without proactive review, S3 bucket policies or Access Control Lists (ACLs) configured with broad access may go unnoticed in your environment. In this post, you learn how to identify and fix over-permissioned S3 buckets across your AWS environment, along with best practice recommendations and automation opportunities to help you prevent security gaps. This post provides a workflow framework and methodology recommendations for your security team to adapt. The focus of this post is on the what and why rather than a prescriptive implementation. You will need to customize the approach based on your organization’s requirements and existing security tooling.

This solution is intended for security engineers, cloud architects, and DevOps teams managing single- or multiple-account AWS environments with Amazon S3 workloads that require access management.

Prerequisites

Before you begin, make sure you have the following in place:

Solution overview

This solution uses a five-phase workflow diagram to detect, remediate, and continuously monitor over-permissioned S3 buckets across your AWS accounts. The following workflow diagram illustrates the high-level end-to-end process for identifying and remediating over-permissioned S3 buckets across your Amazon Web Services (AWS) environment.

Figure 1: Amazon S3 over-permissive access – Detection, remediation, monitoring and cleanup workflow

Figure 1: Amazon S3 over-permissive access – Detection, remediation, monitoring and cleanup workflow

The diagram in Figure 1 consists of five phases:

  1. Setup and prerequisites – Configure AWS Organizations or multi-account access, designate a central security account, deploy AWS Config across all accounts, and enable AWS Security Hub with a central administrator.
  2. Detection and identification – Deploy AWS Config rules (such as s3-bucket-public-read-prohibited and s3-bucket-public-write-prohibited) and run an audit Lambda function that scans each S3 bucket. The function checks three areas: Public Access Block configuration, bucket policy status, and bucket ACL grants. Buckets with issues are added to a risky buckets list. The function then generates a report in CSV and JSON format, uploads it to an output S3 bucket, and sends an SNS alert.
  3. Remediation – Address findings using one or more approaches – Apply restrictive bucket policies to deny public read/write access and restrict access to specific IAM principals; deploy a remediation Lambda function to automatically update bucket policies and disable public access settings; or use CloudFormation StackSets to deploy standardized policies across multiple accounts.
  4. Continuous monitoring – Schedule the audit Lambda function for recurring scans (daily or weekly) using Amazon EventBridge. Use EventBridge to detect policy changes, configure automated notifications for new violations, enable IAM Access Analyzer for S3 to identify external access, and run regular compliance scans.
  5. Resource cleanup – Review and delete resources created during the audit that are no longer needed, including Lambda functions and IAM roles, EventBridge rules, SNS topics and subscriptions, audit output S3 buckets, AWS Config rules, and Security Hub (if enabled only for this audit).

Cost considerations

This section covers the AWS services used in this solution and their associated costs so you can estimate spend before deployment. The primary cost drivers are AWS Config and Security Hub, which scale with the number of accounts and resources you monitor. Lambda, Amazon EventBridge, Amazon SNS, and Amazon S3 typically add minimal costs for most environments. Start with a pilot in one or two accounts to validate costs before scaling.

  • AWS Config – Charges per configuration item recorded and per rule evaluation. Costs scale with the number of accounts and resources tracked.
  • Security Hub – Charges per account per AWS Region for security checks and finding ingestion.
  • Lambda – Charges per request and per GB-second of compute time.
  • EventBridge – Scheduled rules are free. Custom event bus usage might incur charges.
  • Amazon SNS – Charges per notification delivered.
  • Amazon S3 – Storage costs for audit report output files. Minimal for most environments.
  • AWS IAM Access Analyzer – Check the AWS IAM Access Analyzer pricing page to understand which features have costs associated with them.

Check the service pricing pages for current rates. Use the AWS Pricing Calculator to estimate costs for your specific environment before enabling services across all accounts. Consider starting with a pilot in one or two accounts to validate costs before scaling.

Detect and report over-permissioned buckets

This section walks you through setting up the audit environment, deploying the Lambda-based scanner, and generating reports of over-permissioned S3 buckets across your accounts. Follow these steps to identify over-permissioned S3 buckets in your multi-account environment, starting with preparing your environment for an Amazon S3 audit.

To set up the multi-account audit environment:

  1. Set up AWS Organizations or multi-account access. Set up centralized management of your AWS accounts using AWS Organizations or configure cross-account IAM roles.
  2. Choose a central security account. Choose one account as your security/audit account. This account will run the audit Lambda function and collect results from member accounts.
  3. Create an Amazon SNS topic for alerts. Subscribe your security team to receive notifications when over-permissioned buckets are detected. Note the topic Amazon Resource Name (ARN) from the output—you will need it when creating the Lambda execution role (step 6) and the Lambda function (step 9). Confirm the email subscription before testing; Amazon SNS doesn’t deliver alerts until the subscription is confirmed. Learn more in the Amazon SNS Developer Guide.
  4. (Optional): Create an S3 bucket for audit reports. If you plan to use Script v2 for historical reporting and trend analysis, create a dedicated bucket now. Skip this step if you only need real-time alerts using Script v1.
  5. Plan cross-account IAM roles. The central security account needs permission to scan member accounts. Design cross-account roles that:
    1. Grant minimum Amazon S3 read permissions (list buckets, read policies, ACLs, public access configurations).
    2. Include an external ID condition to mitigate the confused deputy problem.
    3. Can be deployed consistently using AWS CloudFormation StackSets.
    4. See the IAM documentation on creating cross-account roles, The confused deputy problem, and IAM security best practices for additional guidance on role configuration and trust policies.

      Note: The specific trust policy and permissions policy for your cross-account roles will depend on organizational requirements. Work with your IAM administrators to grant minimum necessary access for the audit function.

  6. Create the Lambda execution role. Create an IAM role for your Lambda function with the permissions it needs to scan buckets, publish alerts, and write logs. Apply the principle of least privilege—grant only the minimum Amazon S3 read permissions required for the audit (such as, listing buckets, reading bucket policies, ACLs, and public access block configurations), Amazon SNS publish permission for the alert topic created in step 3, Amazon S3 write permission for the output bucket created in step 4 (Script v2), and Amazon CloudWatch Logs permissions. For multi-account scanning, also include sts:AssumeRolepermission for the cross-account role ARNs created in step 5. The AWS Lambda execution role documentation has instructions on creating and configuring execution roles.
  7. To deploy the S3 audit solution Deploy the audit components
    1. Enable AWS Config in member accounts. AWS Config provides compliance monitoring and can detect when S3 buckets are created or modified with public access settings. This will enable the Lambda-based audit to receive real-time detection between scheduled scans. The AWS Config Developer Guide has setup instructions. Deploy pre-defined AWS Config rules to identify overly permissive settings. These managed rules provide automated compliance checking. When AWS Config detects violations, it sends findings to Security Hub (configured in step 8) for centralized visibility alongside the Lambda audit results.
      • s3-bucket-public-read-prohibited
      • s3-bucket-public-write-prohibited
      • Create AWS Config rules for specific permission patterns. For the full list of available rules, see the AWS Config managed rules reference
  8. Enable Security Hub for centralized visibility. Enable AWS Security Hub in member accounts and configure the central security account as the administrator. Security Hub aggregates findings from AWS Config rules (step 7), IAM Access Analyzer (enabled later), and can receive custom findings from your Lambda audit function, providing a single dashboard for Amazon S3 security issues across your organization. See the Security Hub User Guide for setup details.
  9. Deploy the audit Lambda function. Deploy a Python Lambda function using the Boto3 library to list S3 buckets, check their policies, ACLs, and IAM permissions, and identify over-permissioned buckets. See the example scripts that follow.

Important: These code examples aren’t production ready. Adapt them to meet your organization’s requirements and test them in a non-production environment before deployment.

Choose your approach:

  • Script v1 – Best for immediate SNS alerts when issues are detected.
  • Script v2 – Best for historical reports, trend analysis using BI tools.
  • Both scripts – Best for different schedules and ongoing needs.

Audit Lambda function – Example script v1 (Scan and alert)

The following is an example of a Lambda function script for reference purposes. Review, adapt, and test before use in your environment, it scans all S3 buckets in the current account and checks for:

  • Public Access block configuration gaps
  • Bucket policies that allow public access
  • ACL grants to AllUsers

Note: Replace placeholder values with actual values before deployment:

  • <REGION>– Your AWS Region (for example, us-east-1)
  • <ACCOUNT_ID>– Your 12-digit AWS account ID
  • <TOPIC_NAME>– The name of your SNS topic created in step 3
import boto3
import json

def lambda_handler(event, context):
    s3 = boto3.client('s3')
    sns = boto3.client('sns')
    risky_buckets = []
    errors = []

    try:
        buckets = s3.list_buckets()['Buckets']
    except Exception as e:
        return {'statusCode': 500, 'body': f'Failed to list buckets: {str(e)}'}

    for bucket in buckets:
        bucket_name = bucket['Name']
        issues = []

        try:
            # Check Public Access Block — all four settings should be enabled
            try:
                pab = s3.get_public_access_block(Bucket=bucket_name)
                config = pab['PublicAccessBlockConfiguration']
                if not all([
                    config.get('BlockPublicAcls'),      # Block new public ACLs
                    config.get('BlockPublicPolicy'),     # Block new public bucket policies
                    config.get('IgnorePublicAcls'),      # Ignore existing public ACLs
                    config.get('RestrictPublicBuckets')   # Restrict access to public buckets
                ]):
                    issues.append('Public Access Block not fully enabled')
            except s3.exceptions.NoSuchPublicAccessBlockConfiguration:
                issues.append('No Public Access Block configured')

            # Check bucket policy — flag if policy status is public
            try:
                policy_status = s3.get_bucket_policy_status(Bucket=bucket_name)
                if policy_status['PolicyStatus']['IsPublic']:
                    issues.append('Bucket policy allows public access')
            except s3.exceptions.NoSuchBucketPolicy:
                pass  # No bucket policy is acceptable

            # Check bucket ACL
            acl = s3.get_bucket_acl(Bucket=bucket_name)
            for grant in acl.get('Grants', []):
                grantee = grant.get('Grantee', {})
                uri = grantee.get('URI', '')
                # 'AllUsers' = anonymous public access
                # 'AuthenticatedUsers' = any AWS account (still overly permissive)
                if grantee.get('Type') == 'Group' and ('AllUsers' in uri or 'AuthenticatedUsers' in uri):
                    issues.append('Bucket ACL grants public access')
                    break

            if issues:
                risky_buckets.append({'bucket': bucket_name, 'issues': issues})

        except Exception as e:
            errors.append(f'{bucket_name}: {str(e)}')

    # Send alert if risky buckets found
    if risky_buckets:
        message = f'Found {len(risky_buckets)} buckets with public access:\n\n'
        for item in risky_buckets:
            message += f"  {item['bucket']}: {', '.join(item['issues'])}\n"

        sns.publish(
            TopicArn='arn:aws:sns:<REGION>:<ACCOUNT_ID>:<TOPIC_NAME>',
            Subject='S3 Public Access Alert',
            Message=message
        )

    return {
        'statusCode': 200,
        'body': json.dumps({
            'risky_buckets': risky_buckets,
            'errors': errors,
            'total_checked': len(buckets)
        })
    }

Multi-account scanning: This script scans the current account only. To scan across member accounts, see the Multi-account extension section later in this post.

Audit Lambda function – Example script v2 (CSV and JSON report)

The following is an example Lambda function script for reference purposes. Before deploying any script, review error handling, logging, output structure, and permissions. This script generates CSV and JSON output files and uploads them to an S3 bucket for reporting and business intelligence (BI) dashboard integration.

You can deploy both functions with different EventBridge schedules, for example, Script v1 daily for alerts and Script v2 weekly for reports.

Note: Before you deploy this script, replace <OUTPUT_BUCKET_NAME> with the S3 bucket you created for audit reports in step 4.

import boto3
import csv
import json
import os

def lambda_handler(event, context):
    s3 = boto3.client('s3')
    buckets = s3.list_buckets()['Buckets']

    full_access_buckets = []
    for bucket in buckets:
        bucket_name = bucket['Name']
        try:
            bucket_policy = s3.get_bucket_policy(Bucket=bucket_name)['Policy']
            policy = json.loads(bucket_policy)
            for statement in policy['Statement']:
                if (statement['Effect'] == 'Allow'
                    and statement['Principal'] == '*'
                    and 'Action' in statement
                    and 's3:*' in statement['Action']):
                    full_access_buckets.append({'BucketName': bucket_name})
                    break
        except s3.exceptions.ClientError as e:
            if e.response['Error']['Code'] != 'NoSuchBucketPolicy':
                print(f'Error checking bucket policy for {bucket_name}: {e}')

    # Output CSV
    csv_output = os.path.join('/tmp', 'full_access_buckets.csv')
    with open(csv_output, 'w', newline='') as csvfile:
        writer = csv.DictWriter(csvfile, fieldnames=['BucketName'])
        writer.writeheader()
        writer.writerows(full_access_buckets)

    # Output JSON
    json_output = os.path.join('/tmp', 'full_access_buckets.json')
    with open(json_output, 'w') as jsonfile:
        json.dump(full_access_buckets, jsonfile, indent=2)

    # Upload to Amazon S3
    output_bucket = '<OUTPUT_BUCKET_NAME>'
    s3.upload_file(csv_output, output_bucket, 'full_access_buckets.csv')
    s3.upload_file(json_output, output_bucket, 'full_access_buckets.json')

    return {
        'statusCode': 200,
        'body': json.dumps(f'CSV and JSON files uploaded to {output_bucket}')
    }

Important: If this function runs on a schedule, consider implementing a file naming strategy with timestamps to prevent overwriting previous reports or establish a lifecycle policy to manage retention. Include the output bucket in your cleanup procedures when the auditing process is no longer needed.

What if no over-permissioned buckets are found?

If the audit scan returns zero risky buckets, document the clean baseline for future comparison and move to the verification and monitoring phase to so new buckets or policy changes don’t introduce risk over time.

Multi-account extension

The preceding example scripts scan buckets in the current account only. To scan across member accounts in your organization, add the following AssumeRole logic. This function assumes the cross-account IAM role you created during setup, then returns an Amazon S3 client with temporary credentials for each member account.

Note: Before you deploy, configure the following Lambda environment variables:

  • <MEMBER_ACCOUNTS> – Comma-separated list of 12-digit account IDs to scan (for example, 111111111111,222222222222)
  • <CROSS_ACCOUNT_ROLE_NAME> – The IAM role name created in each member account (for example, S3AuditRole)
  • <EXTERNAL_ID> – The external ID configured in the trust policy (for example, s3-audit-external-id)
import boto3
import os

def get_member_s3_clients():
    """
    Assumes the cross-account audit role in each member account
    and returns a list of (account_id, s3_client) tuples.
    """
    sts = boto3.client('sts')
    member_accounts = os.environ.get('<MEMBER_ACCOUNTS>', '').split(',')
    cross_account_role_name = os.environ.get('<CROSS_ACCOUNT_ROLE_NAME>')
    external_id = os.environ.get('<EXTERNAL_ID>')

    clients = []
    for account_id in member_accounts:
        account_id = account_id.strip()
        if not account_id:
            continue

        try:
            assumed_role = sts.assume_role(
                RoleArn=f'arn:aws:iam::{account_id}:role/{cross_account_role_name}',
                RoleSessionName='S3AuditSession',
                ExternalId=external_id
            )

            # Create S3 client with assumed credentials
            s3_client = boto3.client(
                's3',
                aws_access_key_id=assumed_role['Credentials']['AccessKeyId'],
                aws_secret_access_key=assumed_role['Credentials']['SecretAccessKey'],
                aws_session_token=assumed_role['Credentials']['SessionToken']
            )
            clients.append((account_id, s3_client))

        except Exception as e:
            print(f'Failed to assume role in account {account_id}: {e}')

    return clients

To scan each member account, replace the single-account s3.list_buckets() call with a loop over member accounts:

def lambda_handler(event, context):
    all_risky_buckets = []
    all_errors = []

    # Scan each member account
    for account_id, s3_client in get_member_s3_clients():
        try:
            buckets = s3_client.list_buckets()['Buckets']
            for bucket in buckets:
                # ... same scanning logic as the single-account scripts ...
                # Use s3_client instead of s3 for each API call
                pass
        except Exception as e:
            all_errors.append(f'Account {account_id}: {e}')

    # ... same alerting/reporting logic ...

The Lambda execution role in the central security account needs sts:AssumeRole permission for the cross-account role ARNs. Add this to the execution role policy you created in step 5.

Remediate elevated access

This section describes how to fix over-permissioned buckets using account-level controls, bucket policies, and optional automation. Any elevated access that you find needs to be remediated.

Enable Amazon S3 Block Public Access (account level)

Before applying individual bucket policies, enable Amazon S3 Block Public Access at the account level. This prevents buckets in the account from being made public, regardless of individual bucket policies or ACLs. See theS3 Block Public Access documentation for configuration details. See the following example AWS CLI command; replace <ACCOUNT_ID> with the ID of the account you’re using to manage resource access:

aws s3control put-public-access-block \
  --account-id <ACCOUNT_ID> \
  --public-access-block-configuration \
BlockPublicAcls=true,IgnorePublicAcls=true,BlockPublicPolicy=true,RestrictPublicBuckets=true

For multi-account environments, deploy this setting across member accounts using AWS CloudFormation StackSets or AWS Organizations service control policies (SCPs).

Important: Before enabling account-level S3 Block Public Access, check whether any workloads need public bucket access (for example, static website hosting, public dataset sharing). Coordinate with your application teams to identify any exceptions.

Remediate using bucket policies

Implement bucket policies that restrict access to specific IAM users, roles, or accounts. When crafting policies, apply the principle of least privilege and include only the actions and principals required for your use case.

Example S3 bucket policy: deny public read/write access. Modify the resource ARN, actions, and conditions to match your requirements:

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Effect": "Deny",
      "Principal": "*",
      "Action": [
        "s3:PutObject", "s3:PutObjectAcl",
        "s3:GetObject", "s3:GetObjectAcl",
        "s3:DeleteObject"
      ],
      "Resource": "arn:aws:s3:::<BUCKET_NAME>/*",
      "Condition": {
        "StringEquals": {
          "s3:x-amz-acl": ["public-read", "public-read-write"]
        }
      }
    }
  ]
}

Example S3 bucket policy: restrict access to specific IAM principals. Replace <ACCOUNT_ID>, <USERNAME>, and <ROLE_NAME>:

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Sid": "AllowObjectAccess",
      "Effect": "Allow",
      "Principal": {
        "AWS": [
          "arn:aws:iam::<ACCOUNT_ID>:user/<USERNAME>",
          "arn:aws:iam::<ACCOUNT_ID>:role/<ROLE_NAME>"
        ]
      },
      "Action": ["s3:GetObject", "s3:PutObject", "s3:DeleteObject"],
      "Resource": "arn:aws:s3:::<BUCKET_NAME>/*"
    },
    {
      "Sid": "AllowBucketAccess",
      "Effect": "Allow",
      "Principal": {
        "AWS": [
          "arn:aws:iam::<ACCOUNT_ID>:user/<USERNAME>",
          "arn:aws:iam::<ACCOUNT_ID>:role/<ROLE_NAME>"
        ]
      },
      "Action": ["s3:ListBucket", "s3:GetBucketLocation"],
      "Resource": "arn:aws:s3:::<BUCKET_NAME>"
    }
  ]
}

See the Amazon S3 bucket policy documentation for additional examples and guidance.

Automate remediation with Lambda or CloudFormation StackSets (optional):

You can also remediate using Lambda or CloudFormation Stacksets:

  • Create Lambda functions to automatically update bucket policies or disable public access settings for flagged buckets
  • Use CloudFormation StackSets to deploy standardized bucket policies and S3 Block Public Access settings across multiple accounts

Verify your remediation

This section explains how to confirm that your fixes are effective before moving to ongoing monitoring. After applying remediation, verify the fix is effective before setting up ongoing monitoring:

  1. Re-run the audit Lambda function – Confirm the previously flagged buckets no longer appear in the risky buckets list.
  2. Check Security Hub compliance – Verify the compliance status has changed from FAILED to PASSED for Amazon S3-related controls.
  3. Validate with IAM Access Analyzer – Review findings for the remediated S3 buckets. Active findings should resolve automatically after public access is removed.
  4. Test application functionality – Confirm that legitimate workloads continue to function correctly.

Document the verification results for your auditing needs. If any S3 buckets still show issues, investigate whether the policy was applied correctly or if there are conflicting permissions.

Automation opportunities

This section covers optional strategies to automate ongoing detection and maintain your security posture without manual intervention.

  1. (Optional) Schedule recurring scans with Amazon EventBridge
    • Regular security scans help identify new issues arising from configuration changes or newly created S3 buckets. When new security risks are detected, Amazon SNS sends an alert and automatically initiates the remediation phase (Workflow 2 in Figure 1). To avoid repeated alerts, you can configure the audit Lambda function to run on a schedule and compare current results with the previous baseline to generate notifications when new findings are discovered.
    • For ongoing monitoring, you can schedule the audit Lambda function to run on a recurring basis using EventBridge. Create a scheduled rule with a cron expression (for example, daily at 6:00 AM UTC or weekly on Mondays), add the Lambda function as the target, and grant EventBridge permission to invoke it. See Amazon EventBridge scheduling documentation for instructions on creating scheduled rules and configuring targets.
  2. Enable IAM Access Analyzer for Amazon S3
    • IAM Access Analyzer monitors bucket policies, ACLs, and access points to identify buckets accessible from outside your account or organization. Create an analyzer scoped to your organization or individual account, then review findings to identify unintended external access. Findings automatically flow into Security Hub when both services are enabled, giving you a dashboard view for Amazon S3 security findings. See the IAM Access Analyzer documentation for setup and usage instructions.
  3. Automate notifications for policy drift
    • Recurring scans might surface new findings from policy drift or newly created buckets. When new risks are detected, Amazon SNS alert triggers and the remediation cycle repeat (as shown in Workflow 2 in Figure 1) sends email notifications. Configure the audit Lambda function to compare current scan results against the previous baseline and alert on new findings for ongoing reviews.

Clean up

This section lists the resources created during this walkthrough that you should review and remove when they are no longer needed. If the following services were not previously active in your account, leaving them enabled might result in additional ongoing charges. See the Cost considerations section for details. Review and remove unused resources to optimize costs.

Delete or disable the following script-generated resources if they’re not required after outputs are generated. Focus first on Lambda functions and EventBridge rules if you’re not running recurring scans. If you enabled AWS Config or Security Hub specifically for this audit, evaluate whether you need them for other compliance requirements before disabling.

  • Lambda – Functions, IAM roles, and policies created for auditing
  • Amazon EventBridge – Scheduled rules created for recurring audit triggers
  • Amazon SNS – Topics and subscriptions created for notifications
  • Amazon S3 – Buckets containing script-generated audit output files
  • AWS Config – Rules and recorders if no longer needed for compliance
  • Security Hub – Disable if enabled solely for this audit
  • IAM Access Analyzer – Delete the analyzer if no longer needed for ongoing monitoring

Note: Be careful when deleting data and consider temporarily disabling services first to check for dependencies. Only delete resources generated as part of your audit outputs. Verify you have retained any necessary results before proceeding. Verify resources are not used by other workloads before deletion.

Best practices

This section provides recommendations to maintain secure Amazon S3 configurations long-term. To learn more about maintaining secure Amazon S3 configurations, review the AWS documentation links provided in the conclusion. The following recommendations aren’t exhaustive. Adapt and extend them based on your organization’s evolving security requirements and AWS best practices guidance. After you’ve fixed existing issues, these practices help you maintain secure Amazon S3 configurations.

  • Start with account-level controls – Enable S3 Block Public Access at the account level. This prevents buckets from becoming public even if someone misconfigures an individual bucket policy. For multi-account environments, enforce this through AWS Organizations SCPs.
  • Automate detection – Use IAM Access Analyzer to detect external access. Schedule your audit Lambda function with EventBridge to catch new issues weekly or daily, depending on your change frequency. Compare scan results against previous baselines to identify drift.
  • Standardize across accounts – Use CloudFormation StackSets to deploy the same secure configuration to all accounts in your organization, reducing the chance of configuration drift. Use StackSets for IAM roles, AWS Config rules, and S3 Block Public Access settings.

Additional security measures

  • Regularly review and rotate cross-account IAM role credentials and external IDs
  • Implement Amazon S3 server-side encryption (SSE-S3 or SSE-KMS) for data at rest
  • Enable S3 access logging and AWS CloudTrail data events for audit trails

Conclusion

This section summarizes what you accomplished and suggests next steps to maintain your S3 security posture. By implementing the detection, remediation, and monitoring workflow outlined in this post, you can proactively identify and secure over-permissioned S3 buckets across your AWS environment. To maintain your ongoing security posture, enable IAM Access Analyzer for continuous monitoring and schedule recurring audits with EventBridge. To learn more about Amazon S3 security best practices, see Security best practices for Amazon S3

For more information:

If you have feedback about this post, submit comments in the Comments section below.


Hetal Kolekar

Hetal Kolekar

Hetal is a Sr. Technical Account Manager at AWS with more than 21 years of experience in Infrastructure Architecture, Security, Systems Engineering, and Consulting. He excels in leading teams to strengthen their cloud security posture and helps customers scale up their security using AWS services. Hetal is a guitarist and loves playing at church.

Manomayi Vedam

Manonmayi Vedam

Manonmayi is a Senior TAM and Product Owner at AWS, specializing in AI-driven cloud enablement, security, and generative AI risk across Healthcare, Financial Services, Energy, and Public Sector. She co-leads global security programs for Fortune 500 clients, contributes to the NIST Cyber AI Profile RMF and NCCoE, and is a Fellow at SCRS with recognition from GlobeeAwards and IEEE.

Fernando Freitas

Fernando Freitas

Fernando is a Sr. Technical Account Manager at AWS in Salt Lake City, focused on helping customers achieve their desired outcomes with the AWS Cloud. Fernando is passionate about Identity and Security, Training and Education.

Automate certificates with ACME support in AWS Certificate Manager

Post Syndicated from Anthony Harvey original https://aws.amazon.com/blogs/security/automate-certificates-with-acme-support-in-aws-certificate-manager/

Customers tell us that managing TLS certificates at scale is one of their biggest operational concerns. The Certification Authority Browser Forum (CA/Browser Forum) has mandated a phased reduction in maximum certificate validity for public certificates. By March 2027, the maximum validity drops to 100 days. By March 2029, it lasts for 47 days. For an organization managing 1,000 certificates, the final transition means roughly 30 renewal events every day. Renewal and rotations of renewed certificates at that cadence isn’t something manual processes or ticket-driven workflows can sustain at scale.

We recently announced Automated Certificate Management Environment (ACME) protocol support in AWS Certificate Manager (ACM). With this launch, you can use the ACME clients your teams already know, including popular open source tools like certbot, cert-manager, acme.sh, and win-acme, to automate public certificate issuance and renewal for your infrastructure. Customers that are using third-party certificate authorities (CAs) can point their existing ACME-compatible clients at ACM instead of their current CA, with minimal reconfiguration. This applies whether it’s running on Amazon Web Services (AWS), on premises, or in a hybrid environment. Certificates created through ACME are registered in ACM, giving you a unified view of your entire certificate inventory.

This post covers how the feature works, how to get started, and the controls and best practices to help you manage certificate issuance at scale.

Background

ACME is an open source protocol that automates the process of verifying domain ownership and issuing certificates and has become a standard mechanism for certificate automation. While ACM has long provided managed certificate issuance and renewal for AWS-integrated services such as Elastic Load Balancing (ELB), Amazon CloudFront, and Amazon API Gateway, many customers also need to automate certificates for their own infrastructure, including servers they manage in their data centers, Kubernetes clusters, Internet of Things (IoT) fleets, and hybrid environments. Until now, those customers had to turn to external providers. This launch brings the ACM automation model to that same infrastructure, using the standard ACME protocol with AWS managed certificate endpoints.

How it works

The feature introduces a new centrally provisioned and managed resource type: the ACME endpoint. Each endpoint is an AWS resource with a unique ACME directory URL and AWS Identity and Access Management (IAM)-based access controls. You create and manage endpoints through the ACM API or AWS Management Console, and point your existing ACME clients at the endpoint URL. Certificates issued through your endpoint are automatically registered with ACM, appearing in your certificate inventory alongside certificates created by the RequestCertificate and ImportCertificate API calls.

The architecture separates into two planes. In the control plane, PKI administrators use ACM APIs to create ACME endpoints, pre-approve the domains an endpoint is allowed to issue for, and generate external account binding (EAB) credentials. In the data plane, ACME clients register with an endpoint using EAB credentials and request certificates for domains the administrator has already validated. This architecture is how we provide customers the ability to scale. Instead of each client proving domain ownership on every request, a principal with appropriate ACM permissions (typically your PKI administrator) validates domains once at the endpoint level, and then application owners don’t need DNS credentials to get a certificate.

Adding to the data plane, EABs control client access to the endpoints. Each EAB is bound to an IAM role that controls what certificate operations the ACME client can perform, and credentials you generate in ACM are distributed to authorized ACME clients. An ACME client authorized for one endpoint can’t use a different endpoint. This creates security boundaries between environments. For example, a client authorized for your development endpoint can’t obtain certificates from your production endpoint.

Figure 1 shows the ACME request flow through ACM. An ACME client authenticates to an ACME endpoint using EAB credentials. The endpoint routes certificate orders to Amazon Trust Services for issuance. Issued certificates are registered in ACM inventory, where Amazon EventBridge and AWS CloudTrail provide expiration alerting and audit logging.

Figure 1: An ACME architecture and workflow

Figure 1: An ACME architecture and workflow

Getting started

Getting started with the new ACME feature in ACM is straightforward. Use the following steps to create your first ACME-generated certificate.

Prerequisites

  • An AWS account with permissions to create and manage ACM resources
  • An ACME client installed on your infrastructure (for example, Certbot, cert-manager, acme.sh, or others)
  • AWS Command Line Interface (AWS CLI) installed on your device (see this blog post for the console equivalent)
  • Amazon Route 53 hosted zone for your domain, or the ability to create a CNAME record with your DNS provider

Step 1: Create an ACME endpoint

Before you can use ACME clients with ACM, you need to create an ACME endpoint. This endpoint provides the URL that your ACME clients will use to request certificates.

  1. Run the following command from the AWS CLI to create an ACME endpoint:
    aws acm create-acme-endpoint \
      --authorization-behavior PRE_APPROVED \
      --certificate-authority '{"PublicCertificateAuthority":{"AllowedKeyAlgorithms":["EC_prime256v1"]}}

  2. Note the endpoint Amazon Resource Name (ARN) from the response.
    {"AcmeEndpointArn": "arn:aws:acm:us-east-1:123456789012:acme-endpoint/11111111-2222-3333-4444-555555555555"}

  3. Run the following command to retrieve the endpoint URL, replacing the ARN with your endpoint ARN:
    aws acm describe-acme-endpoint \
    --acme-endpoint-arn arn:aws:acm:us-east-1:123456789012:acme-endpoint/11111111-2222-3333-4444-555555555555

  4. Save the output of the ACME EndpointUrl:
    {
        "AcmeEndpoint": {
            "AcmeEndpointArn": "arn:aws:acm:us-east-1:123456789012:acme-endpoint/11111111-2222-3333-4444-555555555555",
            "EndpointUrl": "https://acm-acme-enroll.<region>.api.aws/6666666-7777-8888-9999-000000000000/directory",
            "Status": "ACTIVE",
            "AuthorizationBehavior": "PRE_APPROVED",
            "Contact": "REQUIRED",
            "CertificateAuthority": {
                "PublicCertificateAuthority": {
                    "AllowedKeyAlgorithms": [
                        "EC_prime256v1"
                    ]
                }
            },
            "CreatedAt": "2026-07-14T18:23:58.876000-04:00",
            "UpdatedAt": "2026-07-14T18:23:58.876000-04:00"
        }
    }
    

Step 2: Pre-approve a domain

Before ACME clients can request a certificate, the administrator validates the domain using DNS once at the endpoint level. Use DomainScope to control exactly which certificate patterns are allowed:

  • Enabling only ExactDomain restricts clients to that specific name,
  • Subdomains enabled allows names like api.example.com,
  • Wildcards enabled allows *.example.com.

Leave a scope disabled to block that pattern outright, even if an otherwise-valid ACME request asks for it. For a production endpoint, consider enabling only ExactDomain and Subdomains and leaving Wildcards disabled for a stricter posture.

aws acm create-acme-domain-validation \
--acme-endpoint-arn arn:aws:acm:us-east-1:123456789012:acme-endpoint/11111111-2222-3333-4444-555555555555 \
--domain-name example.com \
--prevalidation-options '{"DnsPrevalidation":{"DomainScope":{"ExactDomain":"ENABLED","Subdomains":"ENABLED","Wildcards":"DISABLED"},"HostedZoneId":"Z1234567890ABC"}}'

If your domain is hosted in Route 53, specifying HostedZoneId lets ACM create the required CNAME record automatically. If your domain is hosted elsewhere, omit it and create the provided CNAME record manually with your DNS provider. Validation typically completes within a few seconds after the record is in place.

You will receive the following response back:

{
    "AcmeDomainValidationArn": "arn:aws:acm:us-east-1:123456789012:acme-endpoint/1111111-2222-3333-4444-555555555555/acme-domain-validation/6666666-8888-9999-0000-11111111111"
}

Step 3: Generate EAB credentials

EAB credentials authenticate your ACME clients to your endpoint. Generate a unique set of credentials for each client or environment to maintain security boundaries.

  1. Run the following command to generate your EAB credentials, adjusting your expiration to fit your organization’s risk profile:
    aws acm create-acme-external-account-binding \
        --acme-endpoint-arn arn:aws:acm:region:111122223333:acme-endpoint/00000000-0000-0000-0000-000000000000 \
        --role-arn arn:aws:iam::111122223333:role/AcmeIssuanceRole \
        --expiration '{"Value": 7, "Type": "DAYS"}'

  2. Note the response from a successful invocation of the command
    {
        "ExternalAccountBinding": {
            "AcmeExternalAccountBindingArn": "arn:aws:acm:region:111122223333:acme-endpoint/00000000-0000-0000-0000-000000000000/acme-external-account-binding/1234567-1234-1234-1234-123456789012",
            "AcmeEndpointArn": "arn:aws:acm:region:111122223333:acme-endpoint/00000000-0000-0000-0000-000000000000",
            "RoleArn": "arn:aws:iam::123456789012:role/service-role/AcmAcmeIssuanceRole-XXXXXXXX",
            "ExpiresAt": "2026-07-21T18:47:50.641000-04:00"
        }
    }
    

  3. Run the following command to retrieve the credentials. You’ll need these values for your ACME client configuration the next step.
    aws acm get-acme-external-account-binding-credentials \
        --acme-external-account-binding-arn arn:aws:acm:region:111122223333:acme-endpoint/00000000-0000-0000-0000-000000000000/acme-external-account-binding/22222222-2222-2222-2222-222222222222

  4. Save the KeyId and MacKey for the next step.
    {
        "KeyId": "xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx",
        "MacKey": "xxxxxxxx-xxxxxxxxxx-xxxxxxxxxxxxxxx"
    }

Step 4: Configure your ACME client

With your endpoint URL and EAB credentials ready, you can now configure your preferred ACME client. The following examples show configuration for two popular clients. As a reminder, the server information was retrieved in step 1, part 4 as the EndpointUrl.

acme.sh:

acme.sh --issue --server https://acm-acme-enroll.us-east-1.api.aws/123457-1234-1234-123456789012/directory \
    --eab-kid <KeyId> --eab-hmac-key <MacKey> \
    --email <EMAIL> \
    -d <DOMAIN> \
    --dns --yes-I-know-dns-manual-mode-enough-go-ahead-please

Certbot:

certbot certonly --standalone --non-interactive --agree-tos \
  --email <EMAIL> \
  --server https://acm-acme-enroll.us-east-1.api.aws/1234567-1234-1234-123456789012/directory \
  --eab-kid <KeyId> \
  --eab-hmac-key <MacKey> \
  -d <DOMAIN>

After the initial registration, your ACME client handles renewals.

Enterprise controls

Other ACME alternatives can provide certificates but don’t give the same amount of control and governance for customers that need to scale their certificate environment. The following controls are available to help reduce risk across your organization.

Domain validation

Customers managing large numbers of domains told us they need a way to prevent unauthorized certificate issuance across their domain space. Domain validation gives you this control. For each domain you validate, you enable the certificate patterns it should be allowed to issue, whether it’s ExactDomain, Subdomains, or Wildcards. For example, if you validate internal.example.com and enable only Wildcards, an ACME client can request *.internal.example.com but a request for internal.example.com itself or api.internal.example.com is rejected. This enforcement happens at the endpoint level, before requests reach the ACM certificate authority, and you can validate multiple domains under a single endpoint, each with its own scope.

Centralized certificate visibility

Certificates issued through your ACME endpoints are registered with ACM. You can use the aws acm list-certificates command to see all your issued certificates.

IAM authorization, CloudTrail audit logging and observability

Endpoint management operations are authorized through IAM and logged to CloudTrail. You can use IAM policies to control which principals can create endpoints, generate EAB credentials, and manage domain constraints.

Best practices

For customers implementing ACME certificates for the first time, consider the following best practices for your organizations.

Segment endpoints along organizational or environment boundaries

The endpoint serves as a useful method of isolation for larger organizations. A large enterprise can create one endpoint per organizational boundary (business unit, subsidiary, or environment) instead of a single shared endpoint company-wide. Each endpoint has its own pre-approved domains and its own set of EABs, so a compromised credential in one business unit has no path to certificates in another.

However, weigh this against your operational overhead as well. A reasonable starting point is one endpoint per environment (dev, staging, andprod) within a business unit, expanding to per-business-unit endpoints only where compliance or organizational requirements call for it.

Manage EAB credentials securely

Anyone holding a validKeyIdandMacKeyfor an endpoint can obtain certificates for any domain pre-approved on that endpoint, so these credentials deserve the same handling you’d give an access key.

  • Avoid hard coding theMacKeywhere possible by using a secret store such as AWS Secrets Manager. Distribute it only to the ACME clients that you authorize to use the endpoint.
  • Set the expiration of the EAB to an acceptable level. While EAB supports long-lived credentials, not all scenarios require an indefinitely long EAB.
  • When creating the role for each EAB, adhere to concept of least privilege. Creating a role per EAB, rather than sharing a role across all bindings, can help reduce risk in your AWS environment.
  • Audit CreateAcmeExternalAccountBinding and GetAcmeExternalAccountBindingCredentials calls in CloudTrail separately. Because retrieving the actual key material is a distinct API call from creating the binding, alerting on retrieval events is a stronger signal of real credential distribution than binding creation alone.

Automate how EABs are associated with clients at runtime

Generate a unique set of EAB credentials for each client or environment rather than sharing one binding across multiple ACME clients. As you begin to scale with multiple endpoints, usesome of the following patterns to reduce operational toil.

  • Name each EAB and its bound IAM role after the client it belongs to (team, application, environment), so the binding’s purpose is obvious from DescribeAcmeExternalAccountBinding output alone, without cross-referencing a spreadsheet.
  • Store each client’s KeyId and MacKey under a secrets path scoped to that client (for example, a Secrets Manager path per team and environment), and let the client’s provisioning pipeline retrieve its own credentials.
  • In Kubernetes, use one ClusterIssuer or namespace-scoped Issuer per EAB rather than one shared issuer across teams. This keeps the client-to-EAB association explicit in cluster config, and lets you revoke one team’s access without touching anyone else’s.
  • For ephemeral infrastructure (build agents, autoscaled fleets), provision EAB credentials as part of your infrastructure-as-code or continuous integration and deployment (CI/CD) pipeline instead of a one-time manual handoff, so credential lifecycle tracks infrastructure lifecycle.

Monitor your deployment of ACME

ACME’s power is through automation, and organizations should monitor their ACME usage for anomalies.

  • Alarm on issuance failures, not just successes. At 45-day certificate validity, a silent renewal failure gives you far less runway to react than the months of time you might be used to with longer-lived certificates.
  • Test renewal automation before you depend on it. Force a manual renewal against a non-production endpoint and confirm your client, monitoring, and on-call runbooks behave as expected, before the CA/Browser Forum’s shortened validity windows turn a failed renewal into a disruptive event for your organization.

Availability and pricing

ACME support in AWS Certificate Manager is available today in all commercial AWS Regions and will be available in AWS GovCloud (US), the China Regions, and the AWS European Sovereign Cloud partitions at a later date. See the ACM pricing page for more information on ACME pricing.

Conclusion

The phased reduction in certificate validity can’t easily be solved without automation. ACME support in ACM gives you that automation through a standard protocol and standard tooling, while keeping the visibility and governance controls your security teams rely on from ACM.

To get started, see the AWS Certificate Manager documentation or follow the getting started guide.

If you have feedback about this post, submit comments in the Comments section below.


Anthony Harvey

Anthony Harvey

Anthony is a Senior Security Specialist Solutions Architect for AWS in the worldwide public sector group. Prior to joining AWS, he was a chief information security officer in local government for half a decade. With his public sector experience, he has a passion for figuring out how to do more with less and leveraging that mindset to enable customers in their security journey.

Chandan Kundapur

Chandan Kundapur

Chandan is a Principal Product Manager on the AWS Certificate Manager (ACM) team. With over 15 years of cybersecurity experience, he has a passion for driving PKI product strategy.

Route Amazon Bedrock Guardrails interventions to Amazon Security Lake

Post Syndicated from Dhananjay Karanjkar original https://aws.amazon.com/blogs/security/route-amazon-bedrock-guardrails-interventions-to-amazon-security-lake/

Security teams investigating AI-related incidents need guardrail intervention data alongside their existing security telemetry. Routing Amazon Bedrock Guardrails violations to Amazon Security Lake makes this possible. With this integration, you can query guardrail events alongside identity, network, and application security data in a single layer. When a guardrail blocks a prompt injection attempt or redacts sensitive data, that intervention carries investigative value comparable to a failed sign-in or a network intrusion alert. Amazon Bedrock publishes this telemetry to Amazon CloudWatch metrics and model invocation logs for operational monitoring. By using Security Lake, organizations can extend this telemetry into their security data lake for unified correlation.

In this post, I show you how to build an automated pipeline that transforms Amazon Bedrock Guardrails intervention events into Open Cybersecurity Schema Framework (OCSF) records and delivers them to Security Lake as a custom source. You can query the data using Amazon Athena or any Security Lake subscriber.

Use case

Consider a financial services organization deploying Amazon Bedrock across multiple business units. Each unit uses guardrails to enforce content policies (blocking harmful content), topic policies (preventing off-topic queries about competitors), sensitive information policies (redacting personally identifiable information (PII) such as account numbers), and prompt injection detection.

The security team needs to:

  • Identify which user accounts trigger the most guardrail interventions and whether those accounts also have unusual AWS Identity and Access Management (IAM) activity
  • Determine if prompt injection attempts correlate with specific source IP addresses that also appear in Amazon Virtual Private Cloud (Amazon VPC) Flow Logs
  • Track the organization-wide trend of guardrail violations across all business units and compare it against the baseline from 30 days ago

With guardrail events routed to Security Lake, a single Athena query covers all three.

Solution overview

The pipeline architecture routes Amazon Bedrock security events to Security Lake as OCSF-compliant records. The same infrastructure—subscription filter, AWS Lambda transformation, Parquet writer, Amazon Simple Storage Service (Amazon S3) partitioning—supports multiple event types by changing the filter pattern and OCSF mapping:

Guardrail interventions (this post) DETECTION_FINDING 2004
Model invocation API calls API_ACTIVITY 6003
Agent guardrail traces DETECTION_FINDING 2004
Token consumption anomalies DETECTION_FINDING 2004

This post demonstrates the guardrail interventions implementation as a working example. The solution captures Amazon Bedrock model invocation logs that contain guardrail trace data and filters for intervention events. It transforms matching events into OCSF-compliant Detection Finding records (class_uid 2004) and delivers them to Security Lake as Parquet files. Guardrail interventions are detection events: the guardrail detected and blocked prohibited content, so OCSF class 2004 (Detection Finding) under the Findings category is the appropriate classification.

Architecture

The following diagram shows the end-to-end pipeline from guardrail intervention to Security Lake ingestion.

Figure 1: Guardrail intervention routing

Figure 1: Guardrail intervention routing

The data flow consists of the following steps:

  1. An application calls Amazon Bedrock (InvokeModel or Converse API) with a guardrail attached.
  2. Amazon Bedrock evaluates the guardrail and logs the invocation (including guardrail trace data) to a CloudWatch Logs log group using model invocation logging. The subscription filter matches log entries where the guardrail action is INTERVENED (blocked or masked content).
  3. The subscription filter delivers matching records to a Lambda function (OCSF Transform).
  4. The Lambda function transforms each intervention event into an OCSF Detection Finding record (class_uid 2004), batches records, and converts them to Zstandard (zstd)-compressed Apache Parquet format. It writes the Parquet file to the Amazon S3 Security Lake bucket using the required partition path (ext/BedrockGuardrails/region=/accountId=/eventDay=/). If the Lambda function fails to process a record, the message routes to an Amazon Simple Queue Service (Amazon SQS) dead-letter queue for later analysis and redrive.
  5. Security Lake manages the ingested Parquet data in the S3 bucket.
  6. AWS Glue crawler detects new partitions and catalogs the Parquet files for query access.
  7. SOC analysts query guardrail violation data alongside other security sources using Athena.

OCSF mapping

The following table shows how Amazon Bedrock Guardrails intervention fields map to OCSF Detection Finding (class_uid 2004) attributes.

OCSF field Source Example value
class_uid Static 2004 (Detection Finding)
category_uid Static 2 (Findings)
severity_id Derived from policy type 3 (Medium) for content/topic; 4 (High) for prompt injection
activity_id Static 1 (Create)
time Invocation log timestamp 1721001600000
cloud.provider Static AWS
cloud.region Invocation log region us-east-1
cloud.account.uid Invocation log accountId 123456789012
actor.user.uid Invocation log identity.arn arn:aws:sts::123456789012:assumed-role/AppRole/session
finding_info.title Derived from policy type ContentPolicy Intervention
finding_info.desc Guardrail trace action/topic Blocked: HATE content detected on INPUT
resource.uid Model ARN arn:aws:bedrock:us-east-1::foundation-model/anthropic.claude-sonnet-4-6-20250514-v1:0
resource.type Static AwsBedrock:Model
metadata.product.name Static Amazon Bedrock Guardrails
metadata.product.vendor_name Static AWS
metadata.version Static 1.3.0
unmapped.guardrail_id Guardrail trace guardrailId my-content-guardrail
unmapped.guardrail_arn Guardrail trace guardrailArn arn:aws:bedrock:us-east-1:123456789012:guardrail/abc123
unmapped.guardrail_version Guardrail trace guardrailVersion 3
unmapped.guardrail_content_source Guardrail trace INPUT or OUTPUT
unmapped.guardrail_policy_type Guardrail trace ContentPolicy, TopicPolicy, SensitiveInformationPolicy, WordPolicy, ContextualGroundingPolicy, PromptAttack

Prerequisites

The following prerequisites are needed to deploy the reference implementation. Before you begin, clone the repository:

git clone https://github.com/aws-samples/sample-bedrock-guardrails-security-lake.git
cd sample-bedrock-guardrails-security-lake

Verify you have the following:

  • An AWS account with AWS Cloud Development Kit (AWS CDK) bootstrapped in the target AWS Region
  • Security Lake enabled in the target Region
  • Python 3.12 or later
  • Node.js 20 or later (for AWS CDK CLI)
  • An existing Amazon Bedrock guardrail (or create one during deployment)
  • Model invocation logging enabled on Amazon Bedrock (with guardrail trace data enabled)

Implementation

The reference implementation deploys three CloudFormation stacks: SecurityLakeSourceStack, TransformPipelineStack and MonitoringStack. The following commands deploy the stacks in dependency order:

cdk deploy SecurityLakeSourceStack \
  -c security_lake_bucket=<your-security-lake-bucket> \
  -c source_location=ext/BedrockGuardrails \
  -c security_lake_enabled=true

cdk deploy TransformPipelineStack \
  -c security_lake_bucket=<your-security-lake-bucket> \
  -c source_location=ext/BedrockGuardrails

cdk deploy MonitoringStack \
  -c security_lake_bucket=<your-security-lake-bucket> \
  -c source_location=ext/BedrockGuardrails

Enable model invocation logging

Model invocation logging captures the guardrail trace data you need. Turn on full request and response logging to a CloudWatch Logs log group. Configure textDataDeliveryEnabled to capture text request and response bodies, which include the guardrail trace output when a guardrail is attached to the invocation.

Register Security Lake custom source

Register BedrockGuardrails as a custom source with Security Lake using the DETECTION_FINDING event class. Security Lake creates the Amazon S3 prefix and IAM role for your source. The stack configures the AWS Glue crawler role for partition discovery.

Create the subscription filter

Create a CloudWatch Logs subscription filter on your model invocation log group with the filter pattern { $.output.guardrailAction = “INTERVENED” }. This captures only the events where a guardrail blocked or modified content, not the successful pass-through events. This reduces Lambda invocations and cost.

Transform to OCSF and write Parquet

The Lambda function performs three operations: parse the CloudWatch Logs event, transform each intervention to an OCSF Detection Finding record (class_uid 2004), and write batched records as Parquet files. The files are written to the Security Lake S3 bucket using the required partition path (ext/BedrockGuardrails/region=<region>/accountId=<accountId>/eventDay=<YYYYMMDD>/).

The transformation maps guardrail trace fields to OCSF attributes as described in the OCSF mapping table. Severity is set to High for prompt injection interventions and Medium for content, topic, or sensitive information interventions. For a concrete before-and-after example, see the sample invocation log and corresponding OCSF output in the companion repository.

Scaling considerations: At low intervention volumes (tens of events per hour), direct Lambda writes produce acceptably sized Parquet files. For higher volumes, consider buffering through Amazon Data Firehose with its native Parquet conversion and 5-minute buffering interval to produce fewer, larger files that optimize Athena query performance.

Multi-account deployment: The partition scheme (accountId=<account>) already supports multi-account environments. Deploy the subscription filter and transform pipeline in each workload account where model invocation logging is enabled. Each pipeline writes cross-account to the delegated-administrator Security Lake bucket. Distribute the pipeline using CloudFormation StackSets across the organization.

Query violations in Athena

After deployment, guardrail violations typically appear in your Security Lake tables within 5–10 minutes, depending on the AWS Glue crawler schedule. You can then run cross-service correlation queries. The following example identifies users who trigger both prompt injection interventions and unusual IAM activity:

WITH guardrail_violators AS (
    SELECT actor.user.uid AS user_arn, COUNT(*) AS violation_count
    FROM "amazon_security_lake_glue_db_us_east_1"."amazon_security_lake_table_us_east_1_bedrockguardrails"
    WHERE eventDay >= '20260701'
      AND unmapped.guardrail_policy_type = 'PromptAttack'
    GROUP BY actor.user.uid
),
iam_failures AS (
    SELECT actor.user.uid AS user_arn, COUNT(*) AS failure_count
    FROM "amazon_security_lake_glue_db_us_east_1"."amazon_security_lake_table_us_east_1_cloud_trail_mgmt_2_0"
    WHERE eventDay >= '20260701'
      AND status_id = 2
    GROUP BY actor.user.uid
)
SELECT g.user_arn, g.violation_count, i.failure_count
FROM guardrail_violators g
JOIN iam_failures i ON g.user_arn = i.user_arn
ORDER BY g.violation_count DESC;

You can also track violation trends by policy type over time to establish baselines and detect spikes. The following query shows the 30-day trend:

SELECT eventDay,
       unmapped.guardrail_policy_type AS policy_type,
       COUNT(*) AS violation_count
FROM "amazon_security_lake_glue_db_us_east_1"."amazon_security_lake_table_us_east_1_bedrockguardrails"
WHERE eventDay >= '20260623'
GROUP BY eventDay, unmapped.guardrail_policy_type
ORDER BY eventDay, violation_count DESC;

The OCSF mapping has been validated against schema version 1.3.0, and the Security Lake AWS Glue crawler correctly detects the partitioned Parquet files for querying.

Alternative for teams not yet using Security Lake: If your organization hasn’t adopted Security Lake, you can query guardrail intervention events directly in CloudWatch Logs Insights using the same subscription filter log group. CloudWatch Logs Insights supports cross-log-group queries, so you can correlate guardrail events with other CloudWatch log sources without the OCSF transformation step. Security Lake adds value when you need to join with non-CloudWatch sources in a single query layer. Examples include Amazon VPC Flow Logs, Amazon Route 53 DNS logs, and third-party findings.

Clean up

To avoid ongoing charges, destroy the stacks in reverse dependency order:

cdk destroy MonitoringStack --force \
  -c security_lake_bucket=<your-security-lake-bucket> \
  -c source_location=ext/BedrockGuardrails

cdk destroy TransformPipelineStack --force \
  -c security_lake_bucket=<your-security-lake-bucket> \
  -c source_location=ext/BedrockGuardrails

cdk destroy SecurityLakeSourceStack --force \
  -c security_lake_bucket=<your-security-lake-bucket> \
  -c source_location=ext/BedrockGuardrails \
  -c security_lake_enabled=true

Conclusion

In this post, you learned how to route Amazon Bedrock Guardrails intervention events to Amazon Security Lake as OCSF-compliant Detection Finding records. This integration extends guardrail telemetry from Amazon CloudWatch into your security data lake. Security analysts can then run cross-service correlation of AI intervention events with IAM, network, and application telemetry.

The pipeline filters for intervention events only, keeping costs low while capturing the security-relevant signals. The records use OCSF event class 2004 (Detection Finding), which integrates with supported Security Lake subscribers such as Amazon OpenSearch Service and third-party SIEM tools.

Clone the reference implementation and adapt the OCSF mapping and subscription filter to your organization’s guardrail configuration.

References

If you have feedback about this post, submit comments in the Comments section below.


Dhananjay Karanjkar

Dhananjay Karanjkar

Dhananjay is a Senior Lead Consultant at AWS Professional Services, specializing in agentic AI systems, multi-agent orchestration, and generative AI security. He holds two US patents and serves as a Responsible AI Champion, with a background spanning financial services, enterprise consulting, and enterprise-scale AI delivery. When not architecting AI solutions, he trains for triathlons, paints oil portraits, and is an avid reader.

Scaling fine-grained access control for enterprise lakehouse using SageMaker Unified Studio and AWS Lake Formation

Post Syndicated from Chintan Agrawal original https://aws.amazon.com/blogs/big-data/scaling-fine-grained-access-control-for-enterprise-lakehouse-using-sagemaker-unified-studio-and-aws-lake-formation/

As enterprise lakehouses grow to thousands of tables across multiple business domains and regions, scaling fine-grained access control becomes a critical governance challenge. Data governance teams spend significant time manually granting table-level permissions, only to face permission drift, inconsistent enforcement, and limited auditability. Without a scalable approach, each new dataset requires manual policy updates, increasing the risk of unauthorized access and slowing time-to-insight for analysts and data scientists.

In this post, we show you how to solve this problem by combining AWS IAM Identity Center, AWS Lake Formation tag-based access control (TBAC), and trusted identity propagation in Amazon SageMaker Unified Studio. You deploy a complete governance architecture using AWS Cloud Development Kit (AWS CDK) that classifies data with LF-Tags, maps IAM Identity Center groups to tag-based policies, and enforces permissions at query time across analytics engines. The solution uses Apache Iceberg tables stored in Amazon Simple Storage Service (Amazon S3) and registered in the AWS Glue Data Catalog.

The core governance challenge

As organizations mature their lakehouse environments, governance complexity increases with each new dataset. Several challenges commonly emerge:

  • Explosive dataset growth: Iceberg-based lakehouses often contain thousands of tables distributed across raw, curated, and conformed zones. Each new dataset introduces additional governance requirements, making table-level permission grants operationally expensive.
  • Multi-domain data ownership: Enterprise lakehouses typically serve multiple business domains such as commercial analytics, clinical research, and regulatory reporting. These domains require strict isolation while still supporting controlled data sharing.
  • Regional data sovereignty: Organizations operating globally must enforce geographic boundaries for sensitive datasets. EU clinical trial data might be restricted by GDPR regulations, whereas US commercial datasets follow different compliance frameworks.
  • Sensitivity-based access controls: Within each domain, datasets vary in sensitivity. Pricing strategies, drug discovery research, and patient-related datasets require stricter access controls than standard operational data.
  • Role explosion: Pure RBAC approaches attempt to encode these dimensions into roles, leading to role proliferation. Manual Lake Formation grants at the table level create permission drift and limited scalability.

To address these challenges, enterprise lakehouse governance must satisfy several criteria:

  • Least-privilege access.
  • Dynamic scalability as new datasets are onboarded.
  • Multi-dimensional enforcement across domain, region, and sensitivity.
  • Auditability traceable to individual users.
  • Automation-ready, configuration-driven workflows.

TBAC addresses each of these challenges directly. Instead of granting permissions on individual tables, you define tag-based policies that automatically apply to any resource matching the tag expression. New datasets inherit access rules through tag inheritance, eliminating manual policy updates (solving explosive dataset growth). Domain and region tags enforce strict isolation between business units (solving multi-domain ownership and regional sovereignty). Sensitivity tags control access within domains without role proliferation (solving sensitivity-based controls and role explosion). The following sections describe the architecture that implements this model and walk you through deploying it end to end.

Reference architecture overview

The governance model integrates identity, metadata, and lakehouse services into a unified access architecture that enforces fine-grained permissions consistently across analytics and machine learning (ML) workloads. The architecture consists of five layers, each handling a distinct responsibility in the access control flow.

The following diagram illustrates the end-to-end architecture, showing how user identity flows from IAM Identity Center through SageMaker Unified Studio to Lake Formation for tag-based policy evaluation against the AWS Glue Data Catalog and Amazon S3 storage layer.

Architecture linking IAM Identity Center, SageMaker Unified Studio, Lake Formation, the Glue Data Catalog, and Amazon S3

Figure 1: End-to-end governance architecture for the enterprise lakehouse

1. Identity and authentication layer: IAM Identity Center manages user identities and group memberships, integrates with corporate identity providers, and provides centralized lifecycle management for enterprise users. IAM Identity Center groups represent business roles and serve as the principals that receive Lake Formation permissions.

2. Unified analytics and ML access layer: Amazon SageMaker Unified Studio serves as the primary interface where analysts, data scientists, and ML engineers discover datasets, run queries, and build ML workflows. Because SageMaker Unified Studio integrates with multiple compute engines, including Amazon Athena, AWS Glue, Amazon EMR, and Amazon Redshift, users can access data using their preferred analytics tools while maintaining consistent governance.

3. Governance and authorization layer: AWS Lake Formation provides fine-grained access control across AWS Glue catalog resources using LF-Tags. Instead of granting permissions directly on databases and tables, Lake Formation evaluates LF-Tag policies dynamically and grants or denies access at query time. Governance teams define access rules once, and Lake Formation automatically applies them to new datasets as they are onboarded.

4. Governance automation layer: Two AWS Lambda functions automate tag assignment and permission provisioning. JSON metadata configuration files drive both pipelines, so governance teams manage access control through configuration rather than manual console operations.

5. Metadata and storage layer: Apache Iceberg tables stored in Amazon S3 form the foundation of the lakehouse. You register these tables in the AWS Glue Data Catalog, which provides centralized metadata management and interoperability across analytics services. Lake Formation evaluates governance decisions at the catalog level rather than independently by each analytics engine.

End-to-end access flow

When a user queries a dataset from SageMaker Unified Studio, the following sequence occurs:

  1. The user authenticates through IAM Identity Center and accesses SageMaker Unified Studio.
  2. SageMaker passes the user’s identity context to downstream analytics services using trusted identity propagation.
  3. The analytics engine requests data access from Lake Formation.
  4. Lake Formation evaluates LF-Tag policies against the user’s IAM Identity Center group membership.
  5. Access is granted or denied dynamically at query time.

Because authorization decisions are centralized in Lake Formation, governance remains consistent regardless of which analytics engine the user employs.

Hybrid RBAC + ABAC governance model

The governance model combines identity context from IAM Identity Center with metadata-driven classification using LF-Tags. The following table summarizes how each layer contributes to the overall governance workflow.

Governance capability IAM Identity Center contribution Lake Formation LF-Tag contribution Governance outcome
Identity context Organizes users into groups aligned with business roles Evaluates permissions using group membership Role-aligned access boundaries
Data classification Provides role eligibility for data access Classifies datasets by domain, region, sensitivity, and layer Attribute-aware authorization
Scalability Simplifies user lifecycle management Automatically applies policies to newly tagged datasets Governance that scales with dataset growth
Operational model Centralizes role lifecycle operations Enables metadata-driven policy automation Reduced administrative overhead

IAM Identity Center defines who can request access, LF-Tags define what datasets are eligible, and Lake Formation enforces policies dynamically at query time.

Enterprise LF-Tag data model

A structured tagging strategy is the foundation of scalable Lake Formation governance. In this solution, the solution classifies datasets across four governance dimensions.

Tag Key Tag Values Purpose Example Usage
region us, eu, global Geographic data location Enforce GDPR compliance for EU data
domain commercial, clinical_research, regulatory Business domain Separate commercial from clinical data
data_class standard, sensitive, regulated Data sensitivity level Restrict access to sensitive pricing data
layer raw, curated, conformed Data processing stage Grant analysts access to curated data only

Together, these dimensions enable multi-dimensional authorization policies that reflect both organizational structure and regulatory requirements.

Tag inheritance and evaluation

LF-Tags can be applied at three resource levels within the Glue Data Catalog: database, table, and column. In this implementation, database-level tags define broad governance attributes (domain, region, layer), table-level tags capture dataset-specific sensitivity (data_class), and column-level tags can further restrict access to individual fields. Lake Formation evaluates the effective tag set at query time by combining inherited and explicitly assigned tags.

For example, a database tagged domain=commercial, region=us, layer=raw automatically applies those tags to all tables within it. A table-level data_class=sensitive tag supplements the inherited tags to distinguish sensitive pricing data from standard sales data. This inheritance model means new tables automatically receive governance coverage without manual tag assignment. To learn more, refer to Lake Formation tag-based access control best practices.

Prerequisites

Before deploying the solution, complete the following setup in the us-east-1 Region. Use the same AWS Region throughout all steps.

  1. AWS account and IAM Identity Center: Enable IAM Identity Center and create test users. Note your Identity Store ID from the IAM Identity Center console under Settings. For setup guidance, see Getting started with IAM Identity Center.
  2. Lake Formation configuration: Complete the following setup in the Lake Formation console:2.1. Change Data Catalog default permissions. In the navigation pane under Administration, choose Data Catalog settings. Uncheck Use only IAM access control for new databases and uncheck Use only IAM access control for new tables in new databases. Choose Save. This makes sure Lake Formation permissions govern access to databases and tables created by the CDK stacks.
    Lake Formation Data Catalog settings with both IAM-only access control checkboxes cleared

    Figure 2: Lake Formation Data Catalog settings with both IAM-only access control checkboxes unchecked

    2.2. Integrate with IAM Identity Center. Complete the prerequisites for IAM Identity Center integration with Lake Formation, including enabling trusted identity propagation.You don’t need to manually create a Lake Formation administrator. The CDK deployment in Step 2: Deploy all stacks automatically registers the required administrators via the LfAdminStack (see lf-admin-stack.ts). S3 data location registration is a post-deployment console step covered after the CDK creates the buckets.

  3. SageMaker Unified Studio: Create a SageMaker Unified Studio domain, select your IAM Identity Center instance for authentication, and enable trusted identity propagation. For a detailed walkthrough, see Accelerate your analytics with Amazon S3 Tables and Amazon SageMaker Lakehouse and enable trusted identity propagation for the domain.
  4. Local tooling: Install AWS Command Line Interface (AWS CLI), Python 3.x, Node.js 18+, AWS CDK CLI (npm install -g aws-cdk), and Git.

Solution overview

Now that you understand the governance model and tag taxonomy, the following section walks you through deploying the complete infrastructure and configuring access control.

The deployment uses AWS CDK (TypeScript) and consists of seven stacks that create the complete governance infrastructure. The CDK app manages stack dependencies automatically, so a single cdk deploy --all command deploys everything in the correct order.

The architecture uses a two-layer data lake pattern. The raw layer stores data as CSV files in Amazon S3, registered as external tables in the AWS Glue Data Catalog. The curated layer uses Apache Iceberg v2 tables for ACID transactions and schema evolution. Three business domains (US Commercial, EU Clinical Research, and Global Regulatory) each have one representative table per layer, giving six tables total.

Lake Formation tag-based access control (TBAC) governs all access using four tag dimensions:

Tag Key Values Purpose
domain commercial, clinical_research, regulatory Business domain isolation
region us, eu Geographic data boundary
data_class standard, sensitive, regulated Sensitivity classification
layer raw, curated Data layer identification

Step 1: Clone the repository and install dependencies

Clone the accompanying repository and install the CDK project dependencies:

git clone https://github.com/aws-samples/sample-aws-smus-governance-automation
cd aws-smus-governance-automation/cdk
npm install

The CDK project is written in TypeScript and uses aws-cdk-lib v2. The lib/ directory contains seven stack definitions, and bin/app.ts wires them together with explicit dependency ordering.

If this is your first CDK deployment in this account and Region, bootstrap the CDK environment. Bootstrapping provisions an S3 bucket and IAM roles that CDK uses to deploy assets:

cdk bootstrap aws://<ACCOUNT_ID>/us-east-1

Step 2: Deploy all stacks

Deploy the entire infrastructure with a single command. Pass your IAM Identity Center Identity Store ID as a CDK context variable:

cdk deploy --all -c identityStoreId=d-xxxxxxxxxx --require-approval never --region us-east-1

CDK will prompt for IAM permission changes on each stack. The --require-approval never flag auto-approves these so the deployment runs unattended.

CDK deploys the seven stacks in dependency order:

  1. LfSetupStack: Lake Formation admin registration + LF-Tags (domain, region, data_class, layer)
  2. GlueRawTablesStack: S3 bucket + three Glue databases + three CSV-backed tables.
  3. GlueCuratedTablesStack: S3 bucket + three Glue databases + three Iceberg v2 tables.
  4. SsoGroupsStack: three IAM Identity Center groups (DataLake-US-Commercial, DataLake-EU-Clinical-Research-Sensitive, DataLake-Regulatory)The three groups map to specific tag combinations that control data access:
    • DataLake-US-Commercial: domain=commercial, region=us, data_class=standard.
    • DataLake-EU-Clinical-Research-Sensitive: domain=clinical_research, region=eu, data_class=sensitive,regulated.
    • DataLake-Regulatory: domain=regulatory (all regions, all data classes within regulatory).

    The following table summarizes the user personas, their group assignments, and the data access each group provides:

  5. AssetTaggingAutomationStack: Tag automation Lambda.
  6. SsoPermissionAutomationStack: Permission automation Lambda.
  7. LfAdminStack: Registers CDK + Lambda roles as Lake Formation admins.

After deployment completes, review the CloudFormation stack outputs. They include S3 bucket names, database names, SSO group IDs, and Lambda function ARNs.

The following figure shows all seven CDK stacks deployed successfully in the CloudFormation console.

CloudFormation console showing all seven CDK stacks in CREATE_COMPLETE status

Figure 3: CloudFormation console showing all seven CDK stacks in CREATE_COMPLETE status

Register S3 data locations with Lake Formation: Now that the S3 buckets exist, register them with Lake Formation. In the Lake Formation console, under Administration, choose Data lake locations, then choose Register location. Register both buckets from the stack outputs (for example, s3://datalake-raw-data-<ACCOUNT_ID>-us-east-1 and s3://datalake-curated-data-<ACCOUNT_ID>-us-east-1). For IAM role, use the default AWSServiceRoleForLakeFormationDataAccess and choose Lake Formation as the permission mode. See Registering an Amazon S3 location for step-by-step instructions.

The following figure shows both data lake S3 locations registered in the Lake Formation console.

Lake Formation Data lake locations page listing the registered raw and curated S3 buckets

Figure 4: Lake Formation Data lake locations page with raw and curated S3 buckets registered

Step 3: Populate sample datasets

The scripts use Amazon Athena to insert sample data. Athena stores query results under the athena-results/ prefix in the shared governance metadata bucket (lf-governance-metadata-<ACCOUNT_ID>-<REGION>) created by the CDK deployment.

Populate the raw and curated tables:

cd ../scripts
python3 populate_raw_layer.py
python3 populate_curated_layer.py

Each script executes INSERT INTO statements through the Athena StartQueryExecution API and waits for completion. You should see success messages for all six tables (three raw, three curated).

After populating the tables, you can verify the data in the Glue Data Catalog. The following figure shows the six tables across the three raw and three curated databases.

AWS Glue Data Catalog showing the six databases and tables created by the deployment

Figure 5: AWS Glue Data Catalog showing the six databases and tables created by the CDK deployment

You can also preview the data by querying a table. The following figure shows sample data from the us_sales_summary table.

Athena query results showing sample commercial rows from the us_sales_summary table

Figure 6: Query results for the us_sales_summary table with sample commercial data

Step 4: Apply LF-Tags to data assets

The following diagram illustrates the governance automation flow, showing how metadata JSON configuration files drive the two Lambda pipelines for asset tagging and SSO permission management.

Governance automation flow with the asset tagging and SSO permission Lambda pipelines

Figure 7: Governance automation flow showing the asset tagging and SSO permission Lambda pipelines

The diagram shows two parallel pipelines, each following three steps:

Asset tagging pipeline (left):

  1. Metadata upload – A data governance administrator uploads metadata JSON files (metadata-raw-tables.json and metadata-curated-tables.json) to the asset-tagging/ prefix in the shared S3 governance metadata bucket. These files define which LF-Tags to assign to each AWS Glue database and table.
  2. Lambda processing – The S3 upload triggers the LakeFormationTagAutomation Lambda function, which reads the metadata and calls the Lake Formation API.
  3. Tag operations – The Lambda creates or updates LF-Tags, then assigns them to the target databases and tables in the AWS Glue Data Catalog.

SSO permission pipeline (right):

  1. Permission upload – Three permission JSON files (one per IAM Identity Center group) are uploaded to the sso-permissions/ prefix. These files define the LF-Tag policy expressions that control data access.
  2. Lambda processing – The upload triggers the LakeFormationSSOPermissionAutomation Lambda function.
  3. Permission operations – The Lambda grants tag-based permissions to the corresponding IAM Identity Center groups through the Lake Formation API.

Both pipelines log execution details to Amazon CloudWatch for monitoring and troubleshooting.

Two metadata JSON configuration files drive the asset tagging Lambda that declaratively define which LF-Tags to apply to each AWS Glue resource:

  • metadata-raw-tables.json: Tag definitions for the three raw layer databases and tables.
  • metadata-curated-tables.json: Tag definitions for the three curated layer databases and tables.

Each entry in these files specifies the following fields:

Field Description Example
catalog_id Your AWS account ID (Glue Data Catalog ID) 123456789012
resource_type DATABASE or TABLE DATABASE
database_name AWS Glue database name raw_us_commercial_db
table_name AWS Glue table name (only for TABLE entries) us_sales_summary
lf_tags Array of LF-Tag key/value pairs to assign [{“TagKey”:“domain”,“TagValues”:[“commercial”]}]
access_type Action to perform (GRANT) GRANT

Parameters you must update before invoking: Replace the catalog_id value in every entry of both files with your own AWS account ID. The database and table names match the resources created by the CDK stacks, so those should not be changed unless you customized the stack parameters.

The following snippet from metadata-raw-tables.json shows a database-level entry and a table-level entry:

[
  {
    "comment": "DATABASE LEVEL TAGS - US Commercial RAW Domain",
    "access_type": "GRANT",
    "resource_type": "DATABASE",
    "catalog_id": "<YOUR_ACCOUNT_ID>",
    "database_name": "raw_us_commercial_db",
    "lf_tags": [
      { "TagKey": "region", "TagValues": ["us"] },
      { "TagKey": "domain", "TagValues": ["commercial"] },
      { "TagKey": "layer", "TagValues": ["raw"] }
    ]
  },
  {
    "comment": "TABLE LEVEL TAGS - US Commercial RAW Table (Standard Access)",
    "access_type": "GRANT",
    "resource_type": "TABLE",
    "catalog_id": "<YOUR_ACCOUNT_ID>",
    "database_name": "raw_us_commercial_db",
    "table_name": "us_sales_summary",
    "lf_tags": [
      { "TagKey": "data_class", "TagValues": ["standard"] }
    ]
  }
]

The Lambda applies tags at two levels: database-level entries assign domain, region, and layer tags, while table-level entries assign the data_class tag (standard, sensitive, or regulated). Because of two-level tagging, new tables added to a tagged database automatically inherit the database-level tags. Only the table-specific data_class tag needs explicit assignment. To learn more about this pattern, refer to Lake Formation tag-based access control best practices.

Invoke the Lambda for both layers:

cd ../lf-asset-tagging-automation
aws lambda invoke \
    --function-name LakeFormationTagAutomation \
    --payload fileb://metadata-raw-tables.json \
    --cli-binary-format raw-in-base64-out \
    response.json

aws lambda invoke \
    --function-name LakeFormationTagAutomation \
    --payload fileb://metadata-curated-tables.json \
    --cli-binary-format raw-in-base64-out \
    response.json

Verify tag assignment using the GetResourceLFTags API:

aws lakeformation get-resource-lf-tags \
    --resource '{"Table":{"DatabaseName":"raw_us_commercial_db","Name":"us_sales_summary"}}' \
    --region us-east-1

You should see domain=commercial, region=us, layer=raw, and data_class=standard in the response.

The following figure shows the LF-Tags assigned to the us_sales_summary table in the Lake Formation console, confirming that both database-level inherited tags and table-level tags are applied correctly.

Lake Formation console showing inherited and table-level LF-Tags on the us_sales_summary table

Figure 8: LF-Tags on the us_sales_summary table showing inherited and table-level tags

Step 5: Provision SSO group permissions

Three permission JSON files (one per IAM Identity Center group) define the LF-Tag policy expressions. Update sso_group with the group UUID from the SsoGroupsStack outputs and identity_center_account_id with your AWS account ID. For detailed configuration, see the repository README.

[
  {
    "sso_name": "DataLake-US-Commercial",
    "sso_group": "<GROUP_UUID_FROM_CDK_OUTPUT>",
    "identity_center_account_id": "<YOUR_ACCOUNT_ID>",
    "resources": [
      {
        "resource_type": "DATABASE",
        "permissions": ["DESCRIBE"],
        "lf_tag_expression": [
          { "TagKey": "domain", "TagValues": ["commercial"] },
          { "TagKey": "region", "TagValues": ["us"] },
          { "TagKey": "layer", "TagValues": ["curated", "raw"] }
        ]
      },
      {
        "resource_type": "TABLE",
        "permissions": ["SELECT", "DESCRIBE"],
        "lf_tag_expression": [
          { "TagKey": "domain", "TagValues": ["commercial"] },
          { "TagKey": "region", "TagValues": ["us"] },
          { "TagKey": "data_class", "TagValues": ["standard"] },
          { "TagKey": "layer", "TagValues": ["curated", "raw"] }
        ]
      }
    ]
  }
]

Apply permissions for each group:

cd ../lf-sso-permission-automation
python3 lambda_function.py us-commercial-permissions.json
python3 lambda_function.py eu-clinical-research-sensitive-permissions.json
python3 lambda_function.py regulatory-permissions.json
aws lakeformation list-permissions \
    --principal '{"DataLakePrincipalIdentifier":"arn:aws:identitystore:::group/<GROUP_ID>"}' \
    --region us-east-1

Step 6: Validate fine-grained access control

With all permissions in place, validate that Lake Formation TBAC enforces the correct access boundaries by signing in to SageMaker Unified Studio as different IAM Identity Center users.

Test as Sarah (US Commercial Analyst) — Sarah belongs to DataLake-US-Commercial, which grants access to standard commercial data only.

SELECT * FROM raw_us_commercial_db.us_sales_summary LIMIT 10;

Sarah sees all rows and columns successfully:

SageMaker Unified Studio results: Sarah’s successful query on us_sales_summary

Figure 9: Sarah’s successful query on us_sales_summary in SageMaker Unified Studio

Querying outside her authorized domain returns an access denied error:

SELECT * FROM raw_eu_clinical_research_db.eu_drug_discovery LIMIT 10;
Access denied error when Sarah queries eu_drug_discovery outside her domain

Figure 10: Access denied when Sarah queries eu_drug_discovery, confirming TBAC enforcement

Test as Dr. Chen (EU Clinical Research Lead) — Dr. Chen can access sensitive and regulated EU clinical research data (eu_drug_discovery) but is denied access to US commercial data (us_sales_summary), confirming regional and domain isolation.

Query results showing Dr. Chen’s successful query on eu_drug_discovery

Figure 11: Dr. Chen’s successful query on eu_drug_discovery

Access denied error when Dr. Chen queries us_sales_summary

Figure 12: Access denied when Dr. Chen queries us_sales_summary

Test as Alex (Regulatory Affairs Specialist) — Alex’s tag expression uses only domain=regulatory without a region constraint, granting cross-regional access to regulatory data while maintaining strict isolation from commercial and clinical research domains.

Query results showing Alex’s successful query on fda_submissions

Figure 13: Alex’s successful query on fda_submissions

Access denied error when Alex queries us_sales_summary

Figure 14: Access denied when Alex queries us_sales_summary

These tests demonstrate that TBAC enforces fine-grained permissions based on user identity, data classification, regional boundaries, and domain separation, without per-table permission grants. As new tables are added and tagged, existing groups automatically gain or are denied access based on their tag expressions. This is the core advantage of TBAC over named resource permissions.

Audit user access with CloudTrail

A key benefit of integrating Lake Formation with IAM Identity Center is the detailed audit trail available through AWS CloudTrail. Filter Event history by Event name GetDataAccess to see every data access event. Each record includes the IAM Identity Center user UUID (userIdentity.onBehalfOf.userId), the specific table accessed (requestParameters.tableArn), and confirmation that trusted identity propagation was used (additionalEventData.LakeFormationTrustedCallerInvocation: true).

CloudTrail GetDataAccess event showing Identity Center user identity and table access details

Figure 15: CloudTrail GetDataAccess event showing Identity Center user identity and table access details

To resolve the user UUID to a human-readable name, query the Identity Store:

aws identitystore describe-user \
    --identity-store-id d-xxxxxxxxxx \
    --user-id <USER_UUID_FROM_EVENT> \
    --region us-east-1

This audit capability provides the detailed access logs required for HIPAA, GDPR, and FDA compliance, showing exactly which users accessed which data and when. Learn about configuring CloudTrail for Lake Formation in Logging Lake Formation API calls with CloudTrail.

Cleanup

Run cdk destroy --all to remove all stacks. Manually delete the retained S3 data buckets (datalake-raw-data-* and datalake-curated-data-*) and revoke any remaining Lake Formation permissions. For detailed cleanup steps, see the repository README.

Conclusion

In this post, we showed you how to implement scalable fine-grained access control for an enterprise lakehouse by combining AWS Lake Formation tag-based access control, IAM Identity Center, and trusted identity propagation in SageMaker Unified Studio. The four-dimension LF-Tag taxonomy, hybrid RBAC + ABAC governance model, and metadata-driven Lambda automation together create a governance architecture where new datasets automatically inherit access policies through tag inheritance, permissions scale without per-table grants, and every data access event is auditable to the individual user through CloudTrail.

To extend this solution, consider adding new business domains, implementing column-level security with LF-Tags, scaling to multi-account architectures with Lake Formation cross-account sharing, or integrating additional analytics services such as Amazon Redshift Spectrum or Amazon EMR.

Get started by deploying the CDK stacks from the accompanying repository. To learn more:


About the authors

Chintan Agrawal

Chintan Agrawal

Chintan is a Solutions Architect with over 7 years of experience, with a specialization in Analytics and Healthcare domain. He possesses a strong enthusiasm for assisting clients in discovering valuable insights from their data. Through his expertise, he constructs innovative solutions that empower businesses to arrive at informed, data-driven choices.

Chaitanya Vejendla

Chaitanya Vejendla

Chaitanya is a Senior Solutions Architect and part of Global Healthcare and Life Sciences industry division at AWS. He focuses on developing strategic plans for building an end-to-end analytical strategy for large biopharma, healthcare, and life sciences organizations. His expertise spans across data analytics, data governance, AI, ML, big data, and healthcare-related technologies.

Event-driven pipeline orchestration with Amazon MWAA and Airflow 3.0

Post Syndicated from Satya Chikkala original https://aws.amazon.com/blogs/big-data/event-driven-pipeline-orchestration-with-amazon-mwaa-and-airflow-3-0/

Data engineering teams running Apache Airflow across multiple AWS accounts face a persistent coordination problem. They have no built-in way to coordinate workflows between their separate Amazon Managed Workflows for Apache Airflow (Amazon MWAA) environments, where each team or business unit manages its own isolated environment. Cross-environment orchestration has traditionally relied on time-based polling, complex custom sensors, or API-based triggers that introduce latency and reliability concerns. The Apache Airflow Datasets feature (introduced in version 2.4) added data-aware scheduling of Directed Acyclic Graphs (DAGs, the workflow definitions that specify tasks and their execution order) within a single Amazon MWAA environment. However, teams running Airflow across multiple accounts still had no way to coordinate workflows between environments.

With Apache Airflow 3.0, now available on Amazon MWAA 3.0, you get event-driven cross-account orchestration that responds to upstream events as they happen, without polling overhead or tight environment coupling. Using Amazon Simple Queue Service (Amazon SQS) as the message broker, Asset Watchers replace polling-based sensors with event-driven triggers. This approach reduces orchestration latency from minutes to seconds and reclaims worker resources previously consumed by polling sensors. It also improves message reliability, because Amazon SQS retains coordination signals even when the consumer environment is temporarily unavailable.

In this post, you learn how to design and deploy cross-account orchestration patterns using asset-based scheduling in Airflow 3.0 with Amazon SQS integration. You learn about Asset Watchers, how to publish asset events from producer DAGs, and how to trigger dependent workflows in downstream Amazon MWAA environments, creating responsive, decoupled pipelines that span multiple accounts.

If you use AI coding assistants to build and deploy infrastructure, the solution repository includes an agent skill built on the Agent Skills standard that encodes the architecture and best practices from this post.

Solution overview

This solution demonstrates a multi-MWAA orchestration architecture where:

  1. Producer Amazon MWAA Environment (Account A) runs data processing workflows that publish asset events to an Amazon SQS queue when datasets are created or updated.
  2. Amazon SQS Queue acts as a message broker, decoupling producer and consumer environments.
  3. Consumer Amazon MWAA Environment (Account B) monitors the Amazon SQS queue using Asset Watchers and automatically triggers downstream DAGs when relevant asset events arrive.

Key benefits

This event-driven approach offers several advantages over traditional polling:

  • No more polling overhead: You replace continuous sensor polling with event-driven Asset Watchers that respond as events arrive.
  • Near real-time response: Downstream DAGs trigger within seconds rather than waiting for a scheduled polling interval.
  • Independent environments: Producer and consumer Amazon MWAA environments have no direct dependencies, so each team can scale and update their environment without affecting the other.
  • Reliable message delivery: Amazon SQS provides durable message delivery, even if the consumer environment is temporarily unavailable.
  • Clear team ownership: You and your team maintain your own Amazon MWAA environment while still coordinating complex cross-account workflows.
  • Faster implementation: Describe requirements in natural language and the agent skill generates deployment-ready producer and consumer DAGs with the best practices from this post built in.

Architecture overview

The following architecture shows how you can connect separate Amazon MWAA environments across AWS accounts so that a completed pipeline in one environment automatically triggers dependent workflows in another, without direct environment coupling or polling overhead.

Producer Amazon MWAA environment publishing asset events to an Amazon SQS queue that a consumer environment monitors with an Asset Watcher to trigger downstream DAGs

Figure 1: Cross-account event-driven orchestration between Amazon MWAA environments using Amazon SQS

Architecture components

The architecture has four main components. The producer DAG defines assets as outlets and publishes events to an Amazon SQS queue when tasks complete successfully. The Amazon SQS queue acts as a durable message broker between accounts, with AWS Identity and Access Management (IAM) policies granting the producer permission to send messages and the consumer permission to receive them. On the consumer side, an Asset Watcher monitors the queue and updates asset state when messages arrive, which automatically triggers the consumer DAG scheduled on that asset.

Prerequisites

Before implementing this solution, you need:

  • Two Amazon MWAA environments running Apache Airflow 3.0 or later, in the same or different AWS accounts. Each environment must have the triggerer component enabled.
  • Intermediate knowledge of IAM policies, including cross-account role trust relationships and resource-based policies.
  • Intermediate knowledge of Apache Airflow DAG authoring, including Python-based DAG definitions and task operators.
  • Basic Python experience (Python 3.8 or later) to read and adapt the provided code samples.
  • An Amazon SQS standard queue with cross-account permissions configured (see the Cross-account IAM section).
  • AWS Command Line Interface (AWS CLI) configured with credentials that have permission to access both Amazon MWAA environments and the Amazon SQS queue.
  • Time to complete: Approximately 90 minutes (following the GitHub repository instructions).
  • Estimated cost: Running two Amazon MWAA environments and an Amazon SQS queue will incur AWS charges. Refer to the Amazon MWAA pricing page and Amazon SQS pricing page to estimate costs for your Region and usage. Remember to delete resources when you finish to avoid ongoing charges.

Implementation

The post includes a GitHub repository where you can deploy the solution described in this post. You will follow the implementation steps from setting up Amazon MWAA environments and cross-account Amazon SQS queues to deploying producer and consumer DAGs with Asset Watchers. This post provides the code samples, including the DAG files, IAM policies, and requirements configuration, for demonstration purposes only. Before deploying to production, verify that you conduct thorough testing, security reviews, and validation against the specific requirements and compliance standards.

Considerations

  • Asset Watchers run as background processes in the Airflow triggerer, not the scheduler. Verify that the triggerer is healthy and running in consumer Amazon MWAA environment before expecting event-driven DAG triggers. If the triggerer is down, Amazon SQS messages will accumulate in the queue but won’t trigger downstream DAGs until the triggerer recovers. For more information, read the Asset Watchers documentation.
  • Amazon SQS messages have a default retention period of 4 days (configurable up to 14 days). If the consumer environment is unavailable for longer than the retention period, messages will be lost. Consider configuring a dead-letter queue to capture messages that fail processing, and adjust the MessageRetentionPeriod based on recovery requirements.
  • Cross-account Amazon SQS access requires both an IAM identity policy on the producer’s execution role and a resource-based policy on the Amazon SQS queue. If either policy is missing or misconfigured, message delivery will silently fail. For guidance on cross-account access patterns, refer to Four ways to grant cross-account access on AWS.
  • Set the Amazon SQS VisibilityTimeout higher than the expected time for the Asset Watcher to process a message. If the timeout is too short, messages might be redelivered and trigger duplicate DAG runs. Review the Amazon SQS visibility timeout documentation when tuning this value.
  • Each Amazon MWAA environment has limits on the number of DAGs, triggerers, and concurrent DAG runs. If you plan to scale to multiple Asset Watchers monitoring different Amazon SQS queues, check the current Amazon MWAA quotas before making design decisions.
  • Asset URIs must match exactly between the Asset Watcher definition and the consumer DAG’s schedule parameter. A mismatch, even in casing or trailing characters, will prevent the consumer DAG from being triggered. Define assets in a single DAG file to avoid inconsistencies.
  • Pin the provider packages apache-airflow-providers-amazon and apache-airflow-providers-common-messaging to versions compatible with Airflow. Incompatible versions might cause import errors that prevent the triggerer from starting. Use a constraints file as described in this post to avoid dependency conflicts.

Agent skills

AI coding assistants are most useful when they have context about your specific architecture and constraints, not only general programming patterns. Agent Skills, originally developed by Anthropic and released as a public standard in December 2025, provides a portable format for this need. SKILL.md files encode procedural knowledge, best practices, and workflows so that compatible AI coding agents can discover and apply them on demand. The standard is now supported by Kiro, Strands Agents, Anthropic Claude Code, OpenAI Codex, Cursor, Gemini CLI, and other tools. The solution provided here includes an agent skill (agent-skill/) built on this standard that encodes the cross-account orchestration architecture and operational best practices from this post. When you tell the AI coding assistant something like “Write cross-account Amazon MWAA DAGs for my orders pipeline”, the skill guides the agent through the complete workflow:

  • Collecting Amazon SQS queue URL.
  • Generating correctly structured producer and consumer DAG files.
  • Optionally deploying them to Amazon MWAA environments.

The skill doesn’t require you to provide AWS account IDs or Amazon MWAA environment names upfront. Instead, it auto-discovers your environments by running aws mwaa list-environments and aws sts get-caller-identity using the locally configured AWS CLI credentials, then asks you to confirm which environment is the producer and which is the consumer.

The skill works in two modes:

  • Sample mode: Generates the reference producer and consumer DAGs for quick cross-account validation, requiring only the Amazon SQS queue URL as input.
  • Custom mode: Adapts the DAG templates to specific business logic. For example, the producer runs an AWS Glue extract, transform, and load (ETL) job and the consumer triggers a data build tool (dbt) model refresh. This mode customizes DAG IDs, task names, schedules, and processing logic while preserving the correct Asset Watcher patterns.

Beyond code generation, the skill includes an auto-deploy flow. This flow discovers existing Amazon MWAA environments, runs pre-flight checks (Amazon Virtual Private Cloud (Amazon VPC) networking, provider versions, triggerer health, and Amazon SQS queue accessibility), uploads DAGs to the correct Amazon Simple Storage Service (Amazon S3) buckets, and verifies end-to-end readiness. Each step that modifies infrastructure requires explicit user confirmation. Also refer to the GitHub repository for instructions on using it.

Best practices

Airflow Asset Watchers with Amazon SQS are not always the right fit. When they are, they introduce operational considerations that differ from sensor-based polling approaches.

This section covers how to choose the right cross-environment orchestration pattern, how to configure the infrastructure that Asset Watchers depend on (IAM, Amazon VPC, dependencies), and how to design producer and consumer DAGs that are reliable in production.

Cross-account IAM

  • Producer execution role needs sqs:SendMessage and sqs:GetQueueUrl scoped to the specific queue ARN to avoid sqs:*.
  • Amazon SQS queue resource policy must allow the producer role for sqs:SendMessage and consumer role for sqs:ReceiveMessage, sqs:DeleteMessage, sqs:GetQueueAttributes, and sqs:GetQueueUrl.
  • Test cross-account access with the AWS CLI before deploying DAGs. Debugging AWS IAM through Airflow task logs is much harder and slower than catching misconfigurations at the CLI level.
  • Enable Amazon SQS server-side encryption for production queues.

Triggerer health

  • Airflow Asset Watchers run in the triggerer, not the scheduler. Verify triggerer health in the Airflow UI after deploying consumer DAGs.
  • The health API can report healthy even when components are broken. Cross-check by verifying Amazon CloudWatch log streams exist for the Triggerer log group.
  • Monitor airflow-<ENV>-Triggerer CloudWatch logs for ClientError, QueueDoesNotExist, or ImportError.
  • Set Amazon CloudWatch alarms on Amazon SQS ApproximateNumberOfMessagesVisible and the depth of your dead-letter queue (DLQ), which captures messages that fail processing after the maximum number of receive attempts.
  • Pin provider versions with a constraints file to prevent dependency conflicts.

Amazon VPC networking

  • Private subnets must route 0.0.0.0/0 to a NAT Gateway. Without it, workers and triggerers silently fail while the web server appears healthy.
  • Use two NAT Gateways (one per Availability Zone) for production high availability.
  • For private routing mode, use Amazon VPC Endpoints (Amazon S3, Amazon SQS, Amazon CloudWatch Logs, and Amazon Elastic Container Registry (Amazon ECR)) instead of NAT.
  • Confirm Amazon CloudWatch log streams exist for Scheduler, Worker, DAGProcessing, and Triggerer. Empty log groups mean containers aren’t running.
  • Security group must allow self-referencing inbound traffic and unrestricted outbound.

Dependency management

  • Pin provider versions with == and use a constraints file. Unpinned versions break on environment updates.
  • Test dependencies locally with MWAA Docker images before deploying.
  • Check the requirements_install_ip log stream after updates. If networking was unavailable at creation, force reinstall with a new requirements-s3-object-version.
  • Review pre-installed base packages before adding to requirements.txt to avoid version conflicts.

Choosing an orchestration pattern

Not every cross-environment dependency warrants an Asset Watcher. Airflow 3.0 offers three main orchestration patterns: Asset Watchers with Amazon SQS, the MwaaTriggerDagRunOperator, and sensor-based polling, each with different trade-offs in response time, coupling, and resource consumption. Use the following table to match your use case to the right pattern before committing to an implementation.

Pattern How it works Response time Coupling Occupies a worker? Good fit
1 Asset Watchers + SQS (this post) Consumer’s triggerer listens on SQS, triggers DAG on message arrival Seconds Loose No Cross-account pipelines. Fan-out. Independent release cycles
2 MwaaTriggerDagRunOperator Producer calls MWAA API to start a DAG in another environment Seconds Tight Yes (with wait_for_completion) Same-account one-to-one triggers
3 Sensors (polling) Consumer periodically checks for a condition Poll interval Medium Yes (unless deferrable) Persistent-state conditions. Intra-environment dependencies
  • Avoid wiring persistent-state triggers (for example, S3KeyTrigger) into Asset Watchers. They fire continuously because the condition never clears.

DAG authoring

  • Minimize module-level code. DAG files are re-parsed every cycle, and heavy imports slow the entire parsing loop.
  • Design tasks so they produce the same result whether they run once or multiple times (a property called idempotency). Duplicate Amazon SQS messages can occur on retries, so prefer UPSERT (insert or update) over INSERT to avoid duplicate records.
  • Keep secrets out of DAG files and message bodies. Use Airflow Connections (aws_conn_id) instead.
  • Test DAG imports locally with python your_dag.py before uploading to S3.
  • Allow time for DAG parsing after S3 upload, or force with dags reserialize.

Producer DAG design

  • Include dag_id, run_id, logical_date, and dataset-specific context in Amazon SQS messages so consumers can route without calling back.
  • Use SqsHook instead of the raw boto3 package. It respects aws_conn_id and integrates with Airflow logging.
  • Let publish failures raise so the Airflow retry mechanism handles redelivery.

Consumer DAG design

  • Access messages through triggering_asset_events, not by reading the queue directly. The Asset Watcher has already consumed the Amazon SQS messages.
  • Validate message payloads defensively. Producers might evolve their schema over time.
  • Use conditional asset scheduling (& / |) for complex multi-asset dependencies.

Clean up resources

To avoid ongoing AWS charges, delete the resources you created as part of this solution when you are done. The GitHub repository includes step-by-step cleanup instructions for removing the Amazon SQS queue, Amazon MWAA environments, IAM roles and policies, and Amazon S3 buckets.

Refer to the cleanup instructions in the GitHub repository to remove the provisioned resources.

Conclusion

Asset-based scheduling in Apache Airflow 3.0, with Asset Watchers, gives you a practical way to coordinate workflows across Amazon MWAA environments without polling overhead or tight coupling. By using Amazon SQS as a reliable message broker, you can build responsive, decoupled data pipelines that span multiple Amazon MWAA environments and AWS accounts without the operational overhead of traditional polling mechanisms.

This approach reduces cross-environment orchestration latency from minutes to seconds, replaces custom sensors with declarative asset-based scheduling, and gives you and your team the flexibility to maintain independent Amazon MWAA environments while still coordinating complex workflows. Amazon SQS durable message delivery reduces the risk of lost signals, even during temporary environment outages.

To get started:

  1. Review the architecture (5 minutes): Open the architecture diagram in the repository and confirm which Amazon MWAA environments will be the producer and which will be the consumer.
  2. Set up the Amazon SQS queue (15 minutes): Create a cross-account Amazon SQS standard queue and apply the IAM identity and resource-based policies from the Cross-account IAM section. Verify access with the AWS CLI before proceeding.
  3. Deploy and validate the DAG examples (30 minutes): Copy the producer and consumer DAG snippets from the Implementation section into Amazon MWAA environments, trigger the producer DAG manually, and confirm the consumer DAG runs automatically.
  4. Run pre-flight checks (20 minutes): Work through the Amazon VPC networking, provider version, and triggerer health checks in the Best Practices section. Confirm Amazon CloudWatch log streams exist for the Triggerer log group before declaring the environment ready.
  5. Optionally, use the agent skills: If you use an AI coding assistant, install the skill from the repository and describe the business logic in natural language to generate deployment-ready DAGs tailored to your pipeline.

As you scale data operations across multiple accounts and AWS Regions, asset-based scheduling with Asset Watchers provides the foundation for building modern, event-driven data architectures on AWS. Start with basic producer-consumer patterns and gradually evolve to complex multi-asset dependencies as orchestration requirements grow.

For more information, refer to


About the authors

Satya Chikkala

Satya Chikkala

Satya is a Senior Solutions Architect at Amazon Web Services, based in Melbourne, Australia. He helps enterprise customers design scalable cloud solutions that drive growth and efficiency. Outside of work, Satya trades virtual clouds for real ones – climbing rock faces, traversing mountain trails, and capturing it all through his camera lens

Corrine Tan

Corrine Tan

Corrine is a Cloud Architect at AWS specialising in data platform design across financial services, government, and startups. With a consulting background, she builds scalable, domain-oriented architectures using cloud-native technologies. Her expertise includes streaming pipelines, Airflow orchestration, data quality, and full-stack systems integrating data, models, and applications, delivering real-time platforms from ingestion to consumption

Haofei Feng

Haofei Feng

Haofei is a Senior Cloud Architect at AWS with over 20 years of expertise in DevOps, IT Infrastructure, Data Analytics, and AI. He specializes in guiding organizations through cloud transformation and generative AI initiatives, designing scalable and secure GenAI solutions on AWS. Based in Sydney, Australia, when not architecting solutions for clients, he cherishes time with his family and Border Collies.

Amazon Redshift multi-Region disaster recovery

Post Syndicated from Werner Gunter original https://aws.amazon.com/blogs/big-data/amazon-redshift-multi-region-disaster-recovery/

Modern enterprises trust Amazon Redshift to power their most demanding analytics workloads and increasingly require multi-Region disaster recovery to protect those workloads against Regional disruptions. From real-time fraud detection and regulatory reporting to customer-facing dashboards processing millions of transactions daily, organizations are designing for resilience from day one. In financial services, for example, regulatory frameworks increasingly mandate geographic redundancy for data infrastructure, making cross-Region disaster recovery (DR) not only a technical consideration but a compliance requirement. A well-designed DR strategy keeps your analytics infrastructure available and responsive regardless of Regional disruptions, protecting revenue streams, maintaining regulatory standing, and preserving customer trust.

In our previous blog post, Implement disaster recovery with Amazon Redshift, we covered node-level recovery, Availability Zone (AZ) recovery, Multi-AZ deployments, cross-Region backup setup, CNAME implementation, Amazon Redshift Spectrum and Redshift Data sharing considerations.

In this post, we walk through the core concepts of cross-Region disaster recovery, introduce a framework for assessing your requirements, and then dive deep into three primary DR strategies for Amazon Redshift: Active-Passive, Active-Active, and a Hybrid approach. For each strategy, we cover architecture, trade-offs, implementation guidance, and cost considerations so you can make an informed decision for your workload.

What is disaster recovery?

Disaster recovery includes the set of policies, tools, and procedures that enable an organization to restore critical systems and data after an incident. It helps maintain business continuity during events such as a regional AWS outage, accidental data deletion, infrastructure failure, or a security event.

Any DR strategy depends on two key metrics:

  • Recovery Point Objective (RPO): The maximum acceptable amount of data loss, measured in time. An RPO of 30 minutes means you can tolerate losing up to 30 minutes of data that you can reproduce from your source systems.
  • Recovery Time Objective (RTO): The maximum tolerance for downtime, before restoring business operations after a disaster is declared. An RTO of 30 minutes means your systems must be fully operational within 30 minutes of a failure.

These two numbers drive all architectural decisions for DR and understanding them helps clarify the trade-offs between various DR strategies.

Assessing your DR requirements

Before selecting a strategy, you need to assess your workload’s criticality and your organization’s tolerance for data loss and downtime. Ask yourself:

  • What is the business impact of downtime? If your Amazon Redshift cluster powers customer-facing applications, regulatory reporting, or real-time risk calculations, even an hour of downtime might be unacceptable. If it powers internal dashboards refreshed daily, a 2-hour RTO might be acceptable.
  • Can data be backfilled from upstream sources? If your data pipeline originates from Amazon Managed Streaming for Apache Kafka (Amazon MSK) or Amazon Simple Storage Service (Amazon S3), you might be able to replay events after a failover, relaxing your RPO requirements. If data is generated in-place or cannot be replayed, you need tighter replication.
  • What are your regulatory obligations? Financial services, healthcare, and government workloads often have explicit RPO/RTO requirements mandated by regulators. These are non-negotiable floors.
  • What is your cost tolerance? Active-active architectures can double your infrastructure spend. Active-passive approaches offer significant savings at the cost of slightly longer recovery times.

The following table serves as a quick reference to match your requirements to a DR strategy:

Requirement Recommended strategy
RPO: 10–30 min, RTO: 1–2 hours, cost-sensitive Active-Passive
RPO: Near-zero, RTO: Minutes, mission-critical Active-Active
Mixed criticality across data tiers Hybrid

The following decision tree helps you select the right disaster recovery strategy based on your workload’s RPO and RTO requirements.

Decision tree for choosing a Redshift DR strategy based on RPO and RTO requirements

Cross-Region best practices

Regardless of which strategy you choose, the following practices apply universally to Amazon Redshift DR implementations.

Use multi-Region AWS KMS keys: Encrypt your Amazon Redshift clusters and S3 data with multi-Region AWS Key Management Service (AWS KMS) keys. This avoids the need to re-encrypt data during failover, which can add significant time to your RTO. Note that AWS KMS allows only one replica of a multi-Region key per AWS Region within the same partition. This is a service-level constraint. In most DR scenarios, a single multi-Region key per Region is sufficient since all resources in that Region can share the same key.

Automate with infrastructure as code: Define all DR Region infrastructure with infrastructure as code (IaC), such as Terraform, AWS CloudFormation, or AWS Cloud Development Kit (AWS CDK). IaC supports consistency between Regions, removes manual configuration errors, and enables rapid provisioning during failover. For organizations using Terraform Enterprise, verify that your workspace configuration supports multi-Region deployments.

Implement comprehensive monitoring. Use Amazon CloudWatch alarms where possible:

Early detection of replication failures is critical. A silent replication failure discovered during a disaster is far worse than one caught proactively. For detailed metrics monitoring configuration, see the Amazon CloudWatch alarms user guide.

Test quarterly. DR plans that aren’t tested regularly are more likely to fail during an actual disaster. Conduct quarterly failover tests that measure actual RTO and RPO against your targets. Validate data consistency post-failover. Document lessons learned and update your runbooks accordingly.

Use Amazon Redshift Spectrum. For cold and warm data tiers, you can query data directly in Amazon S3 without loading it into Amazon Redshift. This can reduce your data restoration requirements during failover. Remember that your cluster and S3 bucket must be in the same Region. Recreate external schemas in the DR Region pointing to your replicated S3 data. For Amazon Redshift Serverless endpoints and Redshift provisioned clusters without Spectrum, the DR strategy relies on snapshot replication and cross-Region restore. The same principles apply regardless of whether you use RA3 or RG (Graviton) node types.

Strategy 1: Active-Passive with snapshot replication

In an active-passive configuration, your primary AWS Region runs the end-to-end workload, including data ingestion, processing, and serving data through Amazon Redshift. Amazon Redshift replicates data to the DR Region using its built-in cross-Region snapshot feature. During a disaster, you restore clusters from replicated snapshots in the DR Region.

RPO: 15 minutes plus time for data replication | RTO: 1–2 hours | Cost: Low

Active-Passive architecture with Amazon Redshift cross-Region snapshot replication to the DR Region

Snapshots in Amazon Redshift provisioned clusters

By default, Amazon Redshift provisioned clusters take a new snapshot every 8 hours, or whenever 5 GB of data changes are detected on any single node, whichever comes first. The 5 GB threshold is evaluated per node independently.

Amazon Redshift offers automated snapshots of your cluster at no extra storage cost in both your primary and DR Regions. You will incur charges for the data transfer when Amazon Redshift copies snapshots across Regions. The initial cross-Region copy is a full snapshot transfer. Subsequent copies are incremental, transferring only the changed blocks since the last snapshot, which significantly reduces transfer time and cost.

When to customize the automatic snapshot schedule

You can override the default and set a custom schedule, with a minimum frequency of once per hour. However, this is only useful in one scenario:

Cluster type Recommendation
≥ 5 GB of changes per node per hour Keep the default — already snapshotting frequently enough
< 5 GB of changes per node per hour Customize the schedule to take snapshots more often

When to use manual snapshots

If you need a guaranteed RPO of less than 1 hour (for example, every 15 minutes), or need to retain backups beyond 35 days, use manual snapshots scheduled at the frequency you want. Manual snapshots incur additional storage charges but are retained until explicitly deleted.

Comparing automatic and manual snapshots

 

Automatic snapshots Manual snapshots
Frequency Every 8 hours or 5 GB change (customizable to run hourly) Any frequency you choose
Best for RPO ≥ 1 hour RPO < 1 hour (for example, 15 min)
Cost No additional cost (included with cluster) Additional storage charges.
Retention 1–35 days (configurable) Until explicitly deleted
Cross-Region copy Supported (incremental) Supported (incremental)

Architecture

The following diagram illustrates the Active-Passive DR architecture.

The Active-Passive strategy keeps compute resources in the DR Region ready to be spun up from snapshots when needed. When replicating data, consider the other services that are part of your end-to-end data pipeline. In the Amazon Redshift data sharing model, the producer cluster creates and owns the data, while consumer clusters read from the producer through data shares. In a DR context, the producer is restored first in the DR Region, then consumer clusters are resumed to serve read workloads.

  • Amazon S3 is frequently used with Amazon Redshift. For complete data resiliency, replicate data in Amazon S3 as well using Amazon S3 Cross-Region Replication (S3 CRR). It continuously replicates your S3 data lake to the DR Region with near-zero lag. For Apache Iceberg tables, we recommend using replication for Amazon S3 Tables, a capability of Amazon S3, to guarantee that both the data and the associated metadata (manifests, snapshots) are replicated consistently to the DR Region.
  • Customers use AWS Glue Data Catalog and AWS Lake Formation to catalog and maintain permissions. Read this post on how to build multi region resilient data architecture using AWS Glue and AWS Lake Formation.
  • Customers often use Amazon DynamoDB alongside Amazon Redshift in data pipeline architectures to track pipeline orchestration state, such as job IDs, processing timestamps, batch completion flags, and ingestion checkpoints that tell your pipeline which data has been processed. Amazon DynamoDB Global Tables replicate this state across both Regions, so pipeline state is available in the DR Region and you know exactly where to resume processing after failover.

DR Region (Passive) components:

  • Amazon Redshift clusters ready to restore from snapshots.
  • AWS Lambda functions with data transformation pipelines code deployed and ready.
  • Amazon MSK infrastructure defined in IaC but not provisioned.
  • Amazon EMR job definitions ready but not running.

Failover sequence (20–60 minutes):

  1. Restore the Amazon Redshift cluster in DR Region, from the latest cross-Region snapshot (this is typically the longest step).
  2. Provision and start Amazon MSK clusters in the DR Region.
  3. Disable S3 event triggers for AWS Glue Catalog (to prevent split-brain metadata updates).
  4. Stand up Amazon EMR and resume data processing.
  5. Resume paused Amazon Redshift consumer clusters.
  6. Recreate external schemas pointing to the DR Region’s AWS Glue Catalog. Note: External schemas, external schema-level permissions, and references to external resources (for example, S3 paths, AWS Glue Catalog databases) included in the Amazon Redshift snapshot, contain references to primary Region resources. Plan to recreate these in your DR Region as part of your failover runbook. Database users, groups, and their internal permissions are replicated with the snapshot. Plan to script external schema recreation as part of your failover runbook.
  7. Update query or application service endpoints to the DR Region.
  8. Update Lambda data transformation pipelines to point to the new producer endpoint.

When to choose Active-Passive

  • You can tolerate 15–20 minutes of data loss.
  • A 1–2 hour RTO is acceptable for your business.
  • Cost optimization is a priority.
  • Data can be backfilled or replayed from upstream sources (for example, Amazon MSK topic retention).

Strategy 2: Active-Active multi-Region

In an Active-Active configuration, both your primary and DR Regions run fully operational data pipelines simultaneously. Data is ingested, processed, and served in both Regions at all times. Failover becomes a matter of redirecting traffic rather than restoring infrastructure. This reduces RTO to minutes.

RPO: Near-zero | RTO: < 1 hour (often minutes) | Cost: High

Architecture

The following diagram illustrates the Active-Active DR architecture. Active-Active requires mirroring your entire pipeline, from ingestion through serving, across both Regions.

Real-time replication layer:

  • Amazon MSK Replicator: Mirrors Kafka topics in real time from the primary Region to the secondary Region. This is the earliest point of replication in the pipeline, so the DR Region processes the same events with minimal lag.
  • Amazon DynamoDB Global Tables: Active state tracking across both Regions keeps pipeline controls and job state synchronized.
  • Active Amazon EMR processing: Both Regions continuously process incoming data, maintaining fresh state in their respective S3 data lakes and AWS Glue Catalogs.
  • Active Amazon Redshift producer clusters: Both Regions continuously ingest processed data, maintaining near-identical warehouse state.
  • Mirrored data transformation pipelines: Data transformation events are actively processed in the DR Region through DynamoDB replication, keeping derived data consistent. In the Active-Active model, both Regions maintain their own Amazon Redshift cluster that independently ingests the same source data, so the DR Region’s Amazon Redshift already has current data. The mirrored pipeline supports the transformation logic and derived datasets stay synchronized.

DR Region (Active) components:

  • Amazon Redshift clusters paused but ready (can be activated in minutes).
  • Any Amazon Redshift data shares synchronized regularly between Regions.
  • External schemas active and synchronized.
  • Query or application service endpoints pre-configured and tested.

Failover sequence (minutes):

  1. Failover Amazon MSK consumers to the DR Region’s Amazon MSK cluster.
  2. Resume Amazon Redshift consumer clusters in the DR Region.
  3. Update query or application service endpoints to point to the DR Region.
  4. Promote the DR Region’s Lambda data transformation pipelines functions to act as primary.

Because the DR Region’s pipeline is already running, there is no infrastructure provisioning delay. Failover is primarily a configuration change.

Cost considerations

Active-Active essentially doubles your infrastructure costs. You are running full Amazon MSK, Amazon EMR, and Amazon Redshift clusters in both Regions simultaneously. For large-scale deployments (1+ PB), this represents a significant ongoing investment. The business case rests on the cost of downtime exceeding the cost of duplicate infrastructure. This is a calculation that often favors Active-Active for customer-facing or regulatory workloads.

When to choose Active-Active

  • You require near-zero RPO with no tolerance for data loss.
  • RTO must be measured in minutes, not hours.
  • Your analytics infrastructure directly impacts customer-facing operations or regulatory compliance.
  • The cost of downtime (financial, reputational, regulatory) exceeds the cost of duplicate infrastructure.
  • You have strict Service Level Agreements (SLAs). For example, zero RPO and full-service functionality within 4 hours including data ingestion.

Strategy 3: Hybrid — tiered DR by data criticality

Not all data in your warehouse is equally critical. Some real-time insights and regulatory reports demand near-zero RPO, while historical trend analyses and archived compliance data can tolerate hours of recovery time. A Hybrid approach applies different DR strategies to different data tiers, optimizing cost while protecting what matters most.

RPO: Varies by tier | RTO: 30 minutes – 2 hours | Cost: Medium

Architecture

The following diagram illustrates the Hybrid DR architecture.

The Hybrid strategy requires a data model that supports clear separation at the schema or table level, with different recovery objectives applied per tier.

Tier 1: Hot data (Active-Active):

  • Real-time dashboards, regulatory reporting, customer-facing analytics.
  • Near-zero RPO through Amazon MSK Replicator and active Amazon Redshift producer in both Regions.
  • RTO: Minutes.

Tier 2: Warm data (Active-Passive):

  • Daily reports, historical trend analysis, internal operational data.
  • RPO: 1 hour through hourly Amazon Redshift snapshots replicated cross-Region.
  • RTO: 1–2 hours.

Tier 3: Cold data (S3 replication only):

  • Archived data, long-term compliance storage, infrequently accessed history.
  • RPO: Hours (S3 CRR with standard replication lag).
  • RTO: 2+ hours (restore from S3 into Amazon Redshift Spectrum or a new cluster).
  • No active Amazon Redshift infrastructure in DR Region for this tier.

Implementation considerations

  • Your data model must support clear separation at the schema or table level to apply different recovery strategies. To achieve different RPO/RTO per data tier, while avoiding unnecessary table level maintenance complexities, consider using separate clusters or namespaces for each tier, or use a combination of cluster snapshots and S3-based backups (UNLOAD) for finer-grained table-level recovery.
  • Workload Management (WLM) queues or separate clusters may be needed to isolate hot, warm, and cold workloads.
  • Monitoring must track replication latency independently for each tier.
  • Failover runbooks must be tier-aware. Operators need to know which systems to restore first.

When to choose Hybrid

  • You have clearly defined data tiers with meaningfully different criticality.
  • Your data model already supports or can be refactored to support hot/cold separation.
  • You want to protect mission-critical data with Active-Active while managing costs for less critical workloads.
  • Your organization has the operational maturity to manage tiered failover procedures.

Testing your DR strategy

Schedule quarterly DR tests that include:

  1. Failover execution following your documented runbook.
  2. RTO measurement from disaster declaration to full operational status.
  3. RPO validation to verify data consistency.
  4. Application testing to confirm connectivity.
  5. Failback procedure documentation.
  6. Lessons learned and runbook updates.

Conclusion

Disaster recovery for Amazon Redshift is not a one-size-fits-all problem. The right strategy depends on your RPO and RTO requirements, your data’s criticality, your ability to replay data from upstream sources, and your cost tolerance.

  • Active-Passive offers a cost-effective path to 10–20 minute RPO and 1–2 hour RTO, suitable for most analytics workloads.
  • Active-Active delivers near-zero RPO and minute-scale RTO for mission-critical services where downtime cost exceeds infrastructure cost.
  • Hybrid lets you apply the right level of protection to the right data, optimizing cost without compromising on what matters most.

Whichever strategy you choose, the fundamentals remain the same: replicate early in the pipeline, automate your infrastructure, monitor replication health continuously, and test your failover procedures regularly. DR is not a project you complete. You maintain it as an ongoing practice.

Next steps


About the authors

Werner Gunter

Werner Gunter

Werner is a Principal Specialist Solutions Architect at Amazon Web Services, based in Berlin, Germany. As a seasoned data professional, he has helped large enterprises worldwide over the past 2 decades, to modernize their data analytics estates.

Nita Shah

Nita Shah

Nita is a Sr. Analytics Specialist Solutions Architect at AWS based out of New York. She has been building enterprise data platforms, data warehousing, and analytics solutions for over 20 years and specializes in Amazon Redshift. She is focused on helping customers design and build enterprise-scale well-architected analytics and decision support platforms

Upgrade Amazon Redshift DC2 clusters to the new Amazon Redshift RG

Post Syndicated from Ricardo Serafim original https://aws.amazon.com/blogs/big-data/upgrade-amazon-redshift-dc2-clusters-to-the-new-amazon-redshift-rg/

When you upgrade your Amazon Redshift DC2 (Dense Compute) clusters to RG instances powered by AWS Graviton, you gain access to capabilities that were never available on DC2. These include managed storage, data sharing, zero-ETL integrations, streaming ingestion, and faster query compilation. You also gain availability zone (AZ) features such as cross-AZ cluster relocation for disaster recovery (DR) and concurrency scaling for writes. RG also adds a built-in data lake engine for querying Apache Iceberg and Parquet tables directly on your cluster nodes.

This post covers the new features you gain when upgrading from DC2 to RG, the node mapping guidance for sizing your new cluster, the upgrade methods available, and validation options including Amazon Redshift Test Drive.

Why upgrade from DC2 to RG instances

As data volumes grow, DC2 customers face a choice: add extra compute nodes only to get more storage, or offload data elsewhere. The local SSD capacity on each node is fixed, and there is no managed storage tier to absorb growth. Both RA3 and RG instances solve this with Amazon Redshift Managed Storage, which decouples storage from compute. You can scale data volume independently of node count, paying only for the storage you use with no fixed ceiling per node. This means you no longer need to over-provision compute to accommodate data growth.

RG is the recommended upgrade path over RA3. RG instances run on AWS Graviton processors, delivering higher throughput for data warehouse and data lake workloads at a lower price per vCPU compared to RA3. Because both RA3 and RG share the same managed storage architecture and feature set, RG provides more performance for less cost. For current pricing details, visit Amazon Redshift pricing.

Amazon Redshift RG instances run on AWS Graviton processors. These processors provide more compute cores and lower memory latency compared to the previous-generation hardware behind DC2. This can translate to faster query execution for data warehouse workloads, particularly for large scans where memory throughput is the bottleneck. Exact performance improvements depend on workload characteristics, cluster size, and query complexity. Use Redshift Test Drive to measure the difference for your specific workload.

Data lake access: New with RG

DC2 clusters can query data in Amazon Simple Storage Service (Amazon S3) through Amazon Redshift Spectrum. However, Spectrum adds a per-TB scanning cost on top of your cluster pricing, and does not support enhanced VPC routing on DC2 provisioned clusters (requiring additional configuration for secure S3 access).

RG addresses these constraints with an integrated data lake engine that processes queries directly on your cluster’s dedicated compute nodes:

DC2 (Spectrum) RG (Integrated Engine)
Data lake query cost Extra $5/TB scanned on top of cluster cost Included in node pricing, no extra charge
Apache Iceberg Queries via Spectrum Native queries on cluster compute, no Spectrum needed
Apache Iceberg Statistics Manual collection JIT-Analyze auto-collects statistics
VPC routing Not compatible with enhanced VPC routing No conflict, runs on the cluster itself

With RG, you can consolidate warehouse and data lake workloads on a single cluster with no extra per-query charges for data lake access.

Features available with RG

Upgrading from DC2 to RG gives you access to the full set of modern Amazon Redshift capabilities. Three of the most impactful for DC2 customers are data sharing, zero-ETL integrations, and managed storage. With data sharing, you can query live data from other Amazon Redshift clusters or accounts without copying or moving data, reducing storage duplication and keeping consumers always up to date. Zero-ETL integrations automatically replicate data from Amazon Aurora, Amazon Relational Database Service (Amazon RDS), and Amazon DynamoDB into Amazon Redshift without building or maintaining ETL pipelines. This reduces operational overhead and data freshness lag. Managed storage scales independently from compute, so you can grow your data without adding nodes and only pay for the storage you use.

Additional capabilities available with RG:

  • Streaming ingestion – ingest data from Amazon Kinesis Data Streams and Amazon Managed Streaming for Apache Kafka (Amazon MSK) in near real-time, so you can build dashboards and analyze the latest data without batch delays.
  • Concurrency scaling for writes – automatically add transient capacity during burst write workloads, so ingest operations don’t slow down your analytical queries.
  • Cross-AZ cluster relocation – relocate your cluster to another Availability Zone with no endpoint changes, supporting disaster recovery without the cost of a standby cluster.
  • Multi-AZ deployments – run your cluster across multiple Availability Zones as a single database delivering high availability (HA) and automatic failover without a passive standby.
  • Faster query compilation – queries compile faster on Graviton processors, reducing cold-start latency for new or modified queries.

RG instance details and node mapping

This table shows the available RG instance configurations:

RG Instance vCPUs Memory
rg.large 2 16 GiB
rg.xlarge 4 32 GiB
rg.4xlarge 16 128 GiB
rg.12xlarge 48 384 GiB

For current pricing, visit Amazon Redshift pricing for more information.

DC2 to RG node mapping guidance

Use this table to determine the recommended starting configuration when upgrading from DC2:

Current Node Type Node Ratio RG Node Type Guidance
dc2.large (1–3 nodes) 1:1 rg.large 1 rg.large for every 1 dc2.large
dc2.large (4 nodes) 4:3 rg.large 3 rg.large for 4 dc2.large
dc2.large (5–15 nodes) 8:3 rg.xlarge 3 rg.xlarge for every 8 dc2.large
dc2.large (16–32 nodes) 10:1 rg.4xlarge 1 rg.4xlarge for every 10 dc2.large
dc2.8xlarge (2–15 nodes) 2:3 rg.4xlarge 3 rg.4xlarge for every 2 dc2.8xlarge
dc2.8xlarge (16–128 nodes) 2:1 rg.12xlarge 1 rg.12xlarge for every 2 dc2.8xlarge

Extra nodes might be needed depending on workload requirements. Add or remove nodes based on the compute requirements of your required query performance. Validate your specific configuration using Redshift Test Drive before migrating production workloads.

Prerequisites

Before starting the upgrade, confirm the following:

  • Snapshot availability — a recent snapshot of your DC2 cluster is required for all upgrade methods. If automated snapshots are disabled, create a manual snapshot before starting. Visit Amazon Redshift snapshots for more information.
  • Network configuration — verify that your virtual private cloud (VPC), subnet groups, and security groups are configured to support the new RG cluster. If you use enhanced VPC routing, confirm your S3 endpoint and route table configuration. Visit Enhanced VPC routing for more information.
  • Cluster version — your DC2 cluster must be running a supported Amazon Redshift version. Check the release notes for minimum version requirements.

Upgrade methods

Three methods are available for migrating from DC2 to RG instances. The right choice depends on your operational constraints: whether you need write access during migration, whether the target configuration supports elastic resize, and how much downtime your workload can tolerate.

Elastic resize is the fastest and most efficient path. Amazon Redshift creates a snapshot, provisions the RG cluster, and redirects the endpoint automatically. The cluster remains in read-only mode for a few minutes during the operation, and the endpoint doesn’t change, meaning no application-side updates are required. This is the recommended method when the target configuration is supported by elastic resize.

Classic resize

Use classic resize when the target configuration is not available through elastic resize, or when you need data slice rebalancing. Downtime is similar to elastic resize (a few minutes of read-only mode in Stage 1). In Stage 2, data redistributes to its original distribution patterns in the background without blocking queries. The advantage of classic resize is that it rebalances data slices evenly across nodes. This matters when you move to a different node type that might require a different number of slices. Stage 2 can take time on busy clusters, and the duration depends on data volume, cluster utilization, and target cluster size. Queries might run slower until redistribution completes.

Snapshot and restore with cluster identifier swap

This method uses snapshot and restore of the existing DC2 cluster to provision a new RG cluster with a different identifier. After validating the new cluster, you swap the cluster identifiers to redirect application traffic without changing the endpoint. This approach provides these benefits:

  • Test and validate the RG cluster while the DC2 cluster continues serving production traffic.
  • Roll back by reversing the identifier swap if issues arise.
  • No application-side endpoint changes required after the swap.

The trade-off is that data written to the source cluster after the snapshot requires manual synchronization before the cutover. If your migration plan includes a write-freeze window, you can take the final snapshot at the start of that window and avoid synchronization entirely.

This AWS Command Line Interface (AWS CLI) command illustrates restoring a DC2 snapshot to an RG cluster:

aws redshift restore-from-cluster-snapshot \
    --cluster-identifier my-cluster-rg \
    --snapshot-identifier my-dc2-snapshot \
    --node-type rg.4xlarge \
    --number-of-nodes 3 \
    --cluster-subnet-group-name my-subnet-group \
    --vpc-security-group-ids sg-abc123 \
    --cluster-parameter-group-name my-param-group \
    --port 5439 \
    --no-publicly-accessible \
    --enhanced-vpc-routing \
    --iam-roles 'arn:aws:iam::111122223333:role/RedshiftRole'

After restoring, validate your workload on the new cluster. When ready, swap the cluster identifiers:

aws redshift modify-cluster \
    --cluster-identifier my-cluster \
    --new-cluster-identifier my-cluster-dc2-old

aws redshift modify-cluster \
    --cluster-identifier my-cluster-rg \
    --new-cluster-identifier my-cluster

Validating your target configuration

Before migrating production clusters, validate that your target RG configuration meets performance requirements. There are several ways to approach this depending on your needs:

Run your existing QA process on a test cluster. Create an RG cluster from a snapshot, then execute the same test suites and validation scripts you would use for any code or infrastructure change. This approach helps confirm basic compatibility and catch regressions.

Use lower environments first. Migrate your development or staging clusters to RG before production. This gives your team hands-on experience with the new instance type and surfaces any configuration differences in a low-risk setting.

Replay production workloads with Redshift Test Drive. For production-level validation with real traffic patterns, Redshift Test Drive is an open source utility that automates workload replay across multiple target configurations. It extracts queries from your source cluster’s audit logs and replays them against the target, then provides a comparison UI for latency, errors, and deviation.

For a detailed walkthrough, read Find the best Amazon Redshift configuration for your workload using Redshift Test Drive.

Best practices

Before migrating, run Amazon Redshift Advisor on your current cluster to identify optimization opportunities such as unused tables, missing sort keys, or distribution style changes. Drop unnecessary tables to reduce data transfer time, and schedule the migration during off-peak hours for minimal business impact. Removing tables that are no longer used (for example, tables with suffixes like _bkp, _tmp, or _old) also speeds up classic resize. These unused tables would otherwise be rebalanced across nodes during Stage 2, adding time to a process that delivers no value for data no one queries.

During migration, communicate the cutover window to stakeholders. Because the DC2 cluster remains active until the identifier swap, coordinate a brief write-freeze period before the final snapshot to minimize data synchronization effort.

After migration, monitor the cluster for 48–72 hours to identify any performance deviations and adjust node count if needed. Update your runbooks and operational documentation with the new cluster details, node types, and any endpoint changes if you used the snapshot and restore method. Once the migration is considered successful you may delete the DC2 cluster.

Conclusion

Upgrading from Amazon Redshift DC2 to RG instances powered by AWS Graviton gives you a Graviton-based architecture with managed storage and improved query performance. It also gives you access to the full suite of Amazon Redshift features that were never available on DC2: data lake queries, data sharing, zero-ETL, faster query compilation, and cross-AZ relocation. The snapshot restore and cluster identifier swap method provides a safe migration path with built-in rollback. Use Redshift Test Drive to validate your target configuration with real workload data before committing.

To get started, review the RG instance availability and pricing, determine your target configuration using the node mapping guidance, and run Redshift Test Drive against your production workload.


About the authors

Ricardo Serafim

Ricardo Serafim

Ricardo is a Senior Analytics Specialist Solutions Architect at AWS. He has been helping companies with Data Warehouse solutions since 2007.

Nita Shah

Nita Shah

Nita is a Sr. Analytics Specialist Solutions Architect at AWS based out of New York. She has been building enterprise data platforms, data warehousing, and analytics solutions for over 20 years and specializes in Amazon Redshift. She is focused on helping customers design and build enterprise-scale well-architected analytics and decision support platforms.

Ankit Sahu

Ankit Sahu

Ankit brings over 18 years of expertise in building innovative data products and services. His diverse experience spans product strategy, go-to-market execution, and digital transformation initiatives. Currently, as Sr. Product Manager at Amazon Web Services (AWS), Ankit is driving the vision and strategy for Amazon Redshift.

Balancing speed and safety: A control framework for AI coding agents

Post Syndicated from Daniel Begimher original https://aws.amazon.com/blogs/security/balancing-speed-and-safety-a-control-framework-for-ai-coding-agents/

AI coding agents are part of the developer toolchain. Tools like Kiro and Claude Code generate features, tests, and code refactors from natural-language prompts. A single agent can open dozens of pull requests (PRs) across your repositories in an afternoon. That productivity comes with a trade-off: agents optimize for task completion at machine speed with no understanding of your organization’s risk.

Through protocols like the Model Context Protocol (MCP), agents also reach beyond the integrated development environment (IDE) to call APIs, query databases, and modify infrastructure and even entire environments, expanding the scope of resources your application security team defends.

This post lays out an application security (AppSec) control framework for AI coding agents. Two pillars organize the framework: author-time controls shape what the agent produces in the IDE; build-time controls verify and gate what reaches production. Your existing secure software development lifecycle (SDLC) controls still apply and are critical to a defense-in-depth security strategy. The framework shows where to layer additional guardrails so AppSec scales with agent-driven development. The framework is tool-agnostic and cloud-agnostic. Throughout, we use AWS services—Kiro in the IDE and AWS CodePipeline in the build—as a running example that you can adapt to your own toolchain.

Risks

Each of the following risks includes a treatment summary. The control framework section later in this post provides implementation details. The risks are ordered by severity with the highest impact risks first.

R001. Prompt and context injection

Agents read untrusted content, such as issue descriptions, web pages, MCP responses, and README files in third-party packages. Text from outside parties can redirect the agent to disclose secrets, open unauthorized PRs, or invoke tools without user consent. This risk, known as prompt injection, is the top risk in the OWASP Top 10 for LLM Applications. Any agent that reads content from outside parties is exposed, with or without MCP, so connecting tools widens the scope of impact.

Treatment: Treat non-developer input as untrusted. A large language model (LLM) can’t reliably separate instructions from data in a single context window, so architect for it: keep the agent that orchestrates trusted actions separate from the one exposed to untrusted content and grant the exposed agent only read-only, least-privilege access. Require human approval for irreversible actions. Use version-control steering files to prevent silent tampering.

R002. Inadvertent data disclosure and overly permissive configurations

Agents optimize for getting work done. Left unchecked, the code they generate can default to wildcard identity and access management policies, open security groups, and unencrypted storage, or embed sensitive values in code rather than referencing a secrets manager. Most coding agents now include safety mechanisms that make these outcomes less likely, but they remain imperfect, so you still need controls to account for the possibility.

Treatment: Security requirements in a steering document, plus policy-as-code scanning (Checkov, cfn-nag) in the IDE and pipeline. See Context as a security control.

R003. Uncontrolled changes reaching production

Ungated code reaching production isn’t new, but AI agents amplify it. Machine-speed generation can propagate a flawed pattern across repositories before it’s identified.

Treatment: Branch protection rules requiring PR approval (a human-in-the-loop checkpoint), pre-commit hooks for security checks, and sandboxed agent runs that prevent direct pushes to protected branches. The right balance between human review and automated speed depends on the risk profile of the change. For many low-risk paths, automated checks alone might suffice, while higher-risk changes warrant a human checkpoint.

R004. Supply chain risks

Agents don’t always distinguish current best practices from outdated patterns. They might recommend deprecated packages, reference library versions with new Common Vulnerabilities and Exposures (CVEs), and hallucinate package names that don’t exist, which can introduce risks of dependency confusion issues.

Treatment: Software Composition Analysis (SCA) in the pipeline (for example, Amazon Inspector code scanning or Dependabot) to flag vulnerable or unexpected dependencies. For additional control, resolve against a scoped registry like AWS CodeArtifact. Even without a fully curated registry, lockfile validation and allow-listing critical packages reduce exposure.

R005. Uncontrolled external access

Through MCP and tool integrations, agents query databases, call APIs, and modify infrastructure. Without constraints on which tools and data an agent can reach, a single misconfigured integration provides unintended access to sensitive resources.

Treatment: Scope MCP servers to least-privilege tools and resources, enforce authn or authz on external connections, and audit tool invocations. The control point is the configuration file. Review it the same way you review AWS Identity and Access Management (IAM) policies.

R006. Hallucinations and incorrect code

Agents produce plausible-looking output. Code that compiles, passes linting, and looks reasonable can still be functionally wrong: misusing APIs, introducing subtle logic errors, or implementing security-sensitive operations incorrectly. Code that passes continuous integration (CI) but is wrong slips through review; code that fails to build is caught immediately.

Treatment: Layer deterministic verification (static application security testing (SAST), unit tests) with non-deterministic review (LLM-assisted screening against the specification). Neither catches everything alone.

R007. Scope creep

Given a bug-fix prompt, an agent might also refactor surrounding code, disable an unreliable test, or reorganize imports. Unrequested changes introduce regressions and complicate review.

Treatment: A reviewed specification document that defines what must change and what must not, paired with a targeted review of the proposed changes. See Specifications as scope boundaries.

The preceding risks share a common thread: agents produce output faster than humans can review it, and they lack context to self-correct.

The following framework addresses this gap. It organizes controls into two pillars: author-time (pre-generation and post-generation of code) and build-time (in the pipeline, before code reaches production). Author-time controls shape what the agent produces. Build-time controls verify it. Neither is sufficient alone; together they reduce the volume and severity of issues that reach human reviewers.

Deterministic compared to non-deterministic mitigations

Deterministic mitigations [D] produce the same result every time. Linters, SAST scanners, secrets detection, and policy-as-code match patterns against rules and define security invariants: no critical findings, no hardcoded secrets, and no wildcard IAM policies. Use them when the condition can be expressed as a rule. Organizations already have these and must continue enforcing them.

Non-deterministic mitigations [ND] use model judgment. They include steering documents, LLM-as-judge review, specification compliance checks, and scope-creep detection, and they evaluate intent rather than patterns. They catch novel issues that rules miss, but are probabilistic. Use them when evaluation requires context or reasoning across files. This is the new layer that AI-generated code demands, because agents produce code that can pass every deterministic check yet remain functionally wrong.

Human review [H] provides the final layer for the risk-based decisions neither tool type can make. Apply it where judgment is needed, not everywhere: routing every change to a person invites consent fatigue, where reviewers approve by reflex and the control loses its value. The default reflex is to route everything back to a human, but that isn’t always the right response—reserve human judgment for the decisions that genuinely need it.

The control framework

The framework organizes controls into two pillars. Author-time controls (Pillar 1) shape what the agent produces in the IDE, before code is generated and just after. Build-time controls (Pillar 2) verify and gate that output in the pipeline, before it reaches production. The controls within each pillar are tagged deterministic [D], non-deterministic [ND], or human [H].

Pillar 1: Author-time controls (pre- and post-generation of code)

Author-time controls work inside the IDE, where the developer and agent still hold full context. They shape the prompt and the generated output before it ever reaches a pull request. The following controls apply at this stage.

Context as a security control [ND]

Control statement: Encode security invariants as natural-language constraints in a steering document that every developer environment consumes at session start. Addresses R002.
Many AI coding agent risks share one root cause: the agent lacks the security context an experienced developer carries implicitly. Your security team sets the policies, such as Amazon Simple Storage Service (Amazon S3) buckets require encryption, API gateways require mutual TLS, and credentials must come from AWS Secrets Manager. Developers don’t always have these requirements available when they’re building. They build what works, not what’s compliant. An AI agent amplifies this gap because it defaults to whatever pattern dominated its training data, with no awareness of your organization’s security posture.

A key mitigation is steering. Security teams write these invariants once as natural-language guidance in a steering document, then distribute them as shareable resources that developers consume in their IDE. The agent loads the file at session start and treats the contents as standing requirements:

  • IAM policies must follow least-privilege principles; no wildcard Amazon Resource Names (ARNs).
  • No hardcoded credentials in source code; use a secrets manager.
  • Security groups must not allow unrestricted inbound access.

This shifts security left, before code generation begins. Steering biases generation toward secure defaults; it doesn’t guarantee them. Treat it as a strong default, paired with the following deterministic gates that block non-compliant code from merging. Security teams define the rules once and every developer environment inherits them automatically. Steering reduces the volume of issues that reach the pipeline, though it doesn’t replace downstream scanning.

How to write effective steering rules: Keep each rule specific and testable, scope it to a concrete risk class, keep the rule set concise so the agent can hold it in context, and iterate from the issues your scanners and reviewers surface.

Specifications as scope boundaries [ND]

Control statement: Require a reviewed specification before code generation begins. Define what must change and what must not. Addresses R007.

Spec-driven workflows turn vague prompts into reviewable specifications before code is generated. This creates a human checkpoint at the design phase, where security decisions are made:

  • Requirements use testable notation that’s auditable before the agent writes a line of code. For example, the Easy Approach to Requirements Syntax (EARS): WHEN [condition] THE SYSTEM SHALL [behavior].
  • Tasks are ordered in implementation steps, each mapped back to a requirement.

For bug fixes, specifications add a critical element: unchanged behavior documentation. This is an explicit list of behaviors that must continue working, giving the agent a written boundary against scope creep.

In this model, the specification becomes the primary artifact, code is a derivative of it. Human review effort concentrates on whether the specification solves the right problem with the right constraints, not on reading implementation diffs line by line.

Controlled tool access using MCP [D + ND]

Control statement: Scope each MCP server to the minimum set of tools the agent needs, and give it a dedicated, scoped-down credential rather than the developer’s own. Maintain an allowlist of reviewed MCP servers. Addresses R005.

MCP servers act as controlled gateways between the agent, the external tools, and data:

  • Dependency management – An MCP server fronting your private package registry resolves dependencies against curated packages, not the public internet. This is a deterministic constraint on supply chain risk.
  • Infrastructure tooling – Visibility into current resource configurations prevents templates that conflict with existing infrastructure.
  • Scoped permissions – Each MCP server exposes a defined set of tools and resources. You choose exactly what the agent can access, supporting least-privilege at the integration layer. You supply that credential through the agent’s configuration (in Kiro, the env block of .kiro/settings/mcp.json). Avoid autoApprove: ["*"], which removes the human approval prompt on every tool call.

IDE code scanning [D]

Control statement: Run real-time static analysis in the IDE so security issues surface while the developer (and agent) still have full context. Addresses R002, R006.

Real-time diagnostics catch syntax errors, type mismatches, and configuration issues as the developer types. A malformed IAM policy is flagged before the agent builds further on it. Security-focused extensions (ESLint security plugins, Checkov, SAST) layer on top for immediate feedback while code is fresh in context.

Hooks: Automated guardrails at the point of action [D + ND]

Control statement: Attach deterministic checks to file-save events and non-deterministic verification to task-completion events. Addresses R002, R007.

  • Shell command hooks [D] – Triggered on file save, these run a linter, formatter, or security scanner and produce the same result every time. They enforce hard rules.
  • AI-powered hooks [ND] – Triggered on task completion. These prompt the agent to verify that the implementation matches the specification and check for any untested edge cases or files that were modified outside the task’s scope.

Pillar 2: Build-time controls (in the pipeline)

Build-time controls run in the pipeline after code is committed and before it reaches production. They verify and gate what the agent produced, catching what author-time controls did not. The following controls apply at this stage.

Layered security scanning [D]

Control statement: Run secrets detection, static analysis, dependency scanning, and infrastructure-as-code scanning in sequence. Fail the build on any critical finding. Addresses R002, R003, R004.

  1. Secrets detection runs first because it’s cheapest and addresses a high-severity class of issue. It scans for hardcoded API keys, database connection strings, and credentials that AI agents might inadvertently include.
  2. SAST scans source code for injection issues, insecure deserialization, and resource leaks. Custom rules can target AI-specific anti-patterns including overly broad exception handling, deprecated APIs, placeholder credentials, dynamic code execution through eval().
  3. Software Composition Analysis (SCA) identifies known CVEs in dependencies. This is critical for AI-generated code, which might reference deprecated packages or hallucinate package names that open you to dependency confusion issues.
  4. Infrastructure as code (IaC) scanning validates AWS CloudFormation, Terraform, and AWS Cloud Development Kit (AWS CDK) templates against security policies before deployment. Catches overly permissive IAM roles, unencrypted storage, and public-facing resources the agent created.

Each stage halts the pipeline on failure. Results export to a standard format (Static Analysis Results Interchange Format (SARIF)) for compliance auditing and flow downstream to human reviewers. The open source Automated Security Helper (ASH) bundles secrets, SAST, SCA, and IaC scanners behind one command that you can run locally and in AWS CodeBuild, emitting SARIF for the gates that follow.

Quality gates [D]

Control statement: Define pass/fail thresholds for each scan type. Block deployment on any critical or high-severity finding. Addresses R003.

Quality gates convert scan results into go/no-go decisions. Define thresholds for each severity: block on critical findings, require justification for highs, and track mediums. The gate is deterministic: if a threshold is breached, the pipeline stops. Exceptions require documented approval.

Differentiate blocking compared to advisory modes: hard failures on main, advisory on feature branches. Avoid gates becoming a friction that teams route around.

AI-assisted review [ND]

Control statement: Use an LLM reviewer to pre-screen every pull request for specification compliance, scope creep, and security anti-patterns before human review. Addresses R001, R006, R007.

  • Specification compliance – Does the implementation match the requirements document?
  • Scope verification – Were files modified outside the task’s stated scope?
  • Security pattern review – Are there logic errors, misused APIs, or insecure patterns that pass SAST but violate intent?

This pre-screening focuses human reviewer attention on genuine risks rather than formatting or obvious issues. On AWS, AWS Security Agent (code review in preview at publication) checks pull requests against AWS-managed and custom security requirements. The reviewer screens and surfaces findings; the merge decision stays with a human.

A critical principle: the agent that wrote the code should not be the agent that reviews it. A separate session helps avoid self-confirmation bias, but a separate session alone doesn’t always avoid the generator’s blind spots, because two sessions of the same model can share them. Where practical, use a different model for review so the reviewer is less likely to inherit the same systematic weaknesses.

Human-in-the-loop review [ND + H]

Control statement: Require human approval on most pull requests, especially those touching security-sensitive or high-blast-radius code. Lower-risk changes might be eligible for agent-assisted or fully automated approval as tooling matures. Provide reviewers with scan results, LLM pre-screening output, and specification context to enable fast, informed decisions. Addresses R003.

Scale review depth to the risk of the change. Low-risk or boilerplate changes can take a lighter-touch review, while security-sensitive or novel-logic changes warrant mandatory deep review and a second reviewer.

Scanners catch known patterns but can’t judge whether code implements the intended business logic. Human review also serves to calibrate trust: teams build intuition about where agents excel (boilerplate, test writing) and where they’ve tended to struggle (novel business logic, security-sensitive operations), recognizing that this frontier shifts as models improve.

Place two approval gates: after security scans (reviewer focuses on correctness and business logic, with scan results as context) and before production deployment (final sign-off after integration testing). Treat human review as a secondary control, not a guarantee: reviewers are themselves non-deterministic and can miss issues, so human review layers on top of the deterministic gates rather than replacing them.

Putting the framework into practice on AWS

The framework is tool-agnostic, but AWS gives you building blocks for each pillar. The following services map directly to the controls described previously: Kiro for author-time guardrails, and CodeBuild and CodePipeline for build-time gates.

Kiro: Structured AI development

Kiro maps to Pillar 1: It puts the author-time controls in the IDE, where the developer and agent still share full context. Each feature in the following list implements one of those controls, configured in-repo under .kiro/ so the guardrails are version-controlled and shared across the team rather than set per developer.

  • Steering documents – Markdown files in .kiro/steering/ load into the agent’s context at session start. Conditional inclusion using fileMatch (for example, ["**/*.tf"]) loads IaC-specific rules only when relevant.
  • Specification-driven workflows – Three-phase specifications (requirements in EARS, design, and tasks) with review checkpoints. Bug-fix specifications capture unchanged behavior explicitly.
  • Agent hooks – Triggered on file save, tool invocation, or task completion. Shell hooks run deterministic checks (linters, tests); Ask Kiro hooks run AI prompts for non-deterministic review. For example, a security pre-commit scanner hook can flag hardcoded credentials when the agent finishes a task.
  • Property-based testing – Guided by a specification or hook, Kiro can generate property-based tests (for example, using the hypothesis library) that exercise hundreds of randomized inputs, probing edge cases a hand-written test suite would miss.
  • MCP integrations – Connect Kiro to private package registries, internal docs, issue trackers, and infrastructure tooling, creating the controlled tool access pattern.

For enterprise environments, Kiro supports AWS IAM Identity Center for single sign-on and provides IP indemnity coverage for subscribers. Check the Kiro documentation for current Region availability.

AWS CodeBuild and AWS CodePipeline: Pipeline controls

CodeBuild runs each scanning tool (checking for secrets, SAST, SCA, and IaC) as a build action. A non-zero exit code fails the action, and the stage halts or rolls back according to its OnFailure setting. Findings export as SARIF to Amazon S3 for compliance, and CodePipeline action variables pass results to downstream approval actions.

  • CodeBuild exit codes halt the pipeline on scan failures
  • AWS Lambda invoke actions evaluate scan results against configurable thresholds and return pass/fail decisions
  • Manual approval actions halt the pipeline, send Amazon Simple Notification Service (Amazon SNS) notifications, and link to review artifacts; decisions and reviewer identity are logged for audit

The following table consolidates the framework into a single view that includes each stage of the SDLC and the deterministic [D] and non-deterministic [ND] controls that apply there. Every stage carries both, a reminder that neither control type is sufficient on its own.

Stage Deterministic [D] Non-deterministic [ND]
IDE (pre-generation) Steering files loaded Steering documents, specification-driven constraints
IDE (post-generation) Shell hooks: Linter, formatter, type checker, and secrets scan AI-powered task completion hooks, context constraints
Pull request SAST, SCA, and IaC scanning LLM PR pre-screening and scope verification
Pipeline (pre-deploy) Full security scan suite, integration tests, and policy-as-code AI-assisted review for human approvers
Post-deploy Runtime monitoring and anomaly detection AI-powered incident triage

Conclusion

This post laid out a framework for adopting AI coding agents at machine speed without letting unreviewed risk reach production. It layers guardrails at two points:

  • Author-time controls – Steering, specs, and scoped tools shape what the agent generates in the IDE.
  • Build-time controls – Scanning, quality gates, and layered review verify it before it reaches production.

No single layer is enough: deterministic gates enforce hard rules, non-deterministic review catches what they miss, and human judgment is reserved for the decisions that need it. Together, they let AppSec scale with agent-driven development.

Where to start this week:

  1. Start with steering and specs – Encode security requirements as steering and use specifications for new features. Highest impact, lowest effort. For a ready-made starting set, the open source Project CodeGuard (a Coalition for Secure AI project under OASIS Open, of which Amazon is a contributing member) publishes reusable steering rules for common risk classes—hardcoded credentials, IaC misconfiguration, supply chain, and MCP security—that you can adapt to your AWS environment.
  2. Add deterministic pipeline gates – Integrate SAST, SCA, and secrets detection. Table-stakes regardless of AI usage.
  3. Calibrate and iterate – Review what controls catch, adjust steering for recurring issues, and expand agent autonomy as trust builds.
  4. Accountability – Developers remain accountable for the security of what they ship. AI agents accelerate development; they don’t transfer ownership.

More information:

If you have feedback about this post, submit comments in the Comments section below.


Daniel Begimher

Daniel Begimher

Daniel is a Senior Security Engineer at AWS, where he built and shipped the company’s first customer-facing AI security agent. He created SIR-Bench, a benchmark for measuring how deeply AI incident-response agents investigate before acting, and Automated Security Helper (ASH), an open source scanner. He co-leads application security technical field community at AWS, and speaks at conferences including AWS re:Invent, re:Inforce, and Cyber Week.

Danny Cortegaca

Danny Cortegaca

Danny is a Principal Security Specialist Solutions Architect and co-leads the Application Security focus area within the AWS Security and Compliance Technical Field Community. He joined AWS in 2021 and partners with some of the largest organizations in the world to help them navigate complex security and regulatory environments. He loves talking about application security with customers and has helped many adopt threat modeling into their practices.

Lowering AWS KMS decrypt API costs in EMR Spark jobs

Post Syndicated from Navaneedha Krishnan Jagathesan original https://aws.amazon.com/blogs/big-data/lowering-aws-kms-decrypt-api-costs-in-emr-spark-jobs/

Modern organizations processing vast amounts of data on Amazon EMR with Apache Spark face a growing cost challenge. As the number of encrypted Amazon Simple Storage Service (Amazon S3) objects grows, AWS Key Management Service (AWS KMS) decrypt API calls multiply rapidly, driving up operational costs. Consider a retail organization processing hundreds of terabytes of customer transaction data daily in S3 encrypted with AWS KMS. Each Spark task accessing an encrypted S3 object triggers an AWS KMS decrypt API call. At scale, these calls accumulate into significant and often unexpected cost increases. This is especially true for workloads that require key auditability and cannot switch to S3 Bucket Keys. S3 Bucket Keys reduce AWS KMS request costs by decreasing the number of calls from S3 to AWS KMS. However, S3 Bucket Keys limit per-object key auditability in AWS CloudTrail, which might not meet the compliance requirements of some organizations.

This post introduces practical techniques to reduce AWS KMS decrypt costs. You can reduce API call volume and lower costs without compromising encryption. It covers three techniques: optimizing file formats (including Apache Iceberg), aggregating data, and using AWS Glue Data Catalog partition indexes.

Optimization techniques

The following sections describe three techniques you can apply independently or together to reduce the number of AWS KMS decrypt API calls.

Use data aggregation

Data aggregation reduces redundant AWS KMS decrypt API calls by consolidating smaller files into larger blocks. When Spark reads many small files from S3, each file triggers its own decrypt call. By combining multiple small files into fewer, larger files, you can reduce the total number of AWS KMS API invocations. This technique is effective for read-heavy workloads that involve numerous small files stored on S3.

You can use AWS CloudTrail to monitor changes in API call frequency and validate the effectiveness of data aggregation in reducing costs.

Step 1: Benchmark baseline performance

Before applying optimizations, establish baseline metrics to quantify improvements.

from pyspark.sql import SparkSession
import time

spark = SparkSession.builder.appName("Baseline Job").getOrCreate()
start_time = time.time()
data = spark.read.format("csv").load("s3://amzn-s3-demo-bucket/data/")
end_time = time.time()
load_time = end_time - start_time
print(f"Load time: {load_time:.2f} seconds")
data_count = data.count()
print(f"AWS KMS calls triggered: {data_count} rows processed")

The following figure shows the number of AWS KMS Decrypt API calls captured in AWS CloudTrail. Use these baseline metrics to compare against optimized results in subsequent steps.

Amazon Athena console displaying CloudTrail log query results with a KMS Decrypt API call events triggered during the baseline Spark job reading unoptimized CSV files from S3

AWS CloudTrail log showing baseline AWS KMS Decrypt API call count

Step 2: Aggregate files using Spark

Consolidating smaller files into fewer, larger files stored in S3 minimizes redundant decrypt API calls.

consolidated_data = data.coalesce(10)
consolidated_data.write.mode("overwrite").parquet("s3://amzn-s3-demo-bucket/optimized-data/")

Step 3: Rerun the job with optimized files

Read the aggregated data created in Step 2 and compare the AWS KMS Decrypt API call count against the baseline metrics from Step 1.

optimized_data = spark.read.parquet("s3://amzn-s3-demo-bucket/optimized-data/")
optimized_data.count()

The following figure shows the AWS CloudTrail logs after reading the aggregated data.

Amazon Athena console displaying CloudTrail log query results with a reduced number of KMS Decrypt API call events after reading aggregated Parquet files

AWS CloudTrail logs in Amazon Athena showing AWS KMS Decrypt API call count after data aggregation

CloudTrail metrics comparison

Track the number of API calls and observe the direct impact of data aggregation on reducing AWS KMS decrypt API calls for the same amount of data.

The following figure compares the AWS KMS Decrypt API call count before and after data aggregation for the same dataset.

Comparison chart showing AWS KMS Decrypt API call count for the same dataset, with a significant reduction after consolidating small CSV files into fewer aggregated Parquet files

AWS KMS Decrypt API call comparison before and after data aggregation

Aggregating small files into fewer large files reduces decrypt calls and shortens load time.

Optimize file formats and compression

Selecting appropriate file formats and applying compression minimizes the amount of data read from S3 and the number of AWS KMS decrypt operations.

Columnar formats (Parquet/ORC)

Columnar file formats like Parquet and ORC let Spark read only the required columns for analysis, which improves performance for analytical queries. For example, you can convert raw CSV data to Parquet to benefit from better I/O efficiency and query optimization.

df = spark.read.format("csv") \
    .option("header", "true") \
    .option("inferSchema", "true") \
    .load("s3://emr-kms-demo/data/")

# Set compression for Parquet files
spark.conf.set("spark.sql.parquet.compression.codec", "snappy")

df.write.format("parquet").save("s3://amzn-s3-demo-bucket/parquet-data/")

Iceberg format

Apache Iceberg is a modern table format designed for large-scale analytic datasets. It supports schema evolution, snapshot isolation, and time travel, making it an excellent choice for data lakes on S3. When used with PySpark, Apache Iceberg simplifies data management by automatically optimizing file layouts, handling partitions, and integrating with Spark catalogs.

The following PySpark example uses Iceberg with Amazon EMR and S3:

pyspark \
  --packages org.apache.iceberg:iceberg-spark-runtime-3.4_2.12:1.4.2 \
  --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
  --conf spark.sql.catalog.spark_catalog=org.apache.iceberg.spark.SparkSessionCatalog \
  --conf spark.sql.catalog.spark_catalog.type=hadoop \
  --conf spark.sql.catalog.spark_catalog.warehouse=s3://amzn-s3-demo-bucket/iceberg-warehouse \
  --conf spark.sql.defaultCatalog=spark_catalog
df = spark.read.format("csv") \
    .option("header", "true") \
    .option("inferSchema", "true") \
    .load("s3://amzn-s3-demo-bucket/data/")

spark.conf.set("spark.sql.parquet.compression.codec", "snappy")

# Write data to Iceberg table
df.writeTo("iceberg_from_emr_data").using("iceberg").create()

# Read from Iceberg table
spark.read.table("iceberg_from_emr_data").show()

Compression

Using compression reduces data size and speeds up reads and writes between S3 and Spark. Note: ZSTD is the recommended and default compression codec for Iceberg, offering better compression ratios. For this demonstration, we use Snappy to illustrate the concept.

spark.conf.set("spark.sql.parquet.compression.codec", "snappy")

Compression not only minimizes I/O and network overhead but also accelerates job execution in distributed Spark environments.

CloudTrail comparison on API calls

The following figure shows the reduction in AWS KMS Decrypt API calls when using optimized file formats with compression.

Comparison showing AWS KMS Decrypt API call count for Parquet files with Snappy compression

AWS KMS Decrypt API calls on optimized file formats with compression

The following table illustrates the reduction in AWS KMS decrypt calls when moving from raw, uncompressed CSV data to optimized Parquet files with compression enabled.

Comparison showing AWS KMS Decrypt API call count for Uncompressed CSV vs Parquet files with Snappy compression

AWS KMS Decrypt API call comparison for CSV and compressed Parquet with Snappy

AWS Glue Data Catalog partition index

Partitioning data helps Spark jobs retrieve subsets of relevant data, reducing scan ranges, and decrypt operations. Using AWS Glue Data Catalog partition indexes reduces scanning overhead and the number of AWS KMS API calls.

Without a partition index, when Spark queries a partitioned table, AWS Glue Data Catalog returns all partitions by calling the GetPartitions API. Spark then reads every S3 object across all returned partitions. Because each S3 object is individually encrypted, Spark must call the AWS KMS Decrypt API once per object. More objects mean more decrypt calls and higher costs. With a partition index, AWS Glue performs server-side partition filtering, returning only matching partitions.

Step 1: Baseline query without partition index

Using an Amazon EMR Spark job:

spark.sql("SELECT * FROM default.`kms-demoevents` WHERE year='2000' AND month='04'").count()

Then check CloudTrail for the kms:Decrypt call count.

The following figure shows the AWS KMS Decrypt API call count when running the baseline query without a partition index. Spark scans all partitions, resulting in a higher number of decrypt calls.

Amazon Athena console displaying CloudTrail log query results showing the total number of AWS KMS Decrypt API calls triggered when querying without a partition index

AWS KMS Decrypt API call count without a partition index

Step 2: Add partition index and rerun the baseline query from Step 1

In the AWS Management Console or through the AWS Command Line Interface (AWS CLI), create partition indexes and rerun the same baseline query from Step 1. Create partition indexes on the year and month columns. Then check CloudTrail for the kms:Decrypt call count.

The following figure shows the AWS KMS Decrypt API call count after adding a partition index. With the partition index, AWS Glue filters partitions server-side, resulting in fewer S3 objects read and fewer decrypt calls.

Amazon Athena console displaying CloudTrail log query results showing a reduced number of AWS KMS Decrypt API calls after adding a partition index on year and month columns, compared to the baseline query without partition index

AWS KMS Decrypt API call count with a partition index

Conclusion

Optimizing EMR Spark jobs ensures cost-effective and efficient processing of encrypted data at scale. S3 Bucket Keys is the most effective way to reduce AWS KMS Decrypt API calls. The techniques covered in this post are additional optimizations that you can use together with S3 Bucket Keys for further cost reduction. You can also use them independently when S3 Bucket Keys cannot be used because of per-object auditability requirements in CloudTrail. Start implementing these strategies today to improve your Spark workload efficiency and achieve cost savings.

We welcome your feedback. If you have questions or suggestions about this post, leave a comment below.


About the author

Naveen Jagathesan

Naveen Jagathesan

Naveen is a Senior Technical Account Manager at AWS and focuses on driving operational excellence for customers. Outside of work, he is an avid gym enthusiast.

Automate Spark Scala migration to 4.x with AWS Spark Upgrade Agent

Post Syndicated from Bezuayehu Wate original https://aws.amazon.com/blogs/big-data/automate-spark-scala-migration-to-4-x-with-aws-spark-upgrade-agent/

If you’re a data worker responsible for managing Apache Spark 3.x workloads on Amazon EMR before, you’ve likely faced the challenge of migrating hundreds of jobs to Spark 4.0 without disrupting production pipelines. In this post, you will learn how to automate Spark 3.x to 4.0 migration using the AWS Spark Upgrade Agent covering API deprecations, behavioral changes, build configuration updates, and job validation. What once took months of manual effort can be completed in hours.

This is part 3 of a three-part series on how the AWS Spark Upgrade Agent can automate and simplify Spark upgrades.

Part 1 introduces the agent’s architecture and capabilities. Part 2 walks through a complete PySpark migration from Spark 3.5 to Spark 4.0 on Amazon EMR Serverless. This post walks through Scala migration from Spark 3.3 (Scala 2.12) to Spark 4.0 (Scala 2.13).

Apache Spark 4.0 on Amazon EMR 8.x delivers improvements like native merge_into() support, enhanced Adaptive Query Execution, improved Python UDF performance through Arrow-based serialization, and major Structured Streaming enhancements. For teams on Spark 2.4 or 3.x, the complexity lies in managing API deprecations, behavioral changes, build configuration updates, and re-validating hundreds of jobs while maintaining production pipelines.

Prerequisites

This post assumes you’ve completed the one-time AWS CloudFormation setup and proxy configuration detailed in the introduction post.

What is the Spark Upgrade Agent?

The AWS Spark Upgrade Agent is a fully managed remote server that automates Spark migration using a Model Context Protocol (MCP) interface for code analysis and transformation. For details, see the introduction post.

Architecture: How it works

The architecture follows a least-privilege security model:

  • Scoped IAM roles — AWS IAM roles are scoped to only MCP server calls, Amazon Simple Storage Service (Amazon S3) staging bucket access, and Amazon EMR job submission.
  • Local source code — Your source code stays local, with only minimal diagnostic information transmitted.
  • Encryption in transit — All data is encrypted in transit.
  • Audit trail — AWS CloudTrail records every tool invocation for full auditability.

Example use case: Enterprise-scale Spark migration

To illustrate the capabilities of the agent at scale, consider a real-world migration scenario from a large company. This company runs a data processing platform with thousands of Spark jobs across Scala, PySpark, and Spark SQL workloads, with a code base spanning Spark 3.3 and 3.5.

The data engineering team faces a migration across three workload types with different complexity profiles:

  • Spark SQL applications: The most portable, but still requiring validation of behavioral changes in the query optimizer and join strategies introduced in Spark 4.0.
  • PySpark workloads: Requiring updates to UDF serialization patterns, Arrow-based optimizations, and DataFrame API changes.
  • Scala applications: The most complex, involving build system updates (Maven and SBT), API deprecations, and recompilation against new Spark 4.0 JARs.

To tackle this, the team uses the Spark Upgrade Agent to migrate all three workload types to Spark 4.0 on Amazon EMR 8.x. For each workload, the agent is invoked directly from Kiro or VS Code with Cline, applying targeted transformations and immediately validating results against a live Amazon EMR 8.x Serverless application running Spark 4.0.

Migrations that would traditionally require months of manual engineering effort complete in a fraction of the time. The agent follows an iterative refinement loop: it performs local validation, then submits the job to a remote Amazon EMR cluster. This loop catches and resolves runtime failures automatically, reducing the need for manual debugging cycles. Build configuration files (pom.xml and build.sbt) are updated automatically by the agent, eliminating a common source of migration errors.

Solution walkthrough

The following sections walk through the complete migration workflow, from initial setup through advanced code transformations and validation.

1. Setup

This section lists what you need before starting. Some items, such as Amazon EMR Serverless applications, can be created during the walkthrough using the agent if they don’t already exist.

Must have before starting:

  • AWS Command Line Interface (AWS CLI) configuration: Your AWS CLI must be configured with a profile that has the necessary permissions. See Configuring the AWS CLI for details.
  • IAM role with Amazon EMR permissions: An AWS CloudFormation template is provided in the setup guide to provision the required IAM role. The role is scoped to the permissions needed for the upgrade process: calling the MCP server, reading and writing to the Amazon S3 staging bucket, and submitting Amazon EMR jobs.
  • Amazon S3 staging bucket for artifacts: Used to store code artifacts and Amazon EMR job outputs during the validation phase.
  • Integrated development environment (IDE) installation: Kiro or VS Code with MCP support (Cline extension). Either IDE can interact with the Spark Upgrade Agent through natural language prompts. Consult the setup guide for Kiro and Cline’s documentation to use the MCP server with Cline.
  • One-click MCP server installation: The dataprocessing-mcp server is installed and configured as described in the setup guide.
  • Amazon EMR Serverless applications: An Amazon EMR Serverless application is required for the validation workflow:
    • Target application (Spark 4.0): An Amazon EMR Serverless application configured with release label emr-spark-8.0.0, used to validate migrated jobs against Spark 4.0 on Amazon EMR 8.0.

1.1 Infrastructure setup (AWS CloudFormation)

Two AWS CloudFormation stacks create the required resources: an AWS IAM role, an Amazon S3 staging bucket, an Amazon EMR Serverless application (Spark 4.0), and its execution role.

Stack 1: AWS IAM role and Amazon S3 staging bucket

The spark-upgrade-mcp-setup template creates the AWS IAM role and Amazon S3 staging bucket required by the upgrade agent. Choose the Launch Stack button for your Region. For additional Regions, see the full Region list.

Region Launch
US East (N. Virginia) Launch Stack
US East (Ohio) Launch Stack
US West (Oregon) Launch Stack
Europe (Ireland) Launch Stack

After deployment, open the AWS CloudFormation Outputs tab, copy the ExportCommand value, and run it in your terminal. This sets SMUS_MCP_REGION, IAM_ROLE, and STAGING_BUCKET_PATH automatically.

The following figure shows the Outputs tab with the ExportCommand value.

AWS CloudFormation console Outputs tab with the ExportCommand value ready to copy

Outputs tab of the AWS CloudFormation stack showing the ExportCommand value

# Sets SMUS_MCP_REGION, IAM_ROLE, and STAGING_BUCKET_PATH
export SMUS_MCP_REGION=<YOUR-REGION> && export IAM_ROLE=arn:aws:iam::<YOUR-ACCOUNT-ID>:role/spark-upgrade-role-* && export STAGING_BUCKET_PATH=<amzn-s3-demo-bucket>

Then configure the AWS CLI profile:

aws configure set profile.spark-upgrade-profile.role_arn ${IAM_ROLE}
aws configure set profile.spark-upgrade-profile.source_profile default
aws configure set profile.spark-upgrade-profile.region ${SMUS_MCP_REGION}

Stack 2: Amazon EMR Serverless target application and execution role

The emr-serverless-target-setup template creates an Amazon EMR Serverless application configured with Spark 4.0 (release label emr-spark-8.0.0) and a shared execution role used for job submission during the validation phase. Deploy it as follows:

git clone https://github.com/aws-samples/sample-amazon-emr-spark4-examples
cd sample-amazon-emr-spark4-examples/scala3/demo_1_spark_change_focus

The Scala sample lives at sample-amazon-emr-spark4-examples/scala3/demo_1_spark_change_focus. The CloudFormation template lives at resources/cloudformation/.

Deploy the CloudFormation template to create the target Amazon EMR Serverless application and a shared execution role:

aws cloudformation deploy \
  --template-file resources/cloudformation/emr-serverless-target-setup.yaml \
  --stack-name spark-emr-serverless-upgrade \
  --region ${SMUS_MCP_REGION} \
  --capabilities CAPABILITY_NAMED_IAM \
  --parameter-overrides \
  StagingBucketName=${STAGING_BUCKET_PATH} \
  TargetReleaseLabel=emr-spark-8.0.0 \
  TargetApplicationName=spark-upgrade-target

This creates an Amazon EMR Serverless target application (Spark 4.0) for upgrade validation, with a shared execution role. The application auto-stops after 15 minutes of idle time, so there is no cost when not in use. To upgrade between different Spark versions, override the SourceReleaseLabel and TargetReleaseLabel parameters with the Amazon EMR release labels that you want.

After the stack completes, note the outputs:

aws cloudformation describe-stacks \
  --stack-name spark-emr-serverless-upgrade \
  --region ${SMUS_MCP_REGION} \
  --query "Stacks[0].Outputs" --output table

This gives you the TargetApplicationId and ExecutionRoleArn needed for the upgrade prompt. Make a note of them.

2. Upgrade

This section covers a complete end-to-end upgrade using a representative ecommerce pipeline, a Scala application that processes order events, applies transformations, and writes results using merge-style upsert patterns. The same workflow applies to Scala and Spark SQL workloads covered in subsequent sections.

Step 1: Clone the sample project

Start by cloning the sample project from the AWS samples repository:

git clone https://github.com/aws-samples/sample-amazon-emr-spark4-examples
cd sample-amazon-emr-spark4-examples/scala3/demo_1_spark_change_focus

The repository includes representative PySpark, Scala, and Spark SQL applications designed to demonstrate common Spark 3.x patterns and their Spark 4.0 equivalents.

Step 2: Open in your IDE and connect to the MCP server

Open the project in Kiro or VS Code with the Cline extension. Verify that the dataprocessing-mcp server is active and connected. You can see it listed as an available MCP server in your IDE’s MCP panel. If you haven’t completed the one-time setup, follow the setup guide before proceeding.

Step 3: Start the upgrade with a natural language prompt

Once connected, initiate the upgrade by entering the following request in the agent interface:

Use the dataprocessing-mcp server to upgrade my local project at <path-to-your-project>.
Upgrade my Spark application from Amazon EMR Serverless version 6.9.0 to Amazon EMR Serverless version 8.0.0.
Use Amazon EMR Serverless app-id <your-app-id> for validation.
Store artifacts at s3://amzn-s3-demo-bucket/spark4-upgrade/

The agent responds by invoking generate_spark_upgrade_plan, analyzing your project structure, identifying incompatible patterns, and presenting a prioritized upgrade plan before proceeding with any code changes.

After you confirm the plan, the agent proceeds autonomously through the remaining phases:

  1. Build configuration update — update_build_configuration rewrites pom.xml, build.sbt, or requirements.txt to target Spark 4.0 dependencies.
  2. Environment validation — Java and Python environments are checked and updated as needed.
  3. Code transformation — fix_upgrade_failure applies targeted fixes for each identified incompatibility, iterating until the project compiles cleanly.
  4. Remote validation — run_validation_job submits the upgraded application to your Amazon EMR Serverless target application and monitors execution through check_job_status.
  5. Data quality check (optional) — get_data_quality_summary compares output between the Spark 3.5 baseline and the Spark 4.0 run, confirming correctness before sign-off.

With the sample Scala ecommerce pipeline cloned from the sample-amazon-emr-spark4-examples repository and your IDE connected to the MCP server, you submitted a natural language prompt. This triggered the agent to analyze the project structure and generate a prioritized five-step upgrade plan, all before making any code changes.

Now that the upgrade plan is in place, the following sections walk through the specific code transformations the agent applies for a Scala workload.

Scala workload migration

This section covers the complete migration of a representative Scala Spark application from Amazon EMR Serverless 6.9.0 (Spark 3.3, Scala 2.12) to Amazon EMR Serverless 8.0.0 (Spark 4.0, Scala 2.13), using the demo_1_spark_change_focus sample from the AWSSpark4AutoUpgradeDemo repository.

Sample project: Ecommerce product change focus pipeline

The sample application processes product catalog change events from Amazon S3, applies enrichment transformations, and writes aggregated results back to Amazon S3. It represents a common pattern in ecommerce data platforms: incremental processing of catalog updates with downstream aggregation. In this example, we use VS Code with Cline, but you can also use Kiro or any other MCP-enabled IDE.

Project structure:

demo_1_spark_change_focus/
├── build.sbt
├── project/
│   ├── build.properties
│   └── plugins.sbt
└── src/
    └── main/
        └── scala/
            └── job_script.scala

The following figure shows the project structure as it appears in the IDE, with the build configuration and Scala source files.

IDE file tree showing the project layout with build.sbt, the project folder, and the Scala source file

Project structure with build.sbt and Scala source files in the IDE

Step 1: Invoke the upgrade agent

Enter the following prompt:

Use the dataprocessing-mcp server to upgrade my local project at
/Users/<your-path>/AWSSpark4AutoUpgradeDemo/resources/scala/demo_1_spark_change_focus/
from Amazon EMR Serverless version 6.9.0 to Amazon EMR Serverless version 8.0.0.
Use Amazon EMR Serverless app-id <your-app-id> for validation.
Store artifacts at s3://amzn-s3-demo-bucket/spark4-upgrade/scala/

The agent reads your project files, identifies the current Spark and Scala versions, and generates a structured upgrade plan (the analysis ID is logged for traceability). The following figure shows the agent generating the upgrade plan.

Agent output listing the detected Spark 3.3 and Scala 2.12 versions and the generated upgrade plan

Agent generating the upgrade plan with Spark and Scala version analysis

When prompted, choose Yes, proceed with the upgrade as planned, then choose Approve on the subsequent permission dialog to allow the agent to update the build configuration.

Step 2: Build configuration update — build.sbt

The first action of the agent is to invoke update_build_configuration, rewriting build.sbt to target Spark 4.0 and Scala 2.13, as shown in the following figure.

Agent output showing the build.sbt dependency changes from Spark 3.3 to Spark 4.0.1

Agent output showing build.sbt updates from Spark 3.3 to Spark 4.0

The build configuration update is generated. The following are the recommended changes for build.sbt.

Changes summary:

  1. Scala version: Updated from 2.12.15 to 2.13.16 (required for Spark 4.0).
  2. Spark dependencies: Updated from 3.3 to 4.0.
    • spark-core: 3.3 → 4.0.
    • spark-sql: 3.3 → 4.0.
  3. Assembly settings: Added configuration for creating uber JARs with proper merge strategies.
  4. Dependency exclusions: Added rules to exclude provided dependencies (Spark, Scala, Hadoop) from assembly.

The following figure shows the updated build.sbt and plugins.sbt files after the configuration changes are saved.

The updated build.sbt and plugins.sbt files open in the editor after the configuration changes

Updated build.sbt and plugins.sbt files saved after configuration changes

Step 3: Java environment check

Amazon EMR Serverless 8.0.0 runs on Java 17. The agent invokes check_and_update_build_environment to verify your local Java Development Kit (JDK) and upgrade it from Java 11 to Java 17, as shown in the following figure.

Agent output verifying the local Java version and recommending an upgrade to JDK 17

Agent verifying Java environment and recommending JDK 17 for Amazon EMR 8.0

Step 4: Scala source code transformations

After updating the build configuration, the agent compiles the project and applies fix_upgrade_failure iteratively to resolve Scala 2.13 and Spark 4.0 breaking changes. Scala 2.13 removed several deprecated collection methods that were available in 2.12. The compilation failed with errors related to the Scala 2.13 syntax change. The .to[Set] syntax needs to be updated to .to(Set) for Scala 2.13. The agent used the fix_upgrade_failure tool to resolve the compilation errors. The following are the key transformations applied to job_script.scala.

The following figure shows the agent applying Scala 2.13 source code transformations to resolve the compilation errors.

Agent output showing the Scala 2.13 source code edits applied to job_script.scala

Agent applying Scala 2.13 source code transformations to resolve compilation errors

// Disable ANSI (American National Standards Institute) (SQL compliance mode)(ANSI) mode to handle overflow and malformed cast operations
spark.conf.set("spark.sql.ansi.enabled", "false")

df.createOrReplaceTempView("airports")
// With ANSI mode disabled, overflow values will be handled gracefully
var new_df = spark.sql("SELECT *, CAST(build_time AS SMALLINT) as numeric_build_time FROM airports")
new_df.show()
new_df.createOrReplaceTempView("airports")

// Migration change: The to[Collection] method was replaced by the to(Collection) method.
val airports_in_us: Set[String] = spark.sql("SELECT name FROM airports WHERE country='USA'").collect().map(_.getString(0)).to(Set)
println(airports_in_us)
val airports_in_us_java: java.util.Set[String] = airports_in_us.asJava

// With ANSI mode disabled, malformed CAST operations will return null instead of failing
new_df = spark.sql("SELECT *, CAST(code AS INT) as numeric_code FROM airports")
new_df.show()
new_df.write
  .mode("overwrite")
  .parquet(outputPath)

Before

val airports_in_us: Set[String] = spark.sql("SELECT name FROM airports WHERE country='USA'").collect().map(_.getString(0)).to[Set]

After

val airports_in_us: Set[String] = spark.sql("SELECT name FROM airports WHERE country='USA'").collect().map(_.getString(0)).to(Set)

Code change explanation:

  • Scala 2.13 changed the collection conversion API. The .to[Collection] syntax was replaced with .to(Collection) using parentheses instead of square brackets.
  • Updated collection conversion from .to[Set] to .to(Set) to comply with Scala 2.13+ syntax requirements.
  • Changed import scala.collection.JavaConverters._ to import scala.jdk.CollectionConverters._ and updated .to[Set] to .to(Set).
  • Renamed the object from Spark3_3_Job to Spark4_0_Job. Updated the Parquet config keys from spark.sql.legacy.parquet.int96RebaseModeInRead/Write to spark.sql.parquet.int96RebaseModeInRead/Write.
  • Added spark.conf.set("spark.sql.ansi.enabled", "false") to handle overflow and malformed cast operations gracefully.
  • The output path was updated to s3://xxxxxxxxx/output.

After the compilation succeeds, the agent builds the assembly JAR. Choose Save to create a report for the build result.

Step 5: Runtime validation

Provide the following information to run the validation job on Amazon EMR Serverless.

Amazon EMR Serverless application ID (target application running Spark 4.0 on Amazon EMR 8.0.0):

  • To create an Amazon EMR application, follow the Amazon EMR documentation, or provide a prompt for the agent to create one for you.
  • Format: 00xxxxxxxxxxxxxxxxxxxxxxxxxx.

Execution role Amazon Resource Name (ARN) (IAM role for the job):

  • Set up the execution role following the IAM role guide.
  • Format: arn:aws:iam::123456789012:role/YourRoleName.

Amazon S3 staging path (for uploading the JAR and storing results):

  • Format: s3://amzn-s3-demo-bucket/path/.

AWS profile (the AWS profile to use for CLI commands, found in your mcp_settings.json file):

  • Example: default, dev, pro.

After you submit this information, the agent uploads the JAR to Amazon S3 and submits the validation job with the following arguments:

{
  "analysis_id": "a8869720-e005-41b1-89f3-620e1c5663c0",
  "application_type": "EMR-Serverless",
  "compute_id": "xxxxxxxxxxxxxx",
  "compute_run_config": {
    "executionRoleArn": "arn:aws:iam::xxxxxxxx:role/data-processing-mcp-role",
    "jobDriver": {
      "sparkSubmit": {
        "entryPoint": "s3://xxxxx/xxxxxxx/xxxxxxxxx-xxxxxxx-xxxxxxxx/xxx-job-assembly-1.0.jar",
        "entryPointArguments": [],
        "sparkSubmitParameters": "--class Spark4_0_Job --conf spark.executor.cores=4 --conf spark.executor.memory=16g --conf spark.driver.cores=4 --conf spark.driver.memory=16g --conf spark.executor.instances=2"
      }
    },
    "configurationOverrides": {
      "monitoringConfiguration": {
        "cloudWatchLoggingConfiguration": {
          "enabled": true,
          "logGroupName": "/aws/emr-serverless"
        }
      }
    }
  },
  "enable_data_quality_check": false,
  "s3_staging_path": "s3://xxxxxxx/xxxxxxxxx/",
  "is_source_version_run": false,
  "spark_metadata_file_put_presigned_url": null,
  "static_artifacts_file_paths": null
}

The agent monitors the job status upon approval.

Agent output showing the Amazon EMR Serverless validation job completing with a success status

Amazon EMR Serverless job validation output showing successful completion

Once you receive a success response, the agent proceeds to generate the upgrade summary for your Spark project.

{"success":true,"message":"EMR SERVERLESS job completed successfully","compute_run_id":"xxxxxxxxxx","compute_id":"xxxxxxxxxx","status":"SUCCESS","application_type":"EMR-Serverless","next_action":"Proceed to next step in the plan"}

After the job succeeds, you can optionally cross-check the results by navigating to your Amazon S3 bucket to verify the output files were written successfully, or by reviewing the job run details in the Amazon EMR Serverless console to confirm execution status and logs. Choose Approve to allow the agent to save the summary to a local file and upload it to Amazon S3.

You can also request ANSI mode enablement through the agent interface, and it applies the necessary updates accordingly.

Relevant Spark changes:

  • Migration rule: Spark 4.0 enables ANSI mode by default. To handle type conversion errors gracefully while keeping ANSI mode enabled, use TRY_CAST instead of CAST.
  • Change description: Enabled ANSI mode and replaced CAST with TRY_CAST for operations that might fail, specifically timestamp-to-smallint overflow and string-to-int malformed value conversions.

Applied changes:

  • Code diff — src/main/scala/job_script.scala:
    • Changed spark.conf.set("spark.sql.ansi.enabled", "false") to spark.conf.set("spark.sql.ansi.enabled", "true").
    • Replaced CAST(build_time AS SMALLINT) with TRY_CAST(build_time AS SMALLINT).
    • Replaced CAST(code AS INT) with TRY_CAST(code AS INT).

The agent compiles the change and follows the previous steps to run the job on Amazon EMR Serverless.

Result: SUCCESS

In this Scala workload migration section, the agent automatically upgraded the ecommerce pipeline from Spark 3.3 and Scala 2.12 on Amazon EMR 6.9.0 to Spark 4.0 and Scala 2.13 on Amazon EMR 8.0.0. It rewrote build.sbt and plugins.sbt, upgraded the JDK from 11 to 17, and applied Scala 2.13 syntax fixes (.to(Set), CollectionConverters), Parquet config key updates, and ANSI mode handling with TRY_CAST replacements. The upgraded JAR was compiled, submitted to Amazon EMR Serverless, and validated with a SUCCESS status, completing the full migration without manual code edits.

Clean up

To avoid ongoing charges, delete the resources created during this walkthrough. Start by emptying the Amazon S3 staging bucket, then delete both AWS CloudFormation stacks in reverse order:

  1. Empty the Amazon S3 staging bucket.
    aws s3 rm s3://${STAGING_BUCKET_PATH} --recursive

  2. Delete the Amazon EMR Serverless application stack.
    aws cloudformation delete-stack --stack-name spark-emr-serverless-upgrade

  3. Delete the MCP setup stack (IAM role and Amazon S3 bucket).
    aws cloudformation delete-stack --stack-name spark-upgrade-mcp-setup

Conclusion

The AWS Spark Upgrade Agent transforms what has traditionally been a months-long, error-prone migration process into an automated, IDE-driven workflow that completes in hours. By combining intelligent code analysis, targeted transformations, and an iterative local-to-remote validation loop, the agent handles the complexity of upgrading Scala workloads from Spark 3.x to Spark 4.0 on Amazon EMR 8.x. The demo_1_spark_change_focus walkthrough demonstrates the ability of the agent to automatically update build configurations, apply Scala 2.13 syntax changes, handle Spark 4.0 breaking changes like ANSI mode defaults, and validate results against live Amazon EMR clusters, all through natural language prompts in your IDE. For teams managing large-scale Spark estates, this approach eliminates manual debugging cycles, reduces migration risk, and unlocks the performance gains of Spark 4.0 without the traditional engineering overhead.

Next steps:

  • If you’re new to the Spark Upgrade Agent, start with the introduction post for a lighter-weight introduction before tackling Scala workloads.
  • For a complete PySpark implementation and demo, refer to Upgrade PySpark from Spark 3.5 to Spark 4.0 with AWS Spark Upgrade Agent.
  • When you are ready for production, review the security model in the Architecture section and the IAM role setup guide to confirm your least-privilege configuration before running against production workloads.

Useful resources:

Have questions or feedback? Share your migration experience in the AWS re:Post community or open an issue in the sample repository. We’d love to hear how the agent performs on your workloads.


About the authors

Bezuayehu Wate

Bezuayehu Wate

Bezuayehu is a Specialist Solutions Architect at AWS, specializing in big data analytics and AI-driven data processing. She works closely with customers to modernize analytics platforms using AWS data and AI services. With a passion for emerging technologies and customer success, she thrives on designing innovative cloud solutions that deliver measurable business impact and drive organizational transformation.

Prasad Nadig

Prasad Nadig

Prasad is a Senior Analytics Specialist Solutions Architect at Amazon Web Services (AWS), specializing in large-scale data analytics and AI. He partners with customers to tackle complex, large-scale data challenges guiding them as they design, migrate, and modernize their analytics platforms into solutions that are scalable, performant, and cost-effective. His expertise spans data lakes, data warehousing, and distributed data processing, with a strong focus on architectural best practices, performance tuning, and cost-optimization strategies that help organizations run analytics efficiently at petabyte scale.

Karthik Prabhakar

Karthik Prabhakar

Karthik is a Data Processing Engines Architect for Amazon EMR at Amazon Web Services (AWS). He specializes in distributed systems architecture and query optimization, working with customers to solve complex performance challenges in large-scale data processing workloads. His focus spans engine internals, cost-optimization strategies, and architectural patterns that enable customers to run petabyte-scale analytics efficiently.

Shubham Mehta

Shubham Mehta

Shubham is a Senior Product Manager at AWS Analytics. He leads generative AI feature development across services such as AWS Glue, Amazon EMR, and Amazon Managed Workflows for Apache Airflow (Amazon MWAA), using AI/ML to simplify and enhance the experience of data practitioners building data applications on AWS.

Keerthi Chadalavada

Keerthi Chadalavada

Keerthi is a Senior Software Development Engineer in the AWS analytics organization. She focuses on combining generative AI and data integration technologies to design and build comprehensive solutions for customer data and analytics needs.

Chuhan Liu

Chuhan Liu

Chuhan is a Software Engineer at AWS Glue. He is passionate about building scalable distributed systems for big data processing, analytics, and management. He is also keen on using generative AI technologies to provide brand-new experience to customers. In his spare time, he likes sports and enjoys playing tennis.

Building a scalable personalized recommendation system on AWS: From batch to real-time

Post Syndicated from Shraddha Anil Naik original https://aws.amazon.com/blogs/big-data/building-a-scalable-personalized-recommendation-system-on-aws-from-batch-to-real-time/

Amazon.com receives millions of visits every day, and behind every product recommendation on our website is a system that needs to process customer signals, run machine learning (ML) models, and deliver results before the next visit. Doing this across global marketplaces for millions of customers at tens of thousands of requests per second, while keeping experimentation fast and infrastructure costs bounded, is an orchestration challenge as much as a machine learning one.

Our team built a system that addresses this challenge. This post shows how we did it using a batch-first architecture with AWS Lake Formation, Amazon Managed Workflows for Apache Airflow (Amazon MWAA), Amazon Athena, AWS Glue, Amazon SageMaker, and Amazon DynamoDB, and how we later extended it with Amazon MemoryDB for real-time vector similarity search when we needed to incorporate more real-time signals.

Architecture overview

Data flow from the Lake Formation data lake and Athena through Airflow-orchestrated pipelines using Glue and SageMaker into DynamoDB for batch serving and MemoryDB for real-time inference

Data flows from the centralized data lake (Lake Formation and Athena) through Airflow-orchestrated pipelines using Glue for processing and SageMaker for ML workloads, into DynamoDB for batch serving and MemoryDB for real-time inference

The data lake foundation: Centralized access with Lake Formation

Every recommendation pipeline starts with data. We built a centralized data lake to create a single source of truth that any pipeline or consumer can access without duplicating data or building bespoke extract, transform, and load (ETL) pipelines.

Golden datasets: Shared once, used everywhere

Before the data lake, each recommendation pipeline independently extracted and transformed its own copy of product catalog, transaction history, and embeddings. This led to subtle inconsistencies: one pipeline might use a slightly different join logic or a stale snapshot, making it difficult to compare model performance or debug discrepancies across pipelines.

Now, we publish curated, validated datasets once, and every consumer (Airflow DAGs, ML notebooks, analytics dashboards) reads from the same tables through the same governed access. This means:

  • New pipelines start faster. A new recommendation model does not need its own data extraction logic. It queries the existing golden datasets from day one.
  • Consistency across models. When we compare model A to model B, we know both models are trained and inferred on the same underlying data.
  • Cross-team collaboration. Multiple teams share the same tables as a single source of truth.
  • Two-way data flow. The data lake serves as both source and destination. Our pipelines read golden datasets as inputs and write computed outputs (model scores, feature sets, intermediate results) back to the lake, where they become inputs for other pipelines. This creates a compounding effect: each new pipeline enriches the lake for the next one.

Why Lake Formation?

Our data is stored in Amazon Simple Storage Service (Amazon S3), partitioned by marketplace. Our data consumers (Airflow pipelines, ML notebooks, analytics tools) live in separate AWS accounts from the data producers. We chose AWS Lake Formation because it adds governance on top of S3 without requiring data migration:

  • Fine-grained cross-account access. Grant table-level permissions per consumer role, without managing bucket policies manually.
  • Schema governance through AWS Glue Data Catalog. Scheduled Glue crawlers infer schemas from files in S3, keeping the catalog current as data evolves.
  • Multi-Region consistency. We deploy identical infrastructure across multiple AWS Regions using AWS Cloud Development Kit (AWS CDK), each Region serving its local marketplaces.

Orchestration: Amazon Managed Workflows for Apache Airflow (Amazon MWAA)

We chose Amazon Managed Workflows for Apache Airflow (Amazon MWAA) as our orchestration layer. MWAA removes the operational burden of managing Airflow infrastructure: automatic scaling of workers, built-in high availability, and managed upgrades mean our team focuses on pipeline logic rather than cluster maintenance. MWAA lets us define complex multi-step workflows with rich dependencies in Python code, and its operator model lets us encapsulate team-specific conventions into reusable building blocks.

Each recommendation pipeline follows a consistent pattern:

The consistent recommendation pipeline pattern moving from data extraction through model training, batch inference, vector search, ranking, and publishing

We built a library of reusable custom Airflow operators, each encapsulating one of our core compute engines. This reduced new pipeline development from weeks to days in our team’s experience.

Athena and AWS Glue: Data access and processing

Amazon Athena is how our pipelines read from Lake Formation. Our custom operator runs SQL queries against the Glue Data Catalog and automatically runs UNLOAD to write results to S3: serverless, no infrastructure to manage, and integrated with the Lake Formation permission model.

AWS Glue handles compute-intensive data transformations through PySpark: joining datasets, filtering, deduplication, aggregation, and formatting ML outputs into final recommendation lists. We configure Glue with Auto Scaling worker pools (for example, 2–50 workers of G.8X type) so jobs scale with data volume per marketplace. All jobs run ephemerally: they read from S3, write to S3, and require no long-running clusters.

Amazon SageMaker powers the ML-intensive stages of our pipelines across three workload types, all orchestrated as steps within our Airflow DAGs:

Training

We use SageMaker Training Jobs to train our recommendation models on GPU instances (for example, ml.g5). Training data is prepared by upstream Glue jobs and staged in S3. Our custom Airflow operator submits the training job, monitors its progress, and registers the resulting model artifact in S3. Once training completes, the artifact is immediately available for batch inference or endpoint deployment within the same DAG run. This means a single DAG can go from raw data to trained model to deployed inference without manual handoffs, and we can retrain it on fresh data every pipeline cycle with zero operator intervention.

Batch inference

SageMaker Batch Transform runs our trained models at scale, generating the outputs that feed into downstream ranking and publishing steps. Our batch inference operator handles job submission, polls for completion, and writes the output location to the DAG’s S3 convention so the next Glue step can pick it up automatically. Batch Transform lets us run inference without provisioning persistent infrastructure, and we can scale instance count and type independently for every pipeline based on data volume.

Generating recommendations for millions of customers requires searching across hundreds of thousands of candidate products per customer. Exact search at this scale is prohibitively expensive, so we use approximate nearest neighbor (ANN) search using FAISS to find similar products efficiently, run as SageMaker Processing Jobs.

The workflow:

  1. Build a FAISS index over the candidate catalog.
  2. Query the index with per-customer vectors to find top-K nearest neighbors.
  3. Return ranked candidate lists per customer.

We distribute query vectors across the SageMaker Processing fleet using S3-based sharding (ShardedByS3Key). Each instance receives the full candidate index but only a fraction of the query vectors. Every instance builds an identical FAISS index, searches its shard of queries, and writes results to S3. The downstream Glue step merges all shards into the final recommendation lists. This lets us scale horizontally by adding instances without changing any code.

Why SageMaker inside MWAA?

Running SageMaker jobs as MWAA tasks (rather than standalone) gives us:

  • End-to-end lineage. Every model training run, inference job, and vector search is tracked as part of a DAG execution. We can trace a recommendation in DynamoDB back to the exact training run, data snapshot, and ANN search that produced it.
  • Retry and failure handling. If a SageMaker job fails (spot instance preemption, transient capacity errors), Airflow retries it automatically with backoff. No manual re-runs.
  • Resource sequencing. Training must finish before inference, inference before ANN search. The Airflow dependency model handles this naturally without polling scripts or step function state machines.
  • Unified monitoring. One Airflow dashboard shows the health of all pipelines: Glue ETL, SageMaker training, SageMaker inference, and DynamoDB publishing. No context-switching between consoles.

Putting it together: A complete pipeline example

Our recommendation generation pipeline illustrates the full flow:

The complete recommendation generation pipeline: data lake extraction, model training, batch inference, vector search, merge and rank with Glue, and publishing to DynamoDB

Note that the operators shown (GlueSQLOperator, VectorSearchOperator, and others) are custom internal operators built on top of the Airflow AWS provider, not open-source libraries. Here is a simplified version of what this looks like in code:

from airflow import DAG
from airflow.models.baseoperator import chain
from airflow.utils.task_group import TaskGroup

dag = DAG("recommendation_generator", schedule="0 18 * * 4")  # Weekly
task_groups = []

for marketplace in [...]:  # global marketplaces
    with TaskGroup(group_id=marketplace, dag=dag) as group:

        # Step 1: Extract data from lake
        candidates = GlueSQLOperator(
            task_id="select_candidates",
            tables=[f"{marketplace}.catalog", f"{marketplace}.products"],
            sql="select_candidates.sql", dag=dag)

        customer_history = GlueSQLOperator(
            task_id="select_customer_history",
            tables=[f"{marketplace}.transactions"],
            sql="select_customer_vectors.sql", dag=dag)

        customers = AthenaSQLOperator(
            task_id="select_customers", database=marketplace,
            query="select_customer_cohort.sql", dag=dag)

        # Step 2: Train model
        train = ModelTrainingOperator(
            task_id="train_model",
            training_input={"task_id": customer_history.task_id},
            instance_type="ml.g5", dag=dag)

        # Step 3: Batch inference
        inference = BatchInferenceOperator(
            task_id="batch_inference",
            model={"task_id": train.task_id},
            input_data={"task_id": candidates.task_id}, dag=dag)

        # Step 4: Vector search via SageMaker Processing Job
        ann_search = VectorSearchOperator(
            task_id="ann_search",
            index_input={"task_id": inference.task_id},
            query_input={"task_id": customer_history.task_id},
            k=60, instance_count=10, dag=dag)

        # Step 5: Merge and rank with Glue
        merge_and_rank = GlueSparkOperator(
            task_id="merge_and_rank",
            script="merge_results.py", dag=dag)

        # Step 6: Publish to DynamoDB
        publish = DynamoDBPublishOperator(
            task_id="publish_recommendations",
            table_name="Recommendations", dag=dag)

        # Task dependencies
        [candidates, customer_history] >> train >> inference
        [inference, customers] >> ann_search >> merge_and_rank >> publish

    task_groups.append(group)

chain(*task_groups)  # Execute marketplaces sequentially

Key patterns in this DAG

  • Marketplace isolation. Each marketplace runs in its own TaskGroup. A failure in one does not block the others.
  • Parallel data extraction. Independent Athena/Glue queries run concurrently before converging at the training step.
  • Sequential marketplace execution. chain() runs marketplaces one at a time to avoid resource contention across large SageMaker and Glue jobs.
  • Reusable operators. We built custom operators like GlueSQLOperator, VectorSearchOperator, and DynamoDBPublishOperator that encode our team’s conventions (cross-account access, S3 staging, retry logic, metrics) into shared building blocks. New pipelines are mostly configuration rather than infrastructure code, and these operators are shared across dozens of pipelines.
  • Implicit data passing. Each operator writes its output to a convention-based S3 path and downstream operators automatically resolve the upstream output location. No hard-coded paths between steps.

This pattern powers hundreds of pipelines across global marketplaces, with each marketplace as an isolated TaskGroup in the DAG.

Serving layer: DynamoDB + ECS

Our Java-based Amazon Elastic Container Service (Amazon ECS) service reads pre-computed recommendations from Amazon DynamoDB at request time:

  1. Receive request with customer ID, marketplace, customer context, and page context.
  2. Read pre-computed recommendations from DynamoDB.
  3. Apply real-time filters (availability, eligibility constraints).
  4. Re-rank based on real-time context signals.
  5. Return the response.

Amazon DynamoDB is designed to provide single-digit millisecond reads at our scale and time-to-live (TTL) for automatic cleanup of stale recommendations.

Extending to real-time

In recommendation system terms, our batch pipeline handles the retrieval stage offline. Although this covers the majority of our traffic, we identified scenarios where weekly freshness was not enough to capture what customers are doing right now. The real-time extension adds an online retrieval and ranking path for signals that cannot wait for the next batch cycle.

To address these, we extended our system with Amazon MemoryDB (which now supports Valkey as its open-source engine) for real-time vector similarity search and SageMaker real-time endpoints for on-demand embedding generation. The same Airflow pipelines that publish to DynamoDB also publish product vectors to MemoryDB (through our MemoryDB publish operator), and the same model in our batch pipeline is deployed to a SageMaker endpoint for single-item inference at request time.

At serving time, when we want to incorporate fresh signals, the service calls the SageMaker endpoint to generate an embedding on the fly, then queries MemoryDB for the nearest neighbors. These fresh signals include recent search queries, cart additions, and other in-session activity that changes faster than our weekly batch cycle. In our workloads, this gives us sub-millisecond vector search latency without re-running the full batch pipeline. Critically, the batch pipeline keeps the MemoryDB product index fresh. Our Airflow DAG includes a MemoryDB publish operator that refreshes the full product vector index weekly, so real-time queries always search against an up-to-date index.

Batch and real-time are not competing approaches: batch handles slow-moving signals (purchase history, catalog relationships) while real-time handles fast-moving ones (current session, new arrivals, trending items). Both paths share the same models, the same data lake, and the same serving service.

For a detailed deep-dive on this real-time architecture, see Real-time personalized recommendations with Amazon SageMaker and Amazon Managed Valkey.

Security and access control

Security is a foundational concern for a system that spans multiple AWS accounts, processes customer behavioral data, and runs across multiple AWS Regions.

Cross-account access: Lake Formation grants are issued to consumer-account AWS Identity and Access Management (IAM) roles, ensuring consumers can query tables without direct S3 bucket access. Each consumer role receives only the permissions it needs for its specific tables.

IAM execution roles: Each compute engine (MWAA, Glue, SageMaker) runs under a dedicated least-privilege IAM role.

Network isolation: MWAA environments and the ECS serving layer are deployed within virtual private clouds (VPCs), with separate VPC configurations per Region.

Reliability and failure handling

Regional isolation: Each AWS Region runs an independent copy of the system. DynamoDB tables are regional, each populated by the local batch pipeline. MemoryDB clusters are regional, with the product vector index refreshed by the local Airflow DAG. This means a regional failure or pipeline delay in one Region does not affect other Regions.

Batch pipeline failures: If a pipeline fails mid-run, the previous DynamoDB data remains live and continues serving recommendations until the next successful run. Airflow retries failed tasks automatically with configurable backoff. TTL on DynamoDB records bounds how long stale data persists. Failed pipeline runs trigger automated alerting so the on-call engineer can investigate.

Real-time path resilience: MemoryDB is deployed in a multi-AZ configuration with automatic failover. The real-time path is an extension of the batch system, and batch recommendations from DynamoDB remain available regardless of real-time path availability.

Lessons learned

  1. Start batch, add real-time incrementally. Batch pipelines are easier to debug, cheaper to operate, and sufficient for most recommendation scenarios. Add real-time path only when you have clear indication that specific customer signals (for example, in-session activity, search queries) need sub-hour freshness to remain relevant.
  2. Match the serving path to signal velocity. Not every signal needs real-time processing. Categorize your signals by how quickly they change, and route accordingly: batch for slow signals, real-time for fast ones, re-ranking for in-between.
  3. Freshness does not always require re-computation. A batch-generated candidate set remains largely valid between runs. What changes is relative relevance. Re-ranking at serving time with recent customer activity gives the impression of real-time without the cost of real-time candidate generation.
  4. A centralized data lake accelerates everything. Golden datasets eliminated weeks of per-pipeline data extraction work, made model comparisons trustworthy, and let new team members ship their first pipeline in days instead of weeks. The upfront investment in Lake Formation governance paid for itself within the first quarter.
  5. Invest in reusable operators. Custom Airflow operators encapsulating Athena/Glue/SageMaker patterns let teams ship new pipelines in days. The operators encode best practices (retry logic, cross-account access, metrics) so pipeline authors can focus on business logic.
  6. Separate compute from storage. S3 as the universal intermediate layer + ephemeral Glue/SageMaker jobs means you pay only for active computation. No idle clusters between weekly pipeline runs.

Conclusion

We built this system because every customer interaction is an opportunity to surface the right product at the right time. Serving tens of thousands of requests per second across millions of customers in global marketplaces, our batch-first architecture uses AWS Lake Formation for governed data access, Amazon MWAA for orchestration, Amazon Athena for data lake queries, AWS Glue for distributed processing, Amazon SageMaker for training, inference, and vector search, and Amazon DynamoDB for low-latency serving. When we needed to incorporate more real-time signals, we added Amazon MemoryDB vector search paired with SageMaker real-time endpoints for on-demand embedding generation, extending into real-time without replacing the batch foundation.

The architecture choices we have made, batch for efficiency and real-time for freshness, all serve the goal of helping customers discover what they need faster. If you are building personalized experiences at scale, we hope these patterns give you a useful starting point.

This would not have been possible without the Everyday Essentials engineering team, whose collective effort turned these ideas into a production system serving customers every day. We are also grateful for the support and guidance from Sam Heyworth, Nirav Desai, and Ankur Datta, and the broader Everyday Essentials leadership team.


About the authors

Shraddha Anil Naik

Shraddha Anil Naik

Shraddha is a Senior Software Engineer on the Everyday Essentials team at Amazon. She specializes in retrieval and recommendation infrastructure that powers personalized experiences for millions of customers.

Sergii Oborskyi

Sergii Oborskyi

Sergii is a Senior Software Engineer on the Everyday Essentials team at Amazon. He builds online recommendation services that serve product recommendations to millions of customers at high throughput and low latency.

Shawn Liu

Shawn is a Senior Machine Learning Engineer on the Everyday Essentials team at Amazon. He develops and evaluates recommendation models on Amazon SageMaker that power personalized product discovery for millions of customers.

Walter Wong

Walter Wong

Walter is a Software Development Manager in the Everyday Essentials Science org at Amazon. His work focuses on customer understanding and personalization, improving product recommendations for millions of customers across Amazon’s everyday essentials catalog.

Accelerating AWS Network Firewall troubleshooting with AWS DevOps Agent

Post Syndicated from Salman Ahmed original https://aws.amazon.com/blogs/security/accelerating-aws-network-firewall-troubleshooting-with-aws-devops-agent/

When an administrator introduces a rule change in AWS Network Firewall and network connectivity is disrupted, pinpointing the cause requires inspecting multiple points in the traffic path. The firewall gives you stateless and stateful rule engines, domain rules, and routing to the firewall endpoint inside your Amazon Virtual Private Cloud (Amazon VPC). A network drop looks the same from the workload no matter where it started. Isolating the cause means correlating the alert and flow logs with the firewall configuration, route tables, and recent API calls in AWS CloudTrail that might have changed them. That manual correlation is exactly where AWS DevOps Agent helps, accelerating root cause analysis so you can restore connectivity in minutes instead of hours.

AWS DevOps Agent does that correlation for you. As your always-available operations teammate, it resolves and proactively prevents operational issues across AWS, multicloud, and on-premises environments. When an Amazon CloudWatch alarm triggers, it reaches the agent through a webhook. The agent then reads the firewall configuration and logs through AWS APIs, ties the drop to recent API activity, and returns a root cause with a mitigation plan you review before you apply it.

This post connects CloudWatch monitoring to DevOps Agent. It walks through three Network Firewall failures from end to end. The first is a domain deny list blocking a legitimate endpoint. The second is a stateless rule priority misconfiguration. The third is an asymmetric cross Availability Zone (AZ) routing drop. Each maps to a different layer, so each leads down a different investigation path. An AWS Cloud Development Kit (AWS CDK) app deploys the whole environment in your own account so you can reproduce each failure and follow along.

The sample workload

As part of this blog post, we provide a CDK stack that deploys both the AWS DevOps Agent Space and a sample workload used to walk through three separate troubleshooting scenarios. A single t3.micro instance in a protected subnet checks its connectivity to a test endpoint on a continuous loop and publishes results to CloudWatch. Traffic takes the internet egress path through Network Firewall, the NAT gateway, and the internet gateway, so the firewall can intercept or drop it. After completing the walkthrough, you can apply the same troubleshooting techniques with DevOps Agent against your own Network Firewall deployments.

The test endpoint runs in a separate VPC deployed by the same CDK app. It serves HTTPS on port 443 and TCP on port 9142, giving each scenario a different protocol layer to exercise: Scenario 1 targets a TLS connection on 443 (matched by Server Name Indication), Scenario 2 targets a TCP connection on 9142, and Scenario 3 exercises the whole egress path.

A live status page shows one card per scenario plus the network topology. The whole stack deploys from a single CDK app across two Availability Zones, each with a firewall endpoint and NAT gateway, which is what makes Scenario 3 possible.

As shown in the following figure, the egress data path runs from the workload through Network Firewall and the NAT and internet gateways to the test endpoint. The alarm pipeline runs from CloudWatch through Amazon Simple Notification Service (Amazon SNS) and the webhook AWS Lambda function to DevOps Agent.

Figure 1: The sample workload

Figure 1: The sample workload

To use this with your own workload, you need a CloudWatch alarm that detects the connectivity problem and the webhook pipeline (SNS topic and Lambda function) that delivers it to DevOps Agent. The agent reads your firewall configuration, logs, and CloudTrail through AWS APIs, so no additional instrumentation is needed on the firewall side.

Prerequisites

To follow along with this post, you need:

Deploy the sample workload

Clone the project and deploy it into us-east-1 with one command (set awsRegion to use another AWS Region).

git clone https://github.com/aws-samples/sample-accelerating-aws-network-firewall-troubleshooting-with-aws-devops-agent.git
cd sample-accelerating-aws-network-firewall-troubleshooting-with-aws-devops-agent
bash scripts/deploy.sh

The script checks prerequisites, installs dependencies, compiles and tests, and bootstraps the CDK if needed. It then deploys all the stacks from a clean baseline and prints the outputs, including the status-page URL and sign-in details.

  1. Open the status-page link (an https://<random-id>.cloudfront.net address).
  2. Sign in using the username and password provided from the CDK output and confirm all three cards show the green Healthy status.
  3. Keep the page open while you run the scenarios.

Connect AWS DevOps Agent

To connect AWS DevOps Agent to the alarm pipeline

  1. In the AWS DevOps Agent console, open the nf-devops-agent-space Agent Space created by the CDK deployment.
  2. Configure the DevOps Agent webhook and download the CSV file with the webhook URL and signing secret.
  3. On the status page, choose Configure webhook, paste the URL and signing secret, and save. The page writes them to the nf-devops-agent-webhook-credentials AWS Secrets Manager secret, so there is no AWS CLI or console step. Until you set it, the bridge Lambda function sees a placeholder and skips delivery.
  4. Verify the path before you run a scenario. In the Lambda console, open nf-devops-agent-webhook and use the Test tab with this event.
    {
      "Records": [
        {
          "Sns": {
            "Message": "{\"AlarmName\":\"TEST-webhook-verification\",\"AlarmDescription\":\"[TEST] Webhook integration test - not a real alarm.\",\"NewStateValue\":\"ALARM\",\"NewStateReason\":\"[TEST] Manual webhook connectivity test. Safe to ignore.\",\"Region\":\"us-east-1\"}"
          }
        }
      ]
    }

  5. A 200 response confirms the path, and a test investigation appears in the DevOps Agent Operator Web App view.

How the alarm pipeline works

Every scenario reaches DevOps Agent the same way. A CloudWatch alarm moves to ALARM and notifies the SNS topic. Amazon SNS invokes a Lambda function. The function reads the webhook URL and signing secret from Secrets Manager, signs an alarm payload, and POSTs it to the DevOps Agent webhook (as shown in Figure 1). Amazon SNS also provides delivery retries, fan-out to other subscribers, and cross-account publishing.

  • Prebuilt Network Firewall metric (Scenario 1) – Alarm-1 watches the DroppedPackets metric, summed across the stateful streams, and triggers when drops rise above a baseline threshold. This requires no workload or custom metric and works on an already-deployed firewall. However, it only tells you that the firewall is dropping packets, not which rule is responsible.
  • Application health metric (Scenarios 2 and 3) – Alarm-2 and Alarm-3 watch a custom metric from a connectivity check. Use this for an alarm tied to user-facing impact or to tell one traffic path from another, which requires running a component that emits the metric.
Alarm Source Triggers when
Alarm-1 Native AWS/NetworkFirewall DroppedPackets The firewall’s dropped-packet count rises above the baseline
Alarm-2 Custom application health metric The port 9142 (TCP) connectivity check to the test endpoint is being dropped
Alarm-3 Custom application health metric The cross Availability Zone connectivity check is being dropped

Run the scenarios

Work through each of the scenarios one at a time, following the same cycle. Interrupt network connectivity, watch the alarm trigger, let DevOps Agent investigate, apply the recommended fix, and confirm recovery before moving on.

The status-page cards follow the live CloudWatch alarm state. A card shows a green dot and the word Healthy when its alarm is clear, and a red dot and the word DROPPED when its alarm triggers. In the DROPPED state the card also adds a Condition: line describing what’s being dropped, which isn’t shown when the card is healthy. Network Firewall applies changes to new flows, so a change shows within a minute or two. Recovery comes from the mitigation DevOps Agent recommends, which you review and apply.

Scenario 1. Domain deny list blocking a legitimate endpoint

At baseline, the rg-domain Suricata domain rule group denies only an unused placeholder, so the test endpoint stays reachable. The rule group inspects the TLS Server Name Indication (SNI) on each outbound connection and drops any that matches a denied domain. The exact rule syntax and console steps follow.

To add the domain deny rule

  1. Go to the Amazon VPC console.
  2. In the navigation pane, under Network Firewall, choose Network Firewall rule groups.
  3. Choose the rg-domain rule group to open its details page.
  4. In the Rules section, choose Edit.
  5. The rules box already contains two baseline placeholder rules (they match blocked.placeholder.invalid, so nothing real is denied). Leave those in place. Find the <app-endpoint-dns> value for Scenario 1 in the deployment script output (a Nework Load Balancer (NLB) DNS name such as NfTest-AppNl-a1b2C3dEf4G5-1234abcd5678efgh.elb.us-east-1.amazonaws.com). On a new line below the existing rules, add a drop rule that matches that DNS name on the TLS SNI, then choose Save.
    drop tls $HOME_NET any -> $EXTERNAL_NET any (ssl_state:client_hello; tls.sni; content:"<app-endpoint-dns>"; startswith; nocase; endswith; msg:"S1 domain denylist"; flow:to_server, established; sid:2000002; rev:1;)

  6. After saving, the rules box holds all three lines. The two placeholders remain, plus the new drop rule for the endpoint DNS name (note the distinct sid 2000002).
Figure 2: Scenario 1 – Firewall rule change blocking the connection

Figure 2: Scenario 1 – Firewall rule change blocking the connection

What happens. The workload’s HTTPS check to the test endpoint times out, the “AWS/NetworkFirewall DroppedPackets metric climbs above baseline, and Alarm-1 moves to ALARM. The Scenario 1 card reads DROPPED (with the condition Firewall dropping the monitored domain on its allow/deny rules), while the Scenario 2 and Scenario 3 cards stay Healthy (Figure 3). On the topology, the alarm pipeline from CloudWatch through Amazon SNS and Lambda to DevOps Agent and the workload-to-firewall inspect lines both turn amber, which the legend defines as collateral / alarm active, because the packets are now dropped at the firewall. To demonstrate the resulting failure, the HTTPS · SNI line from the internet gateway to the test endpoint is shown in red, which the legend defines as dropped (root cause).

Figure 3: Scenario 1 active – Traffic blocked at the firewall

Figure 3: Scenario 1 active – Traffic blocked at the firewall

Let DevOps Agent investigate. The agent runs several lines of investigation in parallel and correlates them:

  1. Reads the DroppedPackets metric and correlates the spike with a simultaneous drop in passed packets, confirming the firewall is actively blocking traffic.
  2. Reads the ALERT log and finds the workload’s TLS connections to the test endpoint blocked by the S1 domain denylist rule.
  3. Compares the current state against a baseline window, where the same endpoint was reachable with no alerts, which shows the block is new.
  4. Searches CloudTrail and surfaces the UpdateRuleGroup call that added the deny rule, identifying the user, role, and timestamp approximately one minute before the drops began.
  5. Reports the root cause as that manual rule-group change. Recommends removing the deny entry or adding an allow exception and enabling FirewallPolicyChangeProtection to prevent unauthorized changes.
  6. Presents this as a plan you review and apply, not an automatic change.

In the DevOps Agent Operator Web App view, the agent first restates the Alarm-1 trigger and confirms the firewall is dropping packets above the threshold (Figure 4).

Figure 4: Scenario 1 – The symptom

Figure 4: Scenario 1 – The symptom

Next, the agent identifies the root cause: a manual update to the rg-domain rule group that added a domain deny rule (SID 2000002) shortly before the alarm fired, blocking TLS connections to the ELB endpoint (Figure 5).

Figure 5: Scenario 1 – The root cause

Figure 5: Scenario 1 – The root cause

Finally, the agent presents a mitigation plan, recommending you remove the problematic deny rule (SID 2000002) to restore connectivity (Figure 6).

Figure 6: Scenario 1 – The mitigation plan

Figure 6: Scenario 1 – The mitigation plan

Note: In a real-world environment, this type of rule typically exists for a reason. Before removing it, verify whether it was intentional but scoped too broadly. If so, refine the rule to block only unauthorized endpoints rather than removing it entirely.

Confirm recovery. Apply the change the agent recommends. After the deny entry is gone, DroppedPackets falls back to baseline, Alarm-1 clears, and the card returns to green. Move on to Scenario 2.

Scenario 2. Stateless rule priority misconfiguration

At baseline, the rg-stateless-priority stateless rule group keeps the allow rule at priority 100 and the drop rule at 200 for the test class, TCP destination port 9142. The workload opens a TCP connection to the test endpoint on this port. Lower priority numbers evaluate first, so the allow rule wins. This scenario uses port 9142 instead of 443 to demonstrate a stateless rule, which matches on the packet’s 5-tuple (protocol, ports, addresses) rather than application content.

Introduce the change. Invert the two rule priorities so the drop rule evaluates before the allow rule. This is the kind of change a rushed rule edit can introduce.

To invert the stateless rule priorities

  1. Go to the Amazon VPC console.
  2. In the navigation pane, under Network Firewall, choose Network Firewall rule groups.
  3. Choose the rg-stateless-priority rule group to open its details page.
  4. In the Rules section, choose Edit.
  5. Raise the (Action: Pass) rule’s priority number so it sits after the (Action: Drop) rule, then choose Save. For example, change the (Action: Pass) rule from 100 to 300 (any number higher than the drop rule’s 200 works). You only need to move one rule, and using 300 avoids a clash with the drop rule that already sits at 200. Network Firewall evaluates the lowest priority number first, so the (Action: Drop) rule at 200 now wins for this traffic class, ahead of the (Action: Pass) rule at 300.
Figure 7: Scenario 2 – Rule priority change blocking the traffic class

Figure 7: Scenario 2 – Rule priority change blocking the traffic class

What happens. The drop rule now wins, the TCP connection to the test endpoint on port 9142 times out, the StatelessRuleFailures metric climbs above baseline, and Alarm-2 moves to ALARM. The Scenario 2 card reads DROPPED (with the condition Stateless rules dropping the monitored traffic class), while the Scenario 1 and Scenario 3 cards stay Healthy (Figure 8). On the topology, the alarm pipeline from CloudWatch through Amazon SNS and Lambda to DevOps Agent and the workload-to-firewall inspect lines both turn amber, which the legend defines as collateral / alarm active, because the packets are now dropped at the firewall. To demonstrate the resulting failure, the TLS :9142 line from the internet gateway to the test endpoint is shown in red, which the legend defines as dropped (root cause).

Figure 8: Scenario 2 active

Figure 8: Scenario 2 active

Let DevOps Agent investigate. A stateless drop happens before traffic reaches the stateful inspection engine, so it produces no ALERT log entries. The agent turns to configuration and flow logs instead:

  1. Reads the stateless rule group state and finds the drop rule at the lower priority number, ahead of the pass rule, so the drop evaluates first.
  2. Reads the flow logs and sees passed packets drop to zero within a minute of the change.
  3. Searches CloudTrail and surfaces the UpdateRuleGroup call that inverted the priorities, identifying the user, role, and timestamp about a minute before the alarm.
  4. Reports the root cause as that priority inversion. Recommends removing the redundant drop rule and managing the rule group through infrastructure-as-code (IaC) to prevent manual misconfigurations.
  5. Presents this as a plan you review and apply, not an automatic change.

In the DevOps Agent Operator Web App view, the agent first restates the Alarm-2 trigger and confirms that a workload connectivity health check is failing because the firewall’s stateless rules are dropping egress (Figure 9).

Figure 9: Scenario 2 – The symptom

Figure 9: Scenario 2 – The symptom

Next, the agent identifies the root cause, using the rule-group state and CloudTrail to pinpoint the conflicting DROP/PASS rules, where the new DROP rule’s lower priority number makes it match first (Figure 10).

Figure 10: Scenario 2 – The root cause

Figure 10: Scenario 2 – The root cause

Finally, the agent presents a mitigation plan, recommending you remove the conflicting DROP rule at priority 200 to restore traffic flow (Figure 11).

Figure 11: Scenario 2 – The mitigation plan

Figure 11: Scenario 2 – The mitigation plan

Confirm recovery. Apply the change the agent recommends. After the allow rule is ahead of the drop rule again, Alarm-2 clears and the card returns to green. Move on to Scenario 3.

Scenario 3. Asymmetric cross Availability Zone routing drop

At baseline, the protected subnet in each Availability Zone routes its egress through the firewall endpoint in that same Availability Zone , and the matching return route uses that same endpoint. One endpoint sees both directions of the flow, so the stateful engine completes the handshake. The workload runs in the protected subnet in us-east-1a (CIDR 10.0.4.0/24), so at baseline its egress and its return both use the us-east-1a firewall endpoint.

Introduce the change. Make the flow asymmetric by sending egress out one Availability Zone endpoint while the return comes back through the other. This takes two route edits, and both are required. With only the first edit the flow can still complete, so the alarm will not trigger until both are saved. It makes no firewall-policy change, mirroring a real multi-Availability-Zone routing mistake.

To create asymmetric cross Availability Zone routing

  1. Go to the Amazon VPC console and choose Route tables in the navigation pane.
  2. Flip the egress. Select the NfNetworkStack/SampleVpc/protectedSubnet1 route table (the us-east-1a protected subnet, where the workload runs). On the Routes tab, choose Edit routes. Its 0.0.0.0/0 route currently targets the us-east-1a firewall endpoint. For the target, choose Gateway Load Balancer Endpoint and select the us-east-1b firewall endpoint, then choose Save changes.
  3. Move the return. Select the NfNetworkStack/SampleVpc/publicSubnet2 route table (the us-east-1b public subnet, where egress now exits). Choose Edit routes, then Add route. For the destination enter the workload CIDR 10.0.4.0/24. For the target, choose Gateway Load Balancer Endpoint and select the us-east-1a firewall endpoint. Choose Save changes.

After both edits, a flow’s egress leaves through the us-east-1b endpoint while its return is directed to the us-east-1a endpoint. Neither endpoint sees the whole flow.

Figure 12: Scenario 3 routing change breaking the flow’s symmetry

Figure 12: Scenario 3 routing change breaking the flow’s symmetry

What happens. A new connection leaves through one endpoint. Its return arrives at the other endpoint, which never saw the connection open, so the handshake fails. Unlike Scenarios 1 and 2, this affects the whole subnet, so all egress stops and Alarm-2 and Alarm-3 both move to ALARM. The AWS/NetworkFirewall DroppedPackets alarm (Alarm-1) stays quiet because no endpoint is making a drop decision. The flow is lost to asymmetric routing rather than counted as a firewall drop. This is why monitoring application connectivity matters. A routing fault is invisible to the firewall’s own drop counter. On the status page, the Scenario 2 card reads DROPPED (with the condition “Stateless rules dropping the monitored traffic class”) and the Scenario 3 card reads DROPPED (with the condition Return traffic dropped by asymmetric cross-Availability-Zone routing), while the Scenario 1 card stays Healthy (Figure 13). On the topology, the alarm pipeline from CloudWatch through Amazon SNS and Lambda to DevOps Agent and the workload-to-firewall inspect lines both turn amber, which the legend defines as collateral / alarm active, while the egress path from the firewall through the NAT gateway and the TLS :9142 and HTTPS · routing lines to the test endpoint turn red, which the legend defines as dropped (root cause).

Figure 13: Scenario 3 – The status page during a path-wide outage

Figure 13: Scenario 3 – The status page during a path-wide outage

Let DevOps Agent investigate. Both Alarm-2 and Alarm-3 fire in the same datapoint. DevOps Agent recognizes them as linked and merges them into a single investigation:

  1. Reads the flow logs and sees bidirectional TLS connections stop abruptly, with only one-way traffic remaining and no flows reaching the established state.
  2. Reads the firewall metrics and sees received and passed packets shift from one Availability Zone to the other at the moment of the change.
  3. Calls DescribeRouteTables and finds the egress route pointing at one Availability Zone firewall endpoint while the return route points at the other.
  4. Searches CloudTrail and surfaces the ReplaceRoute and CreateRoute calls by the same user, about a minute before both alarms fired.
  5. Reports the root cause as that asymmetric routing change. Recommends restoring symmetric same-Availability-Zone routing so egress and return traverse the same endpoint.
  6. Presents this as a plan you review and apply, not an automatic change.

A mitigation plan is a recommendation you review, not an automatic change, and the right fix depends on the intended design. Restoring symmetric routing can mean sending the workload subnet’s egress back through its own-Availability-Zone firewall endpoint (this sample’s architecture) or, in a design that doesn’t inspect this path, back through a NAT gateway. The agent infers a plausible target from what it can observe, so review the specific route it proposes against your intended topology before you apply it. (Connecting your pipeline or infrastructure-as-code, covered in the next section, lets the agent recommend the target that matches your design.)

In the DevOps Agent Operator Web App view, the agent restates the Alarm-3 (AsymmetricFlowFailures) trigger and confirms the workload’s egress to a monitored endpoint is being blocked by the Network Firewall (Figure 14).

Figure 14: Scenario 3 – The symptom

Figure 14: Scenario 3 – The symptom

Next, the agent identifies the root cause: manual route table changes that created cross-AZ asymmetric routing through the network firewall, breaking its symmetric routing requirement (Figure 15)

Figure 15: Scenario 3 – The root cause

Figure 15: Scenario 3 – The root cause

Finally, the agent presents a mitigation plan, recommending you restore symmetric routing by pointing protectedSubnet1‘s default route back to the same Availability Zone firewall endpoint, so one endpoint sees both directions of the flow again (Figure 16).

Figure 16: Scenario 3 – The mitigation plan

Figure 16: Scenario 3 – The mitigation plan

Confirm recovery. Apply the change the agent recommends, after checking the route target matches your intended design. After the workload subnet’s egress and return use the same Availability Zone firewall endpoint again, the control probe recovers, the alarms clear, and every card returns to green.

Further considerations

In production a single change can trigger several alarms at the same time, as Scenario 3 shows. DevOps Agent links related investigations and works them as one, so you review a single root cause. You can validate the linked findings or unlink an alarm to investigate it independently. If you would rather collapse alarms before they reach the agent, you can add correlation logic in the bridge Lambda function, buffering and grouping by firewall. You can also add email, Amazon Simple Queue Service (Amazon SQS), or HTTP subscribers to the SNS topic, or add the webhook Lambda function to a topic you already run. DevOps Agent produces a mitigation plan but does not change your environment on its own.

You can also give the agent more to work with. DevOps Agent connects to source repositories and CI/CD pipelines, integrating with GitHub (including GitHub Enterprise Server and GitLab Self-Managed through a private connection). It can associate AWS resources with deployments of AWS CloudFormation, AWS CDK, Amazon Elastic Container Registry (Amazon ECR) images, and Terraform. With deployed configuration and recent deployment events in view, the agent correlates the disruption against the change that introduced it and recommends a fix matching your intended design. For this sample, that means recommending the workload subnet’s own Availability Zone firewall endpoint rather than a generic symmetric path.

DevOps Agent also supports proactive incident prevention. It analyzes patterns across past investigations and delivers recommendations to prevent similar issues from recurring, including governance recommendations that strengthen deployment processes and pipeline controls. For Network Firewall rule changes, this means the agent can recommend guardrails for your CI/CD pipeline based on the classes of misconfigurations it has already resolved. You can access these recommendations through the Improvements page in the DevOps Agent Operator Web App.

Clean up

Clean up the environment with one command.

bash scripts/destroy.sh

It reverts any active scenario, runs cdk destroy for all stacks, and sweeps for stragglers by the Project = nf-devops-agent tag. The main cost drivers are the two Network Firewall endpoints, the NAT gateways (one in the main VPC for each Availability Zone, one in the test-endpoint VPC), and the test endpoint’s load balancers. Each of these bills at an hourly rate for as long as it’s provisioned, whether or not traffic is flowing, so a stack left running continues to accrue charges around the clock even while idle. Running the scenarios and tearing the stack down the same day limits the cost to a few active hours rather than days of idle hourly charges.

Conclusion

In this post, we showed you how AWS DevOps Agent accelerates troubleshooting for three common network firewall connectivity issues. The first was a domain deny list. The second was a stateless priority inversion. The third was an asymmetric cross-AZ routing drop. For each one, DevOps Agent investigated the drop and returned a root cause with a mitigation plan you approve before applying. The first scenario triggered on a prebuilt Network Firewall metric, and the other two on application health metrics. That shows both ways to alarm on a firewall problem through one pipeline.

The pattern isn’t specific to Network Firewall. The same flow fits any service that emits CloudWatch metrics and logs, such as AWS WAF, security groups, and network ACLs. Clone the sample repository to explore the solution, then apply what you learn to your own firewall, application, and alarms. For more details, see the AWS Network Firewall Developer Guide and the AWS Network Firewall pricing page. Start with the Getting Started with AWS DevOps Agent guide to connect your first webhook.

Salman Ahmed

Salman is a Senior Technical Account Manager at AWS, specializing in helping customers design, implement, and optimize their AWS environments. He combines deep networking expertise with a passion for exploring emerging technologies to help organizations get the most out of their cloud investments. Outside of work, he enjoys photography, traveling, and watching his favorite sports teams.

Automate creating AWS Glue Data Catalog views with AWS SDK for data mesh use case

Post Syndicated from Aarthi Srinivasan original https://aws.amazon.com/blogs/big-data/automate-creating-aws-glue-data-catalog-views-with-aws-sdk-for-data-mesh-use-case/

AWS Glue Data Catalog view is a multi-dialect view that supports querying from multiple SQL query engines, such as Amazon Athena, Amazon Redshift Spectrum, Apache Spark in Amazon EMR and AWS Glue. You can create a Data Catalog view in one account, using an AWS Identity and Access Management (IAM) definer role in the same or different account and use AWS Lake Formation to share the view across multiple accounts. The definer role has the required full SELECT on the base tables to create the view and share it with other users for querying. The Data Catalog assumes the definer role and manages access of the base tables when the view is queried, thus allowing to share a subset of data without sharing the underlying base tables.

AWS Glue now adds AWS SDK support for creating and updating the ATHENA dialect of Glue views. With this addition, you can now create ATHENA and SPARK dialects of Glue views simultaneously, using a cross account IAM definer role. This feature enhances the automation to create and update Glue views, like that of Data Catalog tables. In our earlier blog Create AWS Glue Data Catalog views using cross-account definer roles, we had introduced IAM definer roles in a cross-account use case to create Data Catalog views with SPARK dialects using the APIs – CreateTable() and UpdateTable() – while creating and adding ATHENA dialects using Athena query editor. As a continuation to it, this post shows you how to use the Catalog objects API CreateTable() to programmatically create ATHENA and SPARK dialects using cross-account IAM definer roles, and how to add the ATHENA dialect programmatically for the views that were created earlier with only SPARK dialect.

Cross account definer roles enable enterprise data mesh architectures where multiple accounts are interconnected in a central governance and multiple producers and consumers. The central governance account hosts the database, tables and permissions, while the producer accounts maintain CI/CD pipelines to create and manage those data assets. Having the definer role in producer accounts allows those CI/CD pipelines to be fully managed by IAM roles in the individual accounts.

Key points on creating multi-dialect views using cross-account definer roles

  • ATHENA dialects are validated and asynchronously created. Hence, a cross-account Glue connection is required for validation for every producer account-central governance account pair. This is a one-time setup.
  • SPARK dialects are not validated. Hence SPARK dialect’s create syntax requires SubObjects list of the base tables and StorageDescriptor fields for the columns of the view.
  • Though queries on cross account views can be run using database resource link names, the view definition SQL query for creating the view requires the original database and base table names from the central governance account.
  • If a view has SPARK and ATHENA dialects available, we recommend updating both the dialects of the view simultaneously using update_table() API/SDK, for any changes in the SQL definition of the view or the base table. This will keep both the dialects queryable.
  • Creating and updating both SPARK and ATHENA dialects using cross account definer role is supported using AWS CloudFormation.
  • The Data Catalog view that can be created using cross account IAM definer roles are available in SPARK and ATHENA dialects and currently not supported for Redshift Spectrum dialect.

Prerequisites

We use the same setup used in Create AWS Glue Data Catalog views using cross-account definer roles for the sample database, tables, definer role, resource link, IAM and Lake Formation permissions on those resources and principals between the two AWS accounts. Summarizing the requirements as below.

  • The setup includes a central governance account with Data Catalog database bankdata_icebergdb and two tables transaction_table1 and transaction_table2, a producer account with a Data-Analyst role used as view definer role.
  • Lake Formation permissions on the central account’s database and tables are granted to the producer account Data-Analyst role as per the earlier blog. The definer role in producer account should have database DESCRIBE and CREATE_TABLE permissions, table SELECT and DESCRIBE permission on all columns and rows of the base tables. The IAM permissions required on the definer role are detailed in Prerequisites for creating views. Similarly, follow the earlier blog to create resource link for the shared database and grant Lake Formation permissions on the resource link to the Data-Analyst
  • An Athena data source named centraladmin in the producer account, pointing to the Data Catalog of the central governance account.

Creating ATHENA and SPARK dialects at the same time

Creating both ATHENA and SPARK dialects of a Glue catalog view simultaneously is now supported by the AWS SDK. In the producer account, create a new Glue connection, required for the Athena dialect validation. This is a prerequisite for creating the ATHENA dialect of the Glue catalog view using cross account definer role. Then we create a Glue view with both dialects.

  1. Sign in to the producer account as the Lake Formation admin role, or any role with permission to create AWS Glue connections.
  2. Using an AWS Command Line Interface (AWS CLI) environment, such as AWS CloudShell, create an AWS Glue connection as follows.
    aws glue create-connection --cli-input-json file://athena-validation-connection.json

    The content of athena-validation-connection.json is as follows.

    {
        "CatalogId": "<producer-account-id>",
        "ConnectionInput": {
            "Name": "glue-view-validation-connection",
            "Description": "Glue view Athena cross-account validation connection",
            "ConnectionType": "VIEW_VALIDATION_ATHENA",
            "ConnectionProperties": {
                "WORKGROUP_NAME": "primary",
                "DATA_SOURCE": "centraladmin"
            }
        }
    }

    Note: If you are using Athena for the first time in your account or using Primary workgroup, setup the query results location bucket using Specify a query result location.

  3. Sign out as the Lake Formation admin and sign back in to the producer account as the definer IAM role, Data-Analyst.
  4. Create an AWS Glue view using the create-table CLI command and JSON file, or using the AWS SDK for Python (Boto3) script.
    aws glue create-table --cli-input-json file://create_multipledialects.json

    The content of create_multipledialects.json is as follows.

     {
       "DatabaseName": "rl_bank_iceberg",
       "TableInput": {
         "Name": "view_2dialects_2basetables_fromcli",
         "StorageDescriptor": {
           "Columns": [
             {
               "Name": "transaction_id",
               "Type": "string"
             },
             {
               "Name": "transaction_type",
               "Type": "string"
             },
             {
               "Name": "transaction_amount",
               "Type": "double"
             },
             {
               "Name": "transaction_location",
               "Type": "string"
             },
             {
               "Name": "transaction_date",
               "Type": "date"
             }
         },
         "ViewDefinition": {
           "SubObjects": [
             "arn:aws:glue:us-west-2:<central-account-id>:table/bankdata_icebergdb/transaction_table1",
             "arn:aws:glue:us-west-2:<central-account-id>:table/bankdata_icebergdb/transaction_table2"
            ],
           "IsProtected": true,
           "Representations": [
             {
               "Dialect": "SPARK",
               "DialectVersion": "1.0",
               "ViewOriginalText": "SELECT a.transaction_id, a.transaction_type, a.transaction_amount, b.transaction_location, b.transaction_date FROM bankdata_icebergdb.transaction_table1 a RIGHT JOIN bankdata_icebergdb.transaction_table2 b ON a.transaction_id = b.transaction_id",
               "ViewExpandedText": "SELECT a.transaction_id, a.transaction_type, a.transaction_amount, b.transaction_location, b.transaction_date FROM bankdata_icebergdb.transaction_table1 a RIGHT JOIN bankdata_icebergdb.transaction_table2 b ON a.transaction_id = b.transaction_id"
             },
             {
                "Dialect": "ATHENA",
                "DialectVersion": "3",
                "ViewOriginalText": "SELECT a.transaction_id, a.transaction_type, a.transaction_amount, b.transaction_location, b.transaction_date FROM bankdata_icebergdb.transaction_table1 a RIGHT JOIN bankdata_icebergdb.transaction_table2 b ON a.transaction_id = b.transaction_id",
                "ValidationConnection": "glue-view-validation-connection"
             }
           ]
         }
       }
    }

    Notes about fields in the above CLI input JSON (applies to all SDK):

    • The definer is by default the API caller, but a Definer field can be set to explicitly specify a different IAM role.
    • In the ViewDefinition, database qualifiers are required for SPARK dialect. That is, the SQL definition provided for ViewOriginalText and ViewExpandedText should be in <source_database_name>.<source_table_name> format.
  5. After the view is created, you can inspect the details on the Lake Formation console. The SQL definitions show both ATHENA and SPARK as shown in the following screenshot.

Lake Formation console showing the SQL definitions tab for the new Data Catalog view, with both ATHENA and SPARK dialects listed

If your view creation fails for any of the dialects, you can use the AWS Glue get-table CLI command with --include-status-details to see what the error is and rectify it.

aws glue get-table --database-name <rl_database_name> --name <view_name> --include-status-details

Glue PySpark script

The PySpark script for creating a view with ATHENA and SPARK dialects are provided below. Download and edit the Pyspark script with your bucket name, producer and central account ids, region and relevant Glue resource names: bdb_5773_createview_bothdialects.py

Provide the following settings to run the script in your Glue Studio. For details on running a Spark job in Glue, refer Working with Spark jobs in AWS Glue.

  • Choose Data-Analyst as the job execution IAM role.
  • Choose Glue 5.1 for Glue version.
  • For the Requested number of workers, provide >=4. This is an FGAC Spark driver requirement, which is needed for Glue catalog views. Below screenshot shows these settings.
  • Add the following 2 properties as additional job parameters. A screenshot is shown for reference.
    --datalake-formats = iceberg
    --enable-lakeformation-fine-grained-access=true

    AWS Glue ETL job configuration page showing the additional job parameters set for the multi-dialect view creation script

  • Save and run the Glue job. Check the stdout logs to review the query on the newly created view.

A sample update_table script is also provided below, to illustrate changing the view definition with additional columns. Note the REPLACE keyword:

bdb_5773_updateview_bothdialects.py

Adding ATHENA dialect using SDK to an existing AWS Glue view

You can update an existing AWS Glue view that was created with the SPARK dialect and add the ATHENA dialect using the SDK. The following example uses the update-table CLI command.

aws glue update-table --cli-input-json file://add-athena-dialect.json

The content of add-athena-dialect.json is as follows.

{
    "DatabaseName": "rl_bank_iceberg",
    "ViewUpdateAction": "ADD",
    "TableInput": {
        "Name": "view_sparkfirst_athenanext",
        "ViewDefinition": {
            "Representations": [
                {
                    "Dialect": "ATHENA",
                    "DialectVersion": "3",
                    "ViewOriginalText": "SELECT a.transaction_id, a.transaction_type, a.transaction_amount, b.transaction_location, b.transaction_date FROM bankdata_icebergdb.transaction_table1 a RIGHT JOIN bankdata_icebergdb.transaction_table2 b ON a.transaction_id = b.transaction_id",
                    "ValidationConnection": "glue-view-validation-connection"
                }
            ]
        }
    }
}

Verify the added dialect on the view by reviewing the SQL definitions of the view in Lake Formation console or using GetTable(). If you want to edit the SQL definition or change the base tables of an existing view that has both SPARK and ATHENA dialects, you can do so using the update_table API (using SDK or CLI), with "ViewUpdateAction": “REPLACE” and provide both the dialect definition under ViewDefinition.

You can run queries on the view from the producer account as Data-Analyst. The view can be shared using Lake Formation Tags or named method, just like sharing tables, to additional consumer accounts from the central governance account. The consumer accounts will create a resource link and query the views.

Cleanup

To avoid incurring ongoing costs, clean up the resources you used for this post:

  1. Revoke the Lake Formation permissions granted to the Data-Analyst role and the producer account from the central governance account.
  2. Drop the Data Catalog tables, views, and the database.
  3. Delete the Athena query results from your Amazon Simple Storage Service (Amazon S3) bucket.
  4. Delete the Data-Analyst role from IAM.
  5. Delete the AWS Glue connection and the Athena data source.
  6. Delete the AWS Glue job, if you tried the Python script as an AWS Glue job.

Conclusion

In this post, I demonstrated how to use cross-account IAM definer roles with AWS Glue Data Catalog views, how to create and update ATHENA and SPARK dialects using the Data Catalog CreateTable() and UpdateTable() APIs. The multi-dialect Data Catalog views allow sharing a subset of data from different tables using Lake Formation permissions, including LF-Tags based access control. The cross-account definer roles support multi-account data mesh architectures so that the producer IAM roles can run the CI/CD pipelines in its account. We encourage you to try the feature and share your feedback in the comments.

Acknowledgements: I would like to thank all the team members who worked to add AWS SDK support for creating ATHENA and SPARK dialects together for AWS Glue views – Daniil Arushanov, Wyatt Hawes, Yuxi Wu, Santhosh Padmanabhan and Karthik Devaraj.


About the author

Aarthi Srinivasan

Aarthi Srinivasan

Aarthi is a Senior Big Data Architect working on data, analytics and GenAI topics with the worldwide Specialists Org at AWS. She works with AWS customers and partners to architect data lake solutions, enhance product features, and establish best practices for data governance and analytics services adoption.

Efficient log management with Amazon OpenSearch Service data streams

Post Syndicated from Praveen Krishnamoorthy Ravikumar original https://aws.amazon.com/blogs/big-data/efficient-log-management-with-amazon-opensearch-service-data-streams/

Time series data workloads in Amazon OpenSearch Service can present unique challenges for organizations, especially when dealing with continuously growing datasets. Many customers struggle with heavily loaded single indices that lead to high query latency, degraded performance, and unnecessary costs. In this post, we show you how to implement data streams with Index State Management (ISM) in Amazon OpenSearch Service. This approach automatically manages your time series data lifecycle and optimizes both performance and costs. Data streams distribute incoming data across multiple backing indices, helping to reduce single-index bottlenecks, while ISM policies automate rollover, retention, and storage tiering to help manage costs.

The challenge

While Amazon OpenSearch Service has long provided tools like Index State Management (ISM) for time series data management, many organizations still struggle with implementing optimal patterns for their continuously growing datasets. Common challenges include:

  • Performance degradation from index growth: As single indices grow unbounded, query latency increases, you might find shard sizes more difficult to manage, and you might experience strain on your cluster resources.
  • Manual index management overhead: Without automation, you must invest significant operational effort to manage index lifecycles, rollover, and retention.
  • Complex setup: Coordinating index templates, aliases, and ISM policies manually can be error-prone.
  • Inefficient resource utilization: All data residing in hot storage regardless of access patterns, leading to unnecessarily high costs.

Solution overview

As illustrated in Figure 1, our solution uses Amazon OpenSearch Service data streams combined with Index State Management to automatically distribute data across multiple indices and manage the data lifecycle. A data stream is an abstraction layer that simplifies time series data ingestion. It provides a single, consistent endpoint for writes while automatically managing multiple backing indices behind the scenes. Instead of writing directly to individual indices, applications write to the data stream, which routes data to the appropriate backing index.

Here’s how it works:

  • Data streams provide a single write index when data is first ingested, which can help streamline time series data ingestion.
  • When the backing index ages or grows to meet your defined criteria, ISM automatically performs the rollover operation.
  • You can use ISM policies to automatically transition your aged data to different storage tiers based on rules you define and configure.
  • You can automate the entire process through rules you define in the index template and ISM policies.

Architecture diagram showing time series data flowing into an Amazon OpenSearch Service data stream, which routes writes to multiple backing indices while ISM transitions aged indices from hot to UltraWarm storage

Figure 1 — Time series data workflow using Amazon OpenSearch Service data streams

Data ingests into an index according to your index template configuration. Over time, a data stream creates new indices automatically. Amazon OpenSearch Service manages the lifecycle, transitioning data from hot to warm storage according to the ISM policy configuration you define.

When to use this solution

This approach is ideal when:

Implementation steps

Prerequisites

Before you begin, make sure that you have the following:

The following steps walk through implementing a time series data solution, using a web server logs database as an example. You can run these steps using any of the following:

For this post, we use the Dev Tools Console in OpenSearch Dashboards. To access it:

  • Log in to OpenSearch Dashboards.
  • Navigate to Dev Tools (usually found on the menu under Management).
  • Use the interactive console to run the commands.

Section 1: Create data stream

  1. Create an ISM policy

First, create an ISM policy that defines the rules for index rollover and storage tier transitions. The following policy defines two states (hot and warm) and sets rules for when indices transition between them. The policy triggers a rollover when the document count reaches 1,000 and moves indices to warm storage after 2 minutes.

Note: The rollover and transitions configurations are only for demo purposes.

PUT _plugins/_ism/policies/ds-ism-policy
{
"policy": {
"description": "rollover policy when index is large",
"default_state": "hot",
"ism_template": [
{
"index_patterns": ["webserver-logs-data-stream*"],
"priority": 300
}
],
"states": [
{
"name": "hot",
"actions": [
{
"rollover": {
"min_doc_count": 1000
}
}
],
"transitions": [
{
"state_name": "warm",
"conditions": {
"min_index_age": "2m"
}
}
]
},
{
"name": "warm",
"actions": [
{
"retry": {
"count": 3,
"backoff": "exponential",
"delay": "1m"
},
"warm_migration": {}
}
]
}
]
}
}
  1. Create an index template for the data stream

An index template matches indices by a regex pattern and applies predefined settings and schema at index creation. The following template maps the required timestamp field for the data stream and associates matching indices with the ISM policy we created.

PUT _index_template/webserver-logs-data-stream-template
{
"index_patterns": ["webserver-logs-data-stream*"],
"data_stream": {
"timestamp_field": {
"name": "timestamp"
}
},
"template": {
"settings": {
"plugins.index_state_management.policy_id": "ds-ism-policy"
},
"mappings": {
"properties": {
"timestamp": {
"type": "date"
}
}
}
}
}
  1. Create data stream

In this step, you create the data stream that handles the time series data. The data stream provides a single, unified write target for ingesting data while managing multiple backing indices behind the scenes.

PUT _data_stream/webserver-logs-data-stream
  1. Validate ISM policy mapping

Verify that the ISM policy is correctly associated with the data stream that was created in the earlier step. This validation step confirms that automatic lifecycle management works as expected.

GET _plugins/_ism/explain/webserver-logs-data-stream

Output

Confirm that policy_id matches the ISM policy name you created earlier and that enabled is set to true.

OpenSearch explain output showing the ISM policy_id mapped to the data stream with enabled set to true

Section 2: Ingesting data to data stream

In real-world scenarios, log data is typically collected and streamed directly to Amazon OpenSearch Service data streams. However, to demonstrate rollover and migration scenarios in this post, we take a different approach. We first load sample log data into a standard OpenSearch Service index, then reindex and migrate that data to a data stream.

To get started, run the commands in the Dev Tools console to create an index and populate it with sample log data.

  1. Reindex existing data

This step shows how to migrate existing data in a traditional index to the new data stream. The reindex operation includes a script that confirms each document has a valid timestamp field (timestamp). Documents without a timestamp field are skipped by the reindex operation to maintain data integrity.

POST _reindex
{
"source": {
"index": "webserver-logs"
},
"dest": {
"index": "webserver-logs-data-stream",
"op_type": "create"
},
"script": {
"source": """
try {
// Validate timestamp field exists and has a value
if (ctx._source.timestamp == null || ctx._source.timestamp.empty) {
ctx.op = 'noop';
}
} catch (Exception e) {
// Skip this document on any error
ctx.op = 'noop';
}
"""
},
"conflicts": "proceed"
}
  1. Monitor index rollover

As shown in Figure 2, after reindexing, monitor the creation of backing indices. The ISM policy evaluates indices every 5 minutes by default. Rollover occurs based on the defined conditions (1,000 documents or 2 minutes of age). Verify this using the following command.

GET _cat/indices/.ds-*?v&h=index,status,health,pri,rep,docs.count,store.size,creation.date&s=index

The output should look like the following:

_cat/indices output listing the .ds backing indices with their status, health, and document counts

You can also validate this from OpenSearch Dashboards by navigating to Index Management, Data streams, webserver-logs-data-stream.

OpenSearch Dashboards Index Management page showing the webserver-logs-data-stream and its backing indices

*Figure 2 — Combined view of the _cat/indices CLI output and the OpenSearch Dashboards data stream details, showing a successful index rollover across multiple backing indices*

  1. Validate warm transition

Verify that indices are correctly transitioning from hot to warm storage based on the ISM policy conditions. You can monitor this through OpenSearch Dashboards or API queries in the Dev Tools console.

GET _plugins/_ism/explain/webserver-logs-data-stream

OpenSearch explain output showing backing indices transitioning from the hot state to the warm state

  1. Verify ingested data

Run a search against the data stream to confirm your documents were successfully indexed:

GET webserver-logs-data-stream/_search
{
"size": 1
}

Output

You should see your ingested documents returned in the hits.hits array, with the timestamp field and other fields you defined in the index template. A non-zero hits.total.value confirms data is flowing correctly through the data stream.

Search results showing an ingested document with the @timestamp field and a non-zero hits.total.value

  1. Clean up

If needed, these commands remove the data stream and its template.

DELETE _data_stream/webserver-logs-data-stream
DELETE _index_template/ webserver-logs-data-stream-template

Delete the sample data.

Conclusion

OpenSearch data streams with ISM offer capabilities for managing time series data at scale. Organizations that implement this approach can see improved query performance through distributed load and smaller, time-based backing indices that support efficient time-range queries. Automated index management reduces operational overhead. Storage tiering automatically moves aged data to UltraWarm storage, which significantly lowers costs without sacrificing access to historical data. Combined with better scalability for growing datasets, this solution simplifies index management while delivering improved performance and a more cost-effective, maintainable infrastructure.


About the authors

Praveen Krishnamoorthy Ravikumar

Praveen Krishnamoorthy Ravikumar

Praveen is an Analytics Specialist Solutions Architect at AWS. He helps customers design and implement modern data and analytics platforms that leverage the scalability, flexibility, and innovation of the cloud. He is passionate about solving complex data challenges and enabling organizations to unlock actionable insights from their data.

JP Boreddy

JP Boreddy

JP is a Senior Solutions Architect at Amazon Web Services, based in San Diego, California. He works with ISV customers in the security segment, helping them architect and optimize their workloads on AWS. JP specializes in AI/ML, containers, and cloud infrastructure, with a focus on enabling customers to build scalable, cost-effective solutions. He has been with AWS for over four years.

Aswin Vasudevan

Aswin Vasudevan

Aswin is a Senior Solutions Architect for Security, ISV at AWS. He is a big fan of generative AI and serverless architecture and enjoys collaborating and working with customers to build solutions that drive business value.

Kevin Fallis

Kevin is seasoned leader, architect, and developer with experience across many industry verticals and disciplines such as agriculture, ad tech, financial services, networking, security, telecommunications and of course search technologies. His passion helps others leverage the correct mix of AWS services and open-source solutions to achieve success for their business goals. His after-work activities include family, DIY projects, carpentry, horses, playing drums, and all things music.

Govern Amazon Redshift Data Warehouses Data Across Accounts using Amazon SageMaker Unified Studio

Post Syndicated from Bandana Das original https://aws.amazon.com/blogs/big-data/govern-amazon-redshift-data-across-accounts-with-sagemaker-unified-studio/

Managing data governance across multiple Amazon Redshift clusters in different AWS accounts presents significant challenges. Organizations operating multiple Amazon Redshift clusters across AWS accounts often rely on manual processes for secure data sharing, which increases operational overhead and governance requirements. In this post, we show you how to use Amazon SageMaker Unified Studio to implement cross-account data sharing in Amazon Redshift using data mesh principles. We demonstrate how to build a scalable data mesh architecture that supports secure, auditable data sharing across AWS accounts while reducing operational burden.

Amazon SageMaker Unified Studio as the backbone of our data mesh

Amazon SageMaker Unified Studio is a data and AI development service which brings together functionality and tools from existing AWS Analytics and AI and machine learning (ML) services, including Amazon EMR, AWS Glue, Amazon Athena, Amazon Redshift, Amazon Bedrock and Amazon SageMaker AI. With the service, organizations can catalog, discover, share, and govern data stored across Amazon Web Services (AWS) without relying on manual coordination between AWS accounts.

A data mesh is an architectural approach that treats data as a product, with decentralized ownership by data producers while maintaining centralized governance. This architecture separates source systems, data producers (data publishers), data consumers (data subscribers), and central governance. The solution we present is tailored for cross-AWS account usage, creating a foundation for data governance so you can share data across Amazon Redshift clusters in different AWS accounts.

Our proposed solution addresses the following common challenges that organizations face when sharing data across AWS accounts:

  • Manual, ad-hoc data sharing processes are replaced with automated, event-driven data publishing to the SageMaker Unified Studio catalog.
  • Inconsistent governance across different use cases is resolved through a consistent governance framework with proper access controls.
  • High load on producer Amazon Redshift clusters is reduced through decoupled publishing that lowers the operational burden on data producers.
  • Complex credential management is simplified using AWS Secrets Manager and AWS KMS encryption.
  • Lack of auditable data publishing is addressed with full traceability of access and permissions supported by the SageMaker Unified Studio service.

With this approach, you can help reduce the time and effort required for cross-account data sharing while maintaining security and governance standards.

Architectural overview

The architecture spans three AWS accounts, each with a distinct role in the data mesh:

Central Data Governance Account (Account A) hosts the Amazon SageMaker Unified Studio domain, which serves as the unified catalog and governance layer for data discovery, access control, and subscription management across accounts.

Data Producer (Account B) hosts the source of data and processing workflows. Raw data lands in an Amazon Simple Storage Service (Amazon S3) source bucket and is processed through AWS Glue extract, transform, and load (ETL) jobs or Amazon Redshift auto copy into the Amazon Redshift source database. Amazon Redshift credentials are securely stored in AWS Secrets Manager.

Data Consumer (Account C) hosts the target Amazon Redshift database and analytics workflows. After access is granted, consumers can query shared data and connect downstream visualization tools.

While this diagram shows a single producer and consumer for simplicity, in a real-world deployment there might be hundreds of producer and consumer accounts connecting through the central governance layer. Amazon SageMaker Unified Studio scales to support this by providing a single place for managing data products regardless of the number of participating accounts.

The data sharing workflow is driven by Amazon SageMaker Unified Studio. The data owner publishes data to the catalog, where it becomes discoverable by consumers across accounts. Consumers browse the catalog, subscribe to data products, and the data owner approves the request. After approval, Amazon SageMaker Unified Studio handles the cross-account sharing, granting the consumer access without requiring direct connectivity between producer and consumer Amazon Redshift clusters.

Publishing Amazon Redshift data assets to the data mesh

In a data mesh architecture, data producers need to make their data products discoverable and accessible across the organization. Amazon SageMaker Unified Studio provides a centralized catalog where data assets can be published for consumer subscription.

In practice, this means registering your data sources with the catalog so they can be discovered, governed, and subscribed to by consuming teams. This section walks through the steps required to register Amazon Redshift data sources with SageMaker Unified Studio.

Before you can publish data assets from your producer account, you need to complete several configuration steps across your Amazon Redshift cluster, AWS Secrets Manager, and Amazon SageMaker Unified Studio.

Prerequisites

  • Install the AWS Command Line Interface (AWS CLI) (v2.15+ recommended).
  • Obtain temporary credentials with permissions to administer each account (producer, consumer, and domain account)
  • IAM permissions required: redshift:* on the relevant clusters, secretsmanager:CreateSecret / PutResourcePolicy / TagResource, kms:CreateKey / PutKeyPolicy / TagResource, datazone:* for subscription-target creation, and iam:PassRole for the Amazon Redshift cluster role.
  • Amazon Redshift clusters must use RA3 node types (ra3.xlplus, ra3.4xlarge, or ra3.16xlarge). Data sharing is not supported on other node types.
  • Amazon SageMaker Unified Studio domain must already be created in Account A with the Tooling and LakeHouseCatalog blueprints available.
  • All resources must be in an AWS Region where Amazon SageMaker Unified Studio is available.

Step 1: Account association and blueprint enablement

To implement the data mesh architecture described in the previous section, you need to set up the following accounts and enable the required blueprints. This ensures that the central governance layer can discover and manage data assets across your producer and consumer accounts.

This post uses three separate AWS accounts to illustrate the cross-account data sharing pattern. However, Amazon SageMaker Unified Studio also supports publishing and subscribing to data within a single account or across any number of accounts depending on your organizational setup. Additionally, this walkthrough uses a provisioned Amazon Redshift cluster, but Amazon SageMaker Unified Studio also supports Amazon Redshift Serverless for both publishing and subscribing to data assets.

Step 2: Configure your Amazon Redshift cluster and credentials

  • In the producer account (Account B), the data to be shared resides in an Amazon Redshift cluster.
  • Verify that your Amazon Redshift cluster uses node types from the RA3 family.
  • Add the following tags to your Amazon Redshift cluster.

Amazon Redshift console showing tags added to the cluster

  • Create a superuser in Amazon Redshift for Amazon SageMaker Unified Studio. For the Amazon Redshift cluster, the database user you provide in AWS Secrets Manager must have superuser permissions. With superuser permission, your Amazon Redshift cluster can publish data and subscribe from the data mesh created with Amazon SageMaker Unified Studio, and it manages the subscriptions (access) on your behalf. For reference, see the note section in this QuickStart guide with sample Amazon Redshift data.

Tag key-value pairs configured on the Amazon Redshift cluster

  • Store the user’s credentials in Secrets Manager. Select the credential type, enter the credential values, and choose the AWS Key Management Service (AWS KMS) key with which to encrypt the secret

QuickStart guide note about providing superuser credentials for Amazon SageMaker Unified Studio

Tags on the AWS Secrets Manager secret including the Amazon Redshift cluster ARN

Resource policy added to the AWS Secrets Manager secret for Amazon SageMaker Unified Studio access

  • If your secret is encrypted with a customer managed AWS KMS key, append the key policy with the following statement and add a tag to the key: AmazonDataZoneEnvironment = All. You can skip this step if you’re using an AWS managed KMS key.
{
    "Sid": "AllowSMUSRolesSecretsAccess",
    "Effect": "Allow",
    "Principal": {
        "AWS": "*"
    },
    "Action": [
        "kms:Decrypt",
        "kms:DescribeKey",
        "kms:GenerateDataKey"
    ],
    "Resource": "*",
    "Condition": {
        "StringEquals": {
            "kms:ViaService": "secretsmanager.<<AWS_Region>>.amazonaws.com"
        },
        "StringLike": {
            "aws:PrincipalArn": [
                "arn:aws:iam::<<Data_Producer_Acct_Id(Account B)>>:role/aws-service-role/redshift.amazonaws.com/AWSServiceRoleForRedshift",
                "arn:aws:iam::<<Data_Producer_Acct_Id(Account B)>>:role/<<Redshift_Cluster_IAM_Role_Name>>",
                "arn:aws:iam::<<Data_Producer_Acct_Id(Account B)>>:role/datazone*",
                "arn:aws:iam::<<Data_Producer_Acct_Id(Account B)>>:role/service-role/AmazonSageMaker*"
            ]
        }
    }
},
{
    "Sid": "AllowSMUSRolesCreateGrant",
    "Effect": "Allow",
    "Principal": {
        "AWS": "*"
    },
    "Action": "kms:CreateGrant",
    "Resource": "*",
    "Condition": {
        "Bool": {
            "kms:GrantIsForAWSResource": "true"
        },
        "StringEquals": {
            "kms:ViaService": "secretsmanager.<<AWS_Region>>.amazonaws.com"
        },
        "StringLike": {
            "aws:PrincipalArn": [
                "arn:aws:iam::<<Data_Producer_Acct_Id(Account B)>>:role/aws-service-role/redshift.amazonaws.com/AWSServiceRoleForRedshift",
                "arn:aws:iam::<<Data_Producer_Acct_Id(Account B)>>:role/<<Redshift_Cluster_IAM_Role_Name>>"
            ]
        }
    }
}

Note: Enable automatic rotation. Configure Secrets Manager automatic rotation for this secret with a rotation interval appropriate to your security policy (for example, every 30 days). When implementing rotation, verify that the rotation Lambda function updates the credentials in both Secrets Manager and Amazon Redshift database users simultaneously. Note that Amazon SageMaker Unified Studio retrieves the secret at connection time, so rotation must produce credentials that are valid immediately upon storage: use the alternating-users rotation strategy if you need to avoid downtime during rotation. See the Secrets Manager rotation documentation for setup instructions.

Using Amazon Redshift Serverless?

  • Add the following Tags to the Amazon Redshift Serverless namespace and workgroup.

Tags added to the Amazon Redshift Serverless namespace and workgroup

  • In the Secrets Manager secret, verify the host points to your Serverless endpoint.

AWS Secrets Manager secret showing the host pointing to the Serverless endpoint

  • Add the following tags to the AWS Secrets Manager secret.

Tags added to the AWS Secrets Manager secret for the Redshift Serverless credentials

Publish Amazon Redshift data to the data mesh

With prerequisites complete, you can now register your Amazon Redshift cluster as a data source in Amazon SageMaker Unified Studio.

Step 1: Create an Amazon Redshift type connection

  • Sign in to Account B, navigate to your Amazon SageMaker Unified Studio associated domain, and open the Amazon SageMaker Unified Studio URL.

Amazon SageMaker Unified Studio associated domain sign-in page

Add an Amazon Redshift connection form in Amazon SageMaker Unified Studio

  • The newly created Amazon Redshift connection appears here.

Newly created Amazon Redshift connection listed in Amazon SageMaker Unified Studio

Step 2: Create the data source for your Amazon Redshift data warehouse

Add an Amazon Redshift data source form in Amazon SageMaker Unified Studio

Amazon Redshift data source configuration in Amazon SageMaker Unified Studio

  • For Publishing settings, choose whether assets are immediately discoverable in Amazon SageMaker Catalog.

Publishing settings controlling asset discoverability in the Amazon SageMaker catalog

Using Amazon Redshift Serverless?

When creating the connection and data source, use your workgroupName instead of clusterName. The rest of the data source configuration remains the same.

Step 3: Run the data source and publish the data asset to the data mesh

Data source run configuration in Amazon SageMaker Unified Studio

Data source run results in Amazon SageMaker Unified Studio

  • During creation of data source if you choose Publishing settings such as assets are immediately discoverable, the Amazon Redshift tables and views appear in the catalog as Published, ready for discovery and subscription by data consumers.

Published Amazon Redshift tables and views in the Amazon SageMaker catalog

Data discovery view in the Amazon SageMaker Unified Studio portal

Subscribe Amazon Redshift data through the data mesh

To complete the end-to-end test, you need to set up a consumer Amazon Redshift cluster in Account C.

Step 1: Setting up the consumer cluster

  • Follow the prerequisites from Steps 1 and 2 in the previous section, make sure the cluster and secret are properly tagged as in the following screenshots:
  • Amazon Redshift cluster tags:

Tags applied to the consumer Amazon Redshift cluster

  • Tags for the AWS Secrets Manager secret that stores the user credentials for the Amazon Redshift cluster:

Tags on the AWS Secrets Manager secret storing the consumer cluster credentials

Step 2: Connect the consumer cluster to the data mesh

  • Log into Amazon SageMaker Unified Studio and navigate to your consumer project.
  • In the Compute section of your project, choose Add compute, then choose Connect to existing compute resources.
  • Choose Amazon Redshift Provisioned.
  • Select your consumer Amazon Redshift cluster from the dropdown list and enter the Secrets Manager name.
  • Choose Add compute.
  • Your newly added Amazon Redshift cluster should now show as available.

Consumer Amazon Redshift cluster added as compute in Amazon SageMaker Unified Studio

  • The newly added Amazon Redshift cluster shows an Available state.

Consumer Amazon Redshift cluster showing an Available state

  • In the Data section you can see that objects (table/views) from Amazon Redshift cluster are visible and you can query them.

Data section showing Amazon Redshift tables and views available to query

Step 3: Creating a subscription target

  • Find the tooling environment ID: in your local terminal after obtaining correct credentials as a project member, run this command to find the tooling environment ID.
export REGION='<your-region>'
export SUBSCRIBER_PROJECT_ID='<your-project-id>'
export DOMAIN_ID='dzd-xxxxxxx'

aws datazone list-environments \
  --domain-identifier $DOMAIN_ID \
  --project-identifier $SUBSCRIBER_PROJECT_ID \
  --region $REGION
  • In the response, find and copy the tooling environment ID as shown in the following example.
{
    "items": [
        {
            "projectId": "<PROJECT_ID>",
            "id": "<ENVIRONMENT_ID>",
            "createdBy": "SYSTEM",
            "createdAt": "<TIMESTAMP>",
            "updatedAt": "<TIMESTAMP>",
            "name": "Tooling",
            "awsAccountId": "<AWS_ACCOUNT_ID>",
            "awsAccountRegion": "eu-west-1",
            "provider": "Amazon SageMaker",
            "status": "ACTIVE",
            "environmentConfigurationId": "<ENVIRONMENT_CONFIGURATION_ID>"
        }
    ]
}
  • Locate the Manage Access Role: In Account C, navigate to SageMaker Unified Studio and find the Tooling blueprint. In the Provisioning Tab you will find the Manage Access role and copy the value, as it is needed for the next CLI call.

Provisioning tab showing the Manage Access role for the Tooling blueprint

  • Create the Subscription Target.

With all the information collected, you can create the subscription target for the Amazon Redshift cluster as shown by the CLI call.

export TOOLING_ENV_ID='<tooling-environment-id>'
export AUTHORIZED_PRINCIPAL='datazone_env_<tooling-env-id>'
export MANAGE_ACCESS_ROLE='arn:aws:iam::<account-id>:role/service-role/AmazonSageMakerManageAccess-<domain-id>'

aws datazone create-subscription-target \
  --domain-identifier $DOMAIN_ID \
  --environment-identifier $TOOLING_ENV_ID \
  --name "RedshiftCluster-default-target" \
  --subscription-target-config '[{
    "formName": "RedshiftSubscriptionTargetConfigForm",
    "content": "{\"databaseName\":\"<db-name>\",\"secretManagerArn\":\"arn:aws:secretsmanager:<region>:<account>:secret:<secret-name>\",\"clusterIdentifier\":\"<cluster-id>\",\"schemaName\":\"<schema-name>"}"
  }]' \
  --applicable-asset-types RedshiftViewAssetType RedshiftTableAssetType \
  --manage-access-role $MANAGE_ACCESS_ROLE \
  --provider "Amazon SageMaker" \
  --type RedshiftSubscriptionTargetType \
  --authorized-principals $AUTHORIZED_PRINCIPAL

Using Amazon Redshift Serverless?

Use RedshiftServerlessSubscriptionTargetType as the --type and RedshiftServerlessSubscriptionTargetConfigForm as the formName in the subscription target config. Replace clusterIdentifier with workgroupName in the content JSON.

  • Verify the Subscription Target.

To verify that the subscription target was created successfully, make a last CLI call. You should find in the return a new subscription target with the name RedshiftCluster-default-target.

aws datazone list-subscription-targets \
  --environment-identifier $TOOLING_ENV_ID \
  --domain-identifier $DOMAIN_ID \
  --region $REGION

Step 4: Subscribing to data assets

  • Open the data catalog inside SageMaker Unified Studio and search for the assets you want to subscribe to.

Data catalog search for assets to subscribe to in Amazon SageMaker Unified Studio

Adding multiple databases and schemas

To publish assets from entirely different databases on the same Amazon Redshift cluster, you need to create a separate data source for each database, meaning repeating the steps mentioned in the section before. Each data source points to the same cluster connection but specifies a different database name. This approach gives you independent control over scheduling, publishing settings, and metadata generation for each database’s assets.

On the consumer side, each subscription target is bound to a specific database and schema combination. This is the target location where SageMaker Unified Studio will create views that give the consumer access to subscribed assets. To receive subscribed data in multiple databases or schemas, you create one subscription target per database-schema combination. For example, different teams within the consumer account might want the data materialized in their own schema. The following example shows this pattern:

# Subscription target for the sales schema
aws datazone create-subscription-target \
  --domain-identifier $DOMAIN_ID \
  --environment-identifier $TOOLING_ENV_ID \
  --name "RedshiftCluster-sales-target" \
  --subscription-target-config '[{ "formName": "RedshiftSubscriptionTargetConfigForm", "content": "{\"databaseName\":\"consumer_db\",\"secretManagerArn\":\"arn:aws:secretsmanager:<region>:<account>:secret:<secret-name>\",\"host\":\"<endpoint>\",\"port\":\"5439\",\"schemaName\":\"sales\"}" }]' \
  --applicable-asset-types RedshiftViewAssetType RedshiftTableAssetType \
  --manage-access-role $MANAGE_ACCESS_ROLE \
  --provider "Amazon SageMaker" \
  --type RedshiftSubscriptionTargetType \
  --authorized-principals $AUTHORIZED_PRINCIPAL

# Subscription target for the marketing schema
aws datazone create-subscription-target \
  --domain-identifier $DOMAIN_ID \
  --environment-identifier $TOOLING_ENV_ID \
  --name "RedshiftCluster-marketing-target" \
  --subscription-target-config '[{ "formName": "RedshiftSubscriptionTargetConfigForm", "content": "{\"databaseName\":\"consumer_db\",\"secretManagerArn\":\"arn:aws:secretsmanager:<region>:<account>:secret:<secret-name>\",\"host\":\"<endpoint>\",\"port\":\"5439\",\"schemaName\":\"marketing\"}" }]' \
  --applicable-asset-types RedshiftViewAssetType RedshiftTableAssetType \
  --manage-access-role $MANAGE_ACCESS_ROLE \
  --provider "Amazon SageMaker" \
  --type RedshiftSubscriptionTargetType \
  --authorized-principals $AUTHORIZED_PRINCIPAL

Verifying the audit trail

To substantiate the governance and traceability claims in this architecture, enable AWS CloudTrail in all three accounts with data events for Secrets Manager and KMS. Enable Amazon Redshift audit logging on clusters to capture connection and query activity through STL_CONNECTION_LOG and STL_QUERY. Subscription approvals and rejections are recorded by SageMaker Unified Studio and emitted to CloudTrail under the datazone.amazonaws.com event source. Look for CreateSubscriptionRequest, AcceptSubscriptionRequest, and RejectSubscriptionRequest events.

Clean up

If you deployed this solution for testing or evaluation purposes and no longer need the resources, we recommend cleaning up to avoid unnecessary costs. Amazon Redshift clusters, Secrets Manager secrets, and SageMaker Unified Studio projects all incur charges when left running. The following steps guide you through a structured teardown in the correct order: subscriptions first, then data assets, and finally the infrastructure itself. This order verifies that no orphaned resources remain.

  • Remove all subscriptions
  • Delete your data assets
  • Delete the projects
    • Delete the project within your SageMaker Unified Studio Domain after all subscriptions are removed. Make sure to delete both the consumer and producer projects.
  • Delete the SageMaker Unified Studio Domain in Account A.

Conclusion

In this post, we demonstrated how Amazon SageMaker Unified Studio simplifies cross-account data governance for Amazon Redshift. By implementing this solution, organizations can move away from ad-hoc, non-auditable data sharing processes to a secure, scalable, and fully governed approach. Amazon SageMaker Unified Studio serves as the central governance layer that data producers and consumers use to publish, discover, and subscribe to data products across AWS accounts. This turns a fragmented data landscape into a well-governed data mesh without the need for custom tooling or manual coordination.

With cross-account data sharing and governance in place, the natural next step is to use this well-governed data for machine learning and generative AI workloads. Because Amazon SageMaker Unified Studio brings together data, analytics, and AI capabilities in a single environment, teams can more efficiently transition from discovering and subscribing to data products to building ML models and generative AI applications, all within the same environment. This reduces the traditional friction between data engineering and data science, accelerating time to value. To get started with establishing your organization’s data mesh using Amazon SageMaker Unified Studio, follow the guidance for Setting up Amazon SageMaker Unified Studio.


About the authors

Bandana Das

Bandana Das is a senior Data Architect in Amazon Web Services and specializes in Data and Analytics. She builds event-driven data architectures to support customers in Data management and data-driven decision making. She is also passionate about enabling customers on their Data management journey to the cloud.

Sindi Cali

Sindi Cali is a ProServe Consultant with AWS Professional Services. She supports customers in building data driven applications in AWS.

Anirban Saha

Anirban Saha is a DevOps Architect at AWS, specializing in architecting and implementation of solutions for customer challenges. He is passionate about well-architected infrastructures, automation, data-driven solutions and helping make the customer’s cloud journey as smooth as possible.

Stoyan Stoyanov

Stoyan Stoyanov works for AWS as a DevOps Engineer. He has more than 10 years of experience in software engineering, cloud technologies, DevOps, data engineering, and security.

Viral Thakkar

Viral Thakkar is a Software Engineer at AWS, working on Amazon DataZone and Amazon SageMaker Unified Studio with a primary focus on distributed systems and data governance with deep expertise in building large-scale data analytics and pipelining solutions. He is passionate about tackling complex distributed systems challenges while also creating tools and automated scripts that simplify day-to-day workflows and improve productivity.