Tag Archives: Technical How-to

Standardizing construct properties with AWS CDK Property Injection

Post Syndicated from Marco Frattallone original https://aws.amazon.com/blogs/devops/standardizing-construct-properties-with-aws-cdk-property-injection/

Standardizing CDK construct properties across a large organization requires repetitive manual effort that scales poorly as teams and repositories grow. Development teams working with AWS Cloud Development Kit (AWS CDK) must apply the same configuration properties across similar resources to meet security, compliance, and operational standards but manual configuration leads to drift, maintenance burden, and compliance gaps. In this post, you learn how to use Property Injection, a feature introduced in AWS CDK v2.196.0, to automatically apply default properties to constructs without modifying existing code.

The Challenge of Infrastructure Standardization

Organizations implementing infrastructure as code face a fundamental tension between developer productivity and operational consistency. CDK provides abstractions for defining cloud resources, but ensuring compliance with organizational security policies, compliance requirements, and operational standards requires repetitive manual configuration.
Consider this scenario: an organization’s security policy requires that all SecurityGroups disable outbound traffic by default. Development teams must apply these settings to every SecurityGroup:


new SecurityGroup(stack, 'api-sg', {
  vpc: myVpc,
  allowAllOutbound: false,        // Required by security policy
  allowAllIpv6Outbound: false     // Required by security policy
});

new SecurityGroup(stack, 'db-sg', {
  vpc: myVpc,
  allowAllOutbound: false,        // Same configuration repeated
  allowAllIpv6Outbound: false     // Same configuration repeated
});

This manual approach creates four specific problems:

  • Configuration drift: Teams omit required properties or apply them inconsistently
  • Maintenance burden: Policy updates require coordinated changes across multiple repositories and teams
  • Developer friction: Repetitive configuration tasks slow development velocity and increase cognitive load
  • Compliance gaps: Manual processes introduce human error, creating security or compliance violations

Custom construct libraries address these challenges but require refactoring every construct instantiation in existing code and create learning curves for development teams already familiar with standard CDK patterns.

Introducing Property Injection

AWS CDK Property Injection addresses these challenges by automatically applying default properties to constructs without requiring changes to existing code.

Property Injection is a feature introduced in AWS CDK v2.196.0 that intercepts construct creation and automatically applies organizational defaults. With this approach, you can enforce standards consistently while preserving existing development workflows and code patterns.

After implementing Property Injection, the same SecurityGroup creation requires only the vpc parameter, security defaults are applied automatically:

// Your existing code remains unchanged
new SecurityGroup(stack, 'my-sg', {
  vpc: myVpc
  // Security defaults applied automatically by Property Injection
});

The key benefits of this approach include:

  • Zero-impact adoption: Existing CDK code continues to work without modification
  • Centralized policy management: Standards are defined once and applied automatically
  • Consistent enforcement: Policies are applied uniformly across all applications and teams
  • Reduced maintenance overhead: Policy updates require changes in only one location
This diagram shows the five-step Property Injection process in a clear two-column format. The left column outlines each process step, while the right column shows the corresponding implementation details with properly formatted TypeScript code. The flow demonstrates how CDK intercepts SecurityGroup creation, applies organizational security defaults through property injectors, merges them with developer-specified properties, and creates a fully configured SecurityGroup that meets both developer requirements and organizational standards.

Figure 1: CDK Property Injection Mechanism

Property Injection operates transparently within CDK, intercepting construct creation to apply predefined defaults before merging them with any properties explicitly provided by developers. This ensures that organizational standards are consistently applied while maintaining the flexibility for developers to override defaults when specific use cases require it.

Understanding the Implementation Approach

Property Injection works by implementing the IPropertyInjector interface, which allows you to define default properties for specific construct types. These injectors are registered with CDK stacks and automatically apply their defaults during construct instantiation.
The implementation follows three steps: define the defaults you want to apply, register the injector with your stack, and let CDK handle the automatic application of these defaults to matching constructs.

Implementation Guide

This section shows you how to implement Property Injection for SecurityGroup constructs.

Step 1: Create a Property Injector

Create a class that implements the IPropertyInjector interface:

import { IPropertyInjector, InjectionContext } from 'aws-cdk-lib';
import { SecurityGroup, SecurityGroupProps } from 'aws-cdk-lib/aws-ec2';

export class SecurityGroupDefaults implements IPropertyInjector {
  readonly constructUniqueId: string;

  constructor() {
    this.constructUniqueId = SecurityGroup.PROPERTY_INJECTION_ID;
  }

  inject(originalProps: SecurityGroupProps, context: InjectionContext): SecurityGroupProps {
    return {
      // Apply organizational defaults
      allowAllIpv6Outbound: false,
      allowAllOutbound: false,
      // Original properties override defaults when specified
      ...originalProps,
    };
  }
}

Step 2: Add the Injector to Your Stack

Apply the injector to your CDK stack:

import { Stack } from 'aws-cdk-lib';
import { SecurityGroupDefaults } from './security-defaults';

const stack = new Stack(app, 'MyStack', {
  propertyInjectors: [
    new SecurityGroupDefaults()
  ]
});

Step 3: Use Constructs Normally

Create constructs as usual. The injector applies defaults automatically:

// This SecurityGroup receives the injected defaults:
// - allowAllOutbound: false
// - allowAllIpv6Outbound: false
new SecurityGroup(stack, 'my-sg', {
  vpc: myVpc
});

// You can override defaults when necessary
new SecurityGroup(stack, 'special-sg', {
  vpc: myVpc,
  allowAllOutbound: true  // Overrides the injected default
});
This side-by-side comparison shows the difference between manual configuration and Property Injection. The left side (Before) shows three SecurityGroup definitions, each requiring manual specification of allowAllOutbound: false and allowAllIpv6Outbound: false, leading to repetitive code, inconsistency risk, and maintenance burden. The right side (After) shows the same SecurityGroups created with the VPC parameter alone after a one-time Property Injection setup, demonstrating the DRY principle, consistent defaults, and reduced maintenance.

Figure 2: CDK Code Before vs After Property Injection

Property Injection vs L2 Constructs

You can achieve the same enforcement of default properties by creating custom L2 constructs with built-in defaults. However, Property Injection is better suited for standardizing existing codebases without refactoring, while L2 Constructs are better suited for new projects where you want custom APIs and multi-resource abstractions.

This decision tree guides the selection between Property Injection and L2 Constructs for CDK standardization. Starting with existing CDK applications, it evaluates willingness to accept potential breaking changes from new defaults. If breaking changes are acceptable or no existing code exists, it assesses whether custom APIs, naming improvements, or multi-resource patterns are needed beyond simple defaults. The tree leads to three outcomes: Property Injection (blue) for transparent defaults with existing code compatibility, L2 Constructs (orange) for custom APIs and purpose-built abstractions, or a Hybrid approach (green) combining both techniques for maximum flexibility.

Figure 3: Decision Tree – Property Injection vs L2 Constructs

Implementation Comparison

Consider an application with multiple SecurityGroup instantiations that need standardized security defaults.

L2 Construct approach requires creating a custom construct and updating each instantiation:

// Step 1: Create custom L2 construct
export class SecureSecurityGroup extends SecurityGroup {
  constructor(scope: Construct, id: string, props: SecurityGroupProps) {
    super(scope, id, {
      allowAllOutbound: false,
      allowAllIpv6Outbound: false,
      ...props
    });
  }
}

// Step 2: Update each instantiation throughout your codebase
// Change from:
new SecurityGroup(stack, 'sg1', { vpc: myVpc })
new SecurityGroup(stack, 'sg2', { vpc: myVpc })
new SecurityGroup(stack, 'sg3', { vpc: myVpc })

// To:
new SecureSecurityGroup(stack, 'sg1', { vpc: myVpc })
new SecureSecurityGroup(stack, 'sg2', { vpc: myVpc })
new SecureSecurityGroup(stack, 'sg3', { vpc: myVpc })

Property Injection approach requires one-time stack configuration:

// Step 1: Add injector to stack configuration
stack.propertyInjectors = [new SecurityGroupDefaults()];

// Step 2: Existing SecurityGroup calls receive defaults automatically
new SecurityGroup(stack, 'sg1', { vpc: myVpc })  // Gets defaults
new SecurityGroup(stack, 'sg2', { vpc: myVpc })  // Gets defaults  
new SecurityGroup(stack, 'sg3', { vpc: myVpc })  // Gets defaults

Key Differences

Property Injection works with existing construct calls, requiring no changes to how developers instantiate SecurityGroups or other constructs. This approach overrides constructs from external libraries and can be implemented without modifying existing code. Developers continue using familiar CDK APIs without learning new interfaces.

L2 Constructs require updating all constructor calls throughout your codebase. This approach cannot modify third-party construct creation since you must change each instantiation to use your custom construct. Implementation requires refactoring existing code and developers must learn your custom construct APIs instead of standard CDK interfaces. L2 constructs serve multiple purposes beyond complex business logic – simple L2 constructs provide domain-specific naming conventions and cleaner APIs, while complex L2 constructs orchestrate three or more resources and implement business rules.

When to Choose Each Approach

Choose Property Injection when you need to standardize existing infrastructure. Property Injection excels in scenarios where you already have CDK applications deployed and need to apply consistent defaults retroactively. Property Injection works transparently with existing code, requiring no changes to how developers instantiate constructs. This makes it useful when you have existing CDK applications that you want to standardize without disrupting current development workflows.

Property Injection also solves the challenge of applying defaults to constructs from third-party libraries. Since you cannot modify external library code, Property Injection enforces organizational standards on any construct type, regardless of its source. Additionally, when you want to implement standards without changing existing code, Property Injection operates at the framework level, automatically applying defaults during construct instantiation without requiring developers to modify their existing implementations.

Choose L2 Constructs when you need custom APIs or multi-resource patterns. L2 Constructs provide the right abstraction when you want to create purpose-built interfaces that differ from standard CDK APIs. This includes simple wrappers with domain-specific naming, complex business logic, validation rules, or multi-resource orchestration patterns. L2 Constructs excel when you want to create opinionated APIs that simplify common patterns by hiding complexity behind intuitive interfaces.

L2 Constructs suit new application development where you can design the API from the start. This approach creates purpose-built abstractions that match your organization’s specific use cases and terminology. Unlike Property Injection, which applies defaults to existing construct APIs, with L2 Constructs you can design entirely new APIs that directly represent your business domain and operational patterns.

Implementation Patterns

Stack Integration Methods

The CDK provides two methods for adding Property Injectors to stacks:

This diagram demonstrates two methods for adding Property Injectors to CDK stacks. Method 1 (blue) shows adding injectors directly in the Stack constructor’s propertyInjectors array. Method 2 (orange) shows using PropertyInjectors.of(stack).add() after stack creation. Both methods produce identical results with green checkmarks indicating success. The diagram includes usage examples showing normal SecurityGroup instantiation (blue) that inherits defaults automatically, and override scenarios (orange) where developers explicitly override injected defaults. The bottom section shows the resulting CloudFormation output: default SecurityGroups have empty egress rules (green), while overridden ones include outbound traffic rules (orange).

Figure 4: CDK Stack Integration Methods

Method 1: Stack Constructor

const stack = new Stack(app, 'MyStack', {
  propertyInjectors: [new SecurityGroupDefaults()]
});

Method 2: PropertyInjectors.of()

const stack = new Stack(app, 'MyStack');
PropertyInjectors.of(stack).add(new SecurityGroupDefaults());

Both methods produce the same result. Choose the method that best fits your existing code structure. For more details, see the PropertyInjectors API documentation.

Organization-Wide Implementation

For organization-wide standardization, create a shared library of injectors:

// @myorg/cdk-injectors package
export const ORGANIZATION_INJECTORS: IPropertyInjector[] = [
  new SecurityGroupDefaults(),
  new LambdaFunctionDefaults(),
  new S3BucketDefaults(),
];

// Teams import and use the shared injectors
import { ORGANIZATION_INJECTORS } from '@myorg/cdk-injectors';

const stack = new Stack(app, 'TeamStack', {
  propertyInjectors: ORGANIZATION_INJECTORS
});

Scope Hierarchy

Property Injectors can be applied at different levels in the CDK construct tree:

This diagram illustrates the three-level hierarchy of Property Injector scopes in CDK with integrated resolution examples. The App level (blue) shows a BucketInjector ‘b1’ that applies globally, with an example showing how Stack2 buckets use this injector. The Stage level (green) demonstrates a FunctionInjector ‘f1’ that applies to all stacks within the stage, including an example of Stack1 functions using this injector. The Stack level (orange) shows two stacks: Stack1 with its own BucketInjector ‘b2’ that overrides the app-level injector, and Stack2 with no injectors that inherits from parent scopes. The Resolution Rules box (green) explains that CDK searches from most specific (stack) to most general (app), with the first match winning per construct type. Arrows show the hierarchical relationship between scopes.

    Figure 5: CDK Scope Hierarchy & Injector Resolution
  • App level: Applies to all stacks in the application
  • Stage level: Applies to all stacks within a specific stage
  • Stack level: Applies only to constructs within a specific stack

CDK searches for applicable injectors starting from the construct’s immediate parent scope and moving upward. The first matching injector for each construct type is used.

Best Practices

When implementing Property Injection, begin with high-impact constructs like SecurityGroups, VPCs, and Lambda functions that require repetitive configuration.

These constructs have the highest frequency of misconfiguration and the most direct compliance impact, making them the most valuable targets for early adoption.

Document your defaults by explaining what properties your injectors provide and why. Include examples and link to relevant policies that drive the requirements. With this documentation, developers can understand standards and make informed override decisions.

Write automated tests using CDK testing utilities to verify that injectors apply expected defaults. Test both standard scenarios and cases where developers override properties to prevent regressions when updating injector logic.

Version injectors carefully using semantic versioning principles because changes affect all applications. Coordinate updates across teams and provide migration guides for breaking changes or changes to default values.

Design override mechanisms so that developers can handle edge cases while benefiting from organizational standards. Property Injection operates as defaults, not restrictions, so design injectors to merge gracefully with developer-specified properties.

Limitations and Considerations

Property Injection operates as a default mechanism rather than a compliance enforcement system. Developers retain the ability to override injected properties, which means organizations cannot rely solely on Property Injection for strict compliance requirements. For teams that need mandatory compliance, combine Property Injection with CDK Aspects or AWS Config rules to validate and enforce standards.

The feature works exclusively with L2 constructs, as documented in the official AWS CDK guidance. The IPropertyInjector interface targets specific L2 construct types, and L1 (CloudFormation) constructs use different instantiation patterns that bypass the property injection mechanism entirely. Organizations with L1 construct usage need alternative standardization approaches.

Property Injection introduces debugging complexity because injected properties do not appear directly in application code. Developers troubleshooting construct behavior must understand which injectors apply to specific construct types and how those injectors modify properties. This hidden behavior requires documentation that lists each injector, the properties it sets, and the policy it enforces, along with clear naming conventions to maintain code clarity.

The feature requires CDK v2.196.0 or later, which affects adoption timelines for organizations using older CDK versions. Teams must plan upgrade paths and test compatibility before implementing Property Injection across their applications.

Conclusion

Property Injection provides a mechanism for applying consistent default properties to CDK constructs without requiring changes to existing code. This approach reduces repetitive configuration, improves consistency, and simplifies maintenance of CDK applications.
Property Injection is the right choice for organizations that need to standardize construct configurations across existing codebases while preserving developer workflows. When combined with proper testing and documentation, Property Injection becomes a reliable foundation for infrastructure governance across your organization.

About the authors:

Put Cheung

Put Cheung is a Senior Software Development Engineer at AWS Security. He is a part of a team that is making it easier for builders to configure AWS Resources securely. AWS CDK Property Injection is an important step toward this goal.

Rico Huijbers

Rico Huijbers is a Software Engineer at Amazon Web Services. He is extremely lazy and is therefore on a quest to eradicate the need for repetitive manual work from software engineering. Rico loves working on AWS CDK—it’s the tool he wishes he had 5 years earlier.

Marco Frattallone

Marco Frattallone is a Senior Technical Account Manager at AWS focused on supporting Partners. He works closely with Partners to help them build, deploy, and optimize their solutions on AWS, providing guidance and leveraging best practices. Marco focuses on helping Partners adopt emerging AWS services and translate technical capabilities into business outcomes. Outside work, he enjoys outdoor cycling, sailing, and exploring new cultures.

Set up production-ready monitoring for Amazon MSK using CloudWatch alarms

Post Syndicated from Yashika Jain original https://aws.amazon.com/blogs/big-data/set-up-production-ready-monitoring-for-amazon-msk-using-cloudwatch-alarms/

Organizations running Apache Kafka as their streaming platform need comprehensive monitoring to maintain reliable operations. Without proper visibility into broker health, resource utilization, and data flow metrics, teams risk service disruptions, data loss, and degraded performance that can impact critical business operations. Effective monitoring and alerting are essential to detect anomalies early, from high system load to connectivity issues, enabling teams to take preventive action before problems affect production workloads.

Amazon Managed Streaming for Apache Kafka (Amazon MSK) addresses these monitoring challenges by publishing detailed metrics to Amazon CloudWatch. The service emits metrics at 1-minute intervals for provisioned (Standard) clusters, with flexible monitoring levels (DEFAULT, PER_BROKER, PER_TOPIC_PER_BROKER, or PER_TOPIC_PER_PARTITION) to control granularity and cost. At the DEFAULT level (free), cluster-level metrics are available; higher levels (paid) expose broker-level, per-topic and per-partition metrics.

In this post, I show you how to implement effective monitoring for your MSK clusters using Amazon CloudWatch. You’ll learn how to track critical metrics like broker health, resource utilization, and consumer lag, and set up automated alerts to prevent operational issues. By following these practices, you can work to improve streaming operations reliability, optimize resource usage, and support high availability for your mission-critical applications.

Key metrics to monitor

This article groups important Amazon MSK metrics into logical categories. For each, we highlight key metrics and what they indicate:

  1. Broker Health and Cluster Availability:
    • ActiveControllerCount is a cluster-level metric where each broker reports whether it’s the active controller (1) or not (0). In a healthy cluster, exactly one broker serves as the active controller at any time. When viewing this metric with the average statistic, the value equals 1 divided by the number of brokers. For example, a 3-broker cluster shows 0.33 (1/3). Set CloudWatch alarm thresholds accordingly—for 6 brokers, alert if average falls below 0.166(1/6). When using the sum statistic, the value should always be 1, indicating one active controller regardless of cluster size. If the sum differs from 1, a controller election is in progress—typically during maintenance activities, configuration changes, or rolling restarts.
      Note: For a KRaft-based clusters, the ActiveControllerCount is only exposed on dedicated controller endpoints so the sample count is 3 and only controller will report value of 1. Thus, the average is always 0.33 no matter how many brokers there are in the cluster. To monitor the broker health for Kraft-based clusters, check LeaderCount metric. If a broker is not emitting any metric, then it’s a good indication that broker might be unhealthy.
    • OfflinePartitionsCount (cluster): Number of partitions with no active leader. Non-zero values mean data is temporarily unavailable or unwritable. Trigger alerts if it rises above 0.
    • UnderReplicatedPartitions (per broker): Number of partitions where not all replicas are caught up. This should stay at 0 under normal conditions. Spikes indicate traffic exceeds capacity or replication lag; sustained values often mean a configuration/ACL issue. Refer to Troubleshoot your Amazon MSK cluster
    • UnderMinIsrPartitionCount (per broker): Partitions below the minimum in-sync replica (ISR) count. A non-zero value means potential data loss risk if brokers fail. Monitor to ensure replication is healthy. Refer to Custom configurations
    • GlobalPartitionCount (cluster): Total number of partitions across all topics (leaders only). Useful for capacity planning and sanity checks.
    • PartitionCount (per broker): Number of partitions (including replicas) hosted by a broker. Sudden changes may indicate re-balances. (Excess partitions per broker can degrade performance).
  2. Resource Utilization:
    • CPU: Total broker CPU utilization is defined as CpuUser + CpuSystem. Best practice is to keep average CPU utilization under 60% . Set alarms on the sum of user+system to detect overload.
    • CPUCreditBalance / CPUCreditUsage (per broker): For burstable instance types(T3), tracks earned/spent CPU credits. A declining credit balance or high credit usage warns that the instance may be CPU-starved.
    • Memory: MemoryUsed, MemoryFree (per broker) show RAM usage. Critically, HeapMemoryAfterGC (per broker) reports JVM heap usage (%) after garbage collection. AWS recommends alerting if HeapMemoryAfterGC exceeds 60%, to avoid out-of-memory issues.
    • Disk: Kafka brokers use attached EBS storage for topic data. Monitor KafkaDataLogsDiskUsed (per broker) – percentage of disk used by message logs. Best practice: alarm when data log usage exceeds 85%. Also track RootDiskUsed: the percentage of the root disk used by the broker.
    • EBS I/O: Volume metrics (per broker) such as VolumeQueueLength, VolumeReadOps, VolumeWriteOps, VolumeReadBytes, VolumeWriteBytes indicate I/O latency and throughput. Rising queue lengths or latency (such as VolumeTotalReadTime) suggest disk contention.
    • Network: Basic network stats per broker include NetworkRxPackets, NetworkTxPackets, and errors/drop counts (NetworkRxErrors, NetworkTxErrors, NetworkRxDropped, NetworkTxDropped). Unexpected errors or drops can indicate network issues.
  3. Topic and Partition Activity:
    • Throughput: BytesInPerSec and BytesOutPerSec measure inbound/outbound data rates per broker or per topic. Sustained drops can signal lost producers/consumers; spikes may require scaling.
    • Replication Traffic: ReplicationBytesInPerSec/ReplicationBytesOutPerSec (per topic) show inter-broker replication volume.
    • Consumer Lag: Consumer lag metrics quantify the difference between the latest data written to your topics and the data read by your applications. Amazon MSK provides the following consumer-lag metrics, which you can get through Amazon CloudWatch or through open monitoring with Prometheus: EstimatedMaxTimeLag, EstimatedTimeLag, MaxOffsetLag, OffsetLag, and SumOffsetLag. For information about these metrics, see Amazon MSK metrics for monitoring Standard brokers with CloudWatch.
  4. Client Connections :
    • ConnectionCount (per broker): Total active connections (clients + inter-broker). Sudden drops or sustained high counts (hitting limits) merit attention.
    • ClientConnectionCount (per broker, with auth filter): Active authenticated client connections.
    • ConnectionCreationRate / ConnectionCloseRate (per broker): New or closed connections per second. Spikes in connection churn may indicate client issues.
    • Authentication: IAMNumberOfConnectionRequests and IAMTooManyConnections (per broker) show IAM auth request rates and throttle breaches (limit of 100 simultaneous connections).
  5. Network Bandwidth Metrics:
    • TrafficShaping > 0 (any throttling) metric serves as your primary warning signal. When this value exceeds zero, your MSK cluster is experiencing network throttling at the EC2 layer, with packets being dropped or queued due to exceeded allocations. This throttling manifests as reduced throughput, increased latency, and potential network errors that impact both producer and consumer performance. TrafficShaping issues stem from two possible bandwidth limitations: BwInAllowanceExceeded & BwOutAllowanceExceeded :
    • BwInAllowanceExceeded tracks when inbound aggregate bandwidth surpasses broker maximums.
    • BwOutAllowanceExceeded monitors when outbound aggregate bandwidth exceeds limits.
      Both BwInAllowanceExceeded and BwOutAllowanceExceeded metrics directly contribute to overall network throttling events.
  6. Other Operational Metrics:
    • Thread Pools: RequestHandlerAvgIdlePercent, NetworkProcessorAvgIdlePercent (per broker) show how busy Kafka’s internal thread pools are. Consistently low idle (%) can indicate bottlenecks.
    • ZooKeeper: For ZooKeeper-based MSK clusters, ZooKeeperRequestLatencyMsMean and ZooKeeperSessionState reflect ZK performance (for older Kafka versions that use Zookeeper). For ZooKeeperSessionState, anything other than 1 for 5-10 mins should be alarming as there can be chances broker has an issue or zookeeper is not able to connect to brokers due to some intermittent network issue.
    • Tiered Storage: For clusters with tiered storage enabled, Amazon MSK provides metrics like RemoteFetchBytesPerSec, RemoteCopyBytesPerSec, RemoteLogSizeBytes, and related error/queue metrics. These track offloading to remote storage.
    • Intelligent rebalancing metrics: For MSK Provisioned clusters using Express brokers, Amazon MSK provides two key metrics to monitor rebalancing operations: RebalanceInProgress and UnderProvisioned metrics. See Monitor Intelligent rebalancing metrics

By grouping metrics into these categories, you can build dashboards and alerts that comprehensively cover Amazon MSK health and performance. Amazon CloudWatch also provides automatic dashboards for Amazon MSK.

Let’s take a quick look on how to access CloudWatch automatic dashboard. In the AWS Console, go to the CloudWatch service. When in the CloudWatch console, select Dashboards. Open the Automatic dashboard tab and search for MSK in the Filter Bar.

These dashboards offer per-configured visualizations of key metrics, enabling quick insights into the health and performance of your MSK clusters.

Recommended CloudWatch alarms

Setting alarms on key metrics helps catch issues early. Detecting issues early is crucial in streaming applications where every second counts. A single failing broker can trigger a chain reaction – halting data ingestion, backing up upstream systems, and breaking downstream applications. This can quickly escalate from delayed order processing to lost revenue. Proactive monitoring helps catch and fix problems before they impact your business operations. Based on AWS best practices and experience, consider alarms such as:

Metric (Dimension) Alarm Condition Rationale
ActiveControllerCount (cluster) ≠ 1 (count) Only one active controller should exist. Deviation implies cluster instability.
CPU Utilization (Sum(CPUUser+CPUSystem), per broker) > 60% (average) for 5+ mins Helps maintain headroom for broker load and maintenance. High CPU may slow processing as outlined in the MSK best practices documentation
HeapMemoryAfterGC (broker) > 60% (percentage) Indicates Kafka heap is filling up. Helps prevent OOM by alerting early.
KafkaDataLogsDiskUsed (broker) ≥ 85% (percent) Warns that disk is nearly full. Helps prevent data loss by providing time for scaling or cleanup.
OfflinePartitionsCount (cluster) > 0 (count) Any offline partition means unavailable data. Immediate investigation needed.
UnderReplicatedPartitions (broker) > 0 (count) No replicas lagging under healthy conditions. Spikes or sustained lag can indicate overload or ACL misconfiguration.
UnderMinIsrPartitionCount (broker) > 0 (count) There must be topics with partitions that have either less in-sync replicas than the min.insync.replicas setting or with RF=MinISR. To find these topics whose partitions are under replicated, use command:
<path-to-your-kafka-installation>/bin/kafka-topics.sh –bootstrap-server <bootstrap-server:port> —command-config client.properties –describe –under-min-isr-partitions
ConnectionCount (broker) Sudden drop (e.g. < 90% of baseline) or spike above high threshold Detect client connectivity issues or connection floods. Unexpected drops may mean a broker is unreachable. Refer to Amazon MSK Standard broker quota
CPUCreditBalance (for T3 broker) < some low threshold (e.g. 10 credits) For burstable instances, alerts when credits are nearly exhausted, which degrades performance.
VolumeQueueLength (broker) > 0 (sustained) or rising Indicates I/O operations are queuing, possible disk bottleneck.
NetworkRxErrors/TxErrors (broker) > 0 (count) Any network errors can cause packet loss or disconnections.
IAMTooManyConnections (broker) > 0 (count) Exceeding IAM connection limit (100) blocks new connections.
Consumer Lag (MaxOffsetLag or SumOffsetLag) (per consumer-group/topic) > threshold (depends on SLAs, e.g. growing beyond expected) Alerts on slow consumers so you can scale consumers or investigate backlogs.
TrafficShaping > 0 (any throttling) This is an indication that brokers are exceeding their allocated network bandwidth.

These are illustrative thresholds; adjust them for your workload and SLAs. The remaining metrics listed in the CloudWatch metrics for Standard and Express brokers documentation are susceptible to downstream impact from anomalies in the primary metrics above. It is recommended to enable CloudWatch alarms on a single test cluster first to validate thresholds before extending coverage across your MSK fleet.

Conclusion

In this post, we covered the important CloudWatch metrics and alarms for monitoring Amazon MSK clusters effectively. By implementing these recommended alarms, you can proactively detect and respond to potential issues before they impact your Kafka workloads. To learn more about Amazon MSK monitoring, refer to the Amazon MSK Monitoring Best Practices documentation or explore our Amazon MSK Workshops hands-on experience.


About the authors

Yashika Jain

Yashika Jain

Yashika is a Senior Cloud Analytics Engineer at AWS, specializing in real-time analytics and event-driven architectures. She is committed to helping customers by providing deep technical guidance, driving best practices across real-time data platforms and solving complex issues related to their streaming data architectures.

Standardize Amazon Redshift operations using Templates

Post Syndicated from Nidhi Nayak original https://aws.amazon.com/blogs/big-data/standardize-amazon-redshift-operations-using-templates/

Over the past year, Amazon Redshift has introduced capabilities that simplify operations and enhance productivity. Building on this momentum, we’re addressing another common operational challenge that data engineers face daily: managing repetitive data loading operations with similar parameters across multiple data sources. This intermediate-level post introduces AWS Redshift Templates, a new feature that you can use to create reusable command patterns for the COPY command, reducing redundancy and improving consistency across your data operations.

The challenge: Managing repetitive data operations at scale

Meet AnyCompany, a fictional data aggregation company that processes customer transaction data from over 50 retail clients. Each client sends daily delimited text files with similar structures:

customer transactions | product catalogs | inventory updates

While the data format is largely consistent across clients (pipe-delimited files with headers, UTF-8 encoding), the sheer volume of COPY commands required to load this data has become a development and maintenance overhead.

Their data engineering team faces several pain points:

  • Repetitive parameter specification: Each COPY command requires specifying the same parameters for delimiter, encoding, error handling, and compression settings
  • Inconsistency risks: With multiple team members writing COPY commands, slight variations in parameters lead to data ingestion failures
  • Maintenance overhead: When they need to adjust error thresholds or encoding settings, they must update hundreds of individual COPY commands across their extract, transform, and load (ETL) pipelines
  • Onboarding complexity: New team members struggle to remember all the required parameters and their optimal values

Additionally, a few clients send data in slightly different formats. Some use comma delimiters instead of pipes or have different header configurations. The team needs flexibility to handle these exceptions without completely rewriting their data loading logic.

Introducing Redshift Templates

You can address these challenges by using Redshift Templates to store commonly used parameters for COPY commands as reusable database objects. Think of templates as blueprints for your data operations where you can define your parameters once, then reference them across multiple COPY commands.

Template management best practices

Before exploring implementation scenarios, let’s establish best practices for template management to ensure your templates remain maintainable and secure.

  1. Use descriptive names that indicate purpose:
    CREATE TEMPLATE analytics.csv_client_data_load;
    CREATE TEMPLATE analytics.json_retail_data_load;

  2. Implement least privilege access:
    -- Grant specific permissions to roles
    GRANT USAGE FOR TEMPLATES IN SCHEMA analytics TO ROLE data_engineers;
    GRANT ALTER FOR TEMPLATES IN SCHEMA reporting TO ROLE senior_analysts;
    -- Revoke broad permissions
    REVOKE ALL ON TEMPLATE analytics.csv_load FROM PUBLIC;

  3. Query the system view to track template usage:
    SELECT database_name, schema_name, template_name, 
           create_time, last_modified_time
    FROM sys_redshift_template;

  4. Document each template, including:
    • Purpose and use cases
    • Parameter explanations
    • Ownership and contact information
    • Change history

Solution overview

Let’s explore how AnyCompany uses Redshift Templates to streamline their data loading operations.

Scenario 1: Standardizing client data ingestion

AnyCompany receives transaction files from multiple retail clients with consistent formatting. They create a template that encapsulates their standard loading parameters:

-- Create a reusable template for standard client data loads
CREATE TEMPLATE data_ingestion.standard_client_load
FOR COPY
AS
DELIMITER '|'
IGNOREHEADER 1
ENCODING UTF8
MAXERROR 100
COMPUPDATE OFF
STATUPDATE ON
ACCEPTINVCHARS
TRUNCATECOLUMNS;

This template defines their standard approach:

  • DELIMITER '|' specifies pipe-delimited files
  • IGNOREHEADER 1 skips the header row
  • ENCODING UTF8 facilitates proper character encoding
  • MAXERROR 100 allows up to 100 errors before failing, providing resilience for minor data quality issues
  • COMPUPDATE OFF helps prevent automatic compression analysis during loading for faster performance
  • STATUPDATE ON keeps table statistics current for query optimization
  • ACCEPTINVCHARS replaces invalid UTF-8 characters rather than failing
  • TRUNCATECOLUMNS truncates data that exceeds column width rather than failing

Now, loading data from a standard client becomes remarkably straightforward:

-- Load transaction data from Client A
COPY transactions_client_a
FROM 's3://amzn-s3-demo-bucket/client-a/transactions/'
IAM_ROLE default
USING TEMPLATE data_ingestion.standard_client_load;
-- Load transaction data from Client B
COPY transactions_client_b
FROM 's3://amzn-s3-demo-bucket/client-b/transactions/'
IAM_ROLE default
USING TEMPLATE data_ingestion.standard_client_load;
-- Load product catalog from Client C
COPY products_client_c
FROM 's3:// amzn-s3-demo-bucket/client-c/products/'
IAM_ROLE default
USING TEMPLATE data_ingestion.standard_client_load;

Notice how clean and maintainable these commands are. Each COPY statement specifies only:

  1. The target table
  2. The Amazon Simple Storage Service (Amazon S3) source location
  3. The default AWS Identity and Access Management (IAM) role for authentication
  4. The template reference

The complex formatting and error handling parameters are neatly encapsulated in the template, facilitating consistency across the data loads.

Scenario 2: Handling client-specific variations with parameter overrides

AnyCompany has two clients (Client D, and E) who send comma-delimited files instead of pipe-delimited files. Rather than creating an entirely separate template, they can override specific parameters while still using the template’s other settings:

-- Load data from Client D with comma delimiter (overriding template)
COPY transactions_client_d
FROM 's3://amzn-s3-demo-bucket/client-d/transactions/'
IAM_ROLE default
DELIMITER ','  -- Override the template's pipe delimiter
USING TEMPLATE data_ingestion.standard_client_load;
-- Load data from Client E with comma delimiter and no header
COPY transactions_client_e
FROM 's3://amzn-s3-demo-bucket/client-e/transactions/'
IAM_ROLE default
DELIMITER ','      -- Override delimiter
IGNOREHEADER 0     -- Override header setting
USING TEMPLATE data_ingestion.standard_client_load;

This demonstrates the Redshift Templates parameter hierarchy:

  1. Command-specific parameters (highest priority): Parameters explicitly specified in your COPY command take precedence
  2. Template parameters (medium priority): Parameters defined in the template are used when not overridden
  3. Amazon Redshift default parameters (lowest priority): Default values apply when neither command nor template specifies a value

This three-tier approach provides the perfect balance between standardization and flexibility. You maintain consistency where it matters while retaining the ability to handle exceptions gracefully.

Scenario 3: Simplified template maintenance

Six months after implementing templates, AnyCompany’s data quality team recommends increasing the error threshold from 100 to 500 to better handle occasional data quality issues from upstream systems. With templates, this change is trivial:

-- Update the template to increase error tolerance
ALTER TEMPLATE data_ingestion.standard_client_load
SET MAXERROR TO 500;

This single command instantly updates the error handling behavior for the future COPY operations using this template without needing to hunt through hundreds of ETL scripts or risking missing updates in some pipelines. They can also add new parameters as their requirements evolve:

-- Add compression parameter to improve load performance
ALTER TEMPLATE data_ingestion.standard_client_load
ADD GZIP;

To remove a template when it’s no longer needed:

DROP TEMPLATE data_ingestion.standard_client_load;

Scenario 4: Environment-specific templates for development and production

AnyCompany maintains separate templates for development and production environments, with different error tolerance levels:

-- Development template with lenient error handling
CREATE TEMPLATE data_ingestion.dev_client_load
FOR COPY
AS
DELIMITER '|'
IGNOREHEADER 1
ENCODING UTF8
MAXERROR 1000        -- More lenient for testing
COMPUPDATE OFF
STATUPDATE OFF;      -- Skip stats updates in dev
-- Production template with strict error handling
CREATE TEMPLATE data_ingestion.prod_client_load
FOR COPY
AS
DELIMITER '|'
IGNOREHEADER 1
ENCODING UTF8
MAXERROR 50          -- Stricter for production
COMPUPDATE OFF
STATUPDATE ON;       -- Keep stats current in prod

This approach helps ensure that data quality issues are caught early in production while allowing flexibility during development and testing.

Key benefits

The key benefits of using templates include:

  • Consistency and standardization: Templates help maintain consistency across different operations by making sure that the same set of parameters and configurations are used every time. This is particularly valuable in large organizations where multiple users work on the same data pipelines.
  • Ease of use and timesaving: Instead of manually specifying the parameters for each command execution, users can reference a pre-defined template. This saves time and reduces the chances of errors caused by manual input.
  • Flexibility with parameter overrides: While templates provide standardization, they don’t sacrifice flexibility. You can override a template parameter directly in your COPY command when handling exceptions or special cases.
  • Simplified maintenance: When changes need to be made to parameters or configurations, updating the corresponding template propagates the changes across the instances where the template is used. This significantly reduces maintenance effort compared to manually updating each command individually.
  • Collaboration and knowledge sharing: Templates serve as a knowledge base, capturing best practices and optimized configurations developed by experienced users. This facilitates knowledge sharing and onboarding of new team members, reducing the learning curve and facilitating consistent usage of proven configurations.

Additional use cases across industries

Templates can be used across industries.

Financial services: Standardizing regulatory data loads

A financial institution needs to load transaction data from multiple branches with consistent formatting requirements:

-- Create template for branch transaction loads
CREATE TEMPLATE compliance.branch_transaction_load
FOR COPY
AS
FORMAT CSV
DELIMITER ','
IGNOREHEADER 1
ENCODING UTF8
DATEFORMAT 'YYYY-MM-DD'
TIMEFORMAT 'YYYY-MM-DD HH:MI:SS'
MAXERROR 0           -- Zero tolerance for compliance data
COMPUPDATE OFF;
-- Load data from different branches
COPY branch_transactions_east
FROM 's3://amzn-s3-demo-source-bucket/east-branch/transactions/'
IAM_ROLE default
USING TEMPLATE compliance.branch_transaction_load;
COPY branch_transactions_west
FROM 's3://amzn-s3-demo-source-bucket/west-branch/transactions/'
IAM_ROLE default
USING TEMPLATE compliance.branch_transaction_load;

Healthcare: Loading patient data with strict standards

A healthcare analytics company standardizes their patient data ingestion across multiple hospital systems:

-- Create template for HIPAA-compliant data loads
CREATE TEMPLATE healthcare.patient_data_load
FOR COPY
AS
FORMAT CSV
DELIMITER '|'
IGNOREHEADER 1
ENCODING UTF8
ACCEPTINVCHARS
TRUNCATECOLUMNS
MAXERROR 10
COMPUPDATE OFF;
-- Apply to different hospital systems
COPY hospital_a_patients
FROM 's3://amzn-s3-demo-destination-bucket/hospital-a/patients/'
IAM_ROLE default
USING TEMPLATE healthcare.patient_data_load;
COPY hospital_b_patients
FROM 's3://amzn-s3-demo-destination-bucket/hospital-b/patients/'
IAM_ROLE default
USING TEMPLATE healthcare.patient_data_load;

Retail: JSON data loading standardization

A retail company processes JSON-formatted product catalogs from various suppliers:

-- Create template for JSON product data
CREATE TEMPLATE retail.json_product_load
FOR COPY
AS
FORMAT JSON 'auto'
TIMEFORMAT 'auto'
ENCODING UTF8
MAXERROR 100
COMPUPDATE OFF;
-- Load from different suppliers
COPY products_supplier_a
FROM 's3://amzn-s3-demo-logging-bucket/supplier-a/products/'
IAM_ROLE default
USING TEMPLATE retail.json_product_load;
COPY products_supplier_b
FROM 's3://amzn-s3-demo-logging-bucket/supplier-b/products/'
IAM_ROLE default
USING TEMPLATE retail.json_product_load;

Conclusion

In this post, we introduced Redshift Templates and showed examples of how they can standardize and simplify your data loading operations across different scenarios. By encapsulating common COPY command parameters into reusable database objects, templates help remove repetitive parameter specifications, facilitate consistency across teams, and centralize maintenance. When requirements evolve, a single template update propagates quickly across the operations, reducing operational overhead while maintaining flexibility to override parameters for use cases.

Start using Redshift Templates to transform your data ingestion workflows. Create your first template for your most common data loading pattern, then gradually expand coverage across your pipelines. Your team will immediately benefit from cleaner code, faster onboarding, and simplified maintenance. To learn more about Redshift Templates and explore additional configuration options, see the Amazon Redshift documentation.

Automate AWS Lambda Runtime Upgrades with AWS Transform custom

Post Syndicated from Venugopalan Vasudevan original https://aws.amazon.com/blogs/devops/automate-aws-lambda-runtime-upgrades-with-aws-transform-custom/

Introduction

Organizations carry a growing burden of technical debt — aging codebases, outdated runtimes, and legacy frameworks that slow innovation, increase security risk, and inflate maintenance costs. Addressing this debt requires tackling a wide range of code transformation challenges: version upgrades, runtime migrations, framework transitions, and language translations, all of which must be repeated across multiple codebases. Today, most organizations perform these tasks manually, consuming 20–30% of enterprise software development effort. Even where automation exists, it’s typically narrow in scope and requires significant upfront investment — leaving most organizations unable to scale transformations effectively and, as a result, unable to meaningfully reduce the technical debt that continues to compound over time.

AWS Transform custom addresses this gap — an intelligent AI agent that learns organization-specific code transformations, executes them consistently at scale, and improves from developer feedback, without requiring specialized automation expertise. This blog explores how AWS Transform custom tackles one of the most pressing transformation challenges today.

Managing Lambda Runtime Lifecycles

AWS Lambda follows a runtime deprecation policy that aligns with the end of community long-term support for programming languages, with all current deprecation schedules available in the AWS Lambda runtime documentation. As upstream language maintainers deprecate runtime versions – including Python 3.8, Node.js 14, and Java 8  , organizations face the critical challenge of upgrading hundreds or thousands of Lambda functions before these runtimes reach end-of-life.

When a Lambda runtime reaches end of life, functions lose access to security patches and technical support, leaving applications potentially exposed to known vulnerabilities and compliance risks. Performance degrades as optimizations in newer runtimes go unrealized, and technical debt compounds — the longer you wait, the harder and more expensive the migration becomes.

For organizations managing hundreds or thousands of Lambda functions across multiple runtimes and languages, the effort is amplified by the scale of coordination required, the manual burden of testing and validation, and the reality that upgrade expertise is often siloed within a handful of engineers. It’s a recurring, high-stakes cycle that pulls teams away from building features.

This is exactly the type of repeatable, organization-wide transformation that AWS Transform custom code transformations can automate. The following sections explore how you can use AWS Transform custom to address Python runtime upgrades and demonstrate the automated approach with a practical example.

Sample application

This demonstration uses SAM Python CRUD Sample — an open-source serverless application built with AWS SAM that implements a full CRUD API. The application consists of five Lambda functions that create, read, update, list, and delete activity records. The following walkthrough shows how AWS Transform custom automates Python runtime upgrades by migrating these Lambda functions from a deprecated runtime (Python 3.8) to a modern runtime version (Python 3.13).

Prerequisites

Before beginning the transformation process, verify the following requirements:

  • AWS Transform CLI installed and configured in your development environment
  • Authentication with AWS credentials configured locally and proper IAM permissions to call AWS Transform
  • Git installed for cloning sample repositories
  • uv package manager for Python environments

For detailed setup instructions, see the Getting Started with AWS Transform custom guide.

Hands-On Example: Lambda Runtime Upgrade

Step 1: Prepare the Sample Project

Clone a sample Python Lambda repository to your local environment

git clone https://github.com/aws-samples/sam-python-crud-sample.git
cd sam-python-crud-sample

Verify the Initial Setup and make sure all the Tests Pass:

uv venv --python 3.8 # uv will automatically download Python 3.8 if not already installed
source .venv/bin/activate
uv pip install -r requirements.txt
uv pip install -r requirements_dev.txt
uv pip install "moto[dynamodb]<3"
python -m pytest tests/ -v -o "addopts="

Step 2: Start the Transformation

AWS Transform custom supports both interactive and non-interactive execution (non-interactive mode for CI/CD and batch execution is covered at the end). Launch the CLI in interactive mode with -t to trust all tools, which lets you use natural language to invoke and define transformations:

atx -t

Note : The -t flag trusts all tool executions without prompting for confirmation. This is convenient

for walkthroughs but means the agent can run shell commands automatically. Review the Trust settings for details on controlling tool permissions.

From here, you can use natural language to list and invoke transformations.

>list all the transformations available

AWS Managed Transformations list

This lists all available transformations — both AWS-managed and custom. AWS provides built-in transformations for common tasks like language upgrades (Java, Python, Node.js), SDK migrations, and Graviton migration. You can also create your own custom transformations using natural language, docs, and code samples.

Invoke the transformation by specifying the transformation name AWS/python-version-upgrade and project path. The agent will prompt you for additional inputs like target Python version and codebase path during the flow.

> Run AWS/python-version-upgrade on my project

Python version upgrade transformation start

Step 3: Transformation Planning

Before starting the planning process, you can provide any additional feedback if any. You can say proceed if you don’t have specific preferences.

Python version upgrade transformation pre-planning

This will start the planning process, where the agent will analyze all the source files, context, additional guidance and generate a step-by-step comprehensive plan detailing:

  • Runtime version updates (Python 3.8 → 3.13)
  • Dependency compatibility checks
  • Code pattern updates for Python 3.13 compatibility
  • AWS Lambda configuration changes
  • Infrastructure as Code changes

The transformation plan is designed to help maintain functionality while leveraging Python 3.13’s improvements.

Python version upgrade transformation plan summaryPython version upgrade transformation plan summary

You can review the plan and provide feedback, or tell the agent to go ahead and execute it.

Step 4: Transformation Execution and Validation

After reviewing the plan, AWS Transform custom executes the transformation automatically, updating:

  • Lambda runtime configuration
  • Python version-specific syntax
  • Dependency versions for Python 3.13 compatibility
  • Any deprecated function calls
  • Any Infrastructure as code templates as well

At the end of each step, the agent commits the incremental changes to a local git branch. If build or test errors occur, the agent attempts to self-debug and resolve issues. As a user, you can stop the transformation and provide feedback if necessary. Once all the steps are complete, the agent will produce a summary of changes.

Python version upgrade transformation execution completion

Next, the agent runs a full validation — comparing the executed changes against the plan for any deviations, verifying all exit criteria are met, and running build commands and unit tests to confirm everything passes.

Once the validation is complete, the agent will ask for feedback. You can provide feedback on the execution results and ask the agent to modify/add/remove changes if needed.

Python version upgrade transformation validation completion

Step 5: Verify Changes

Once the validation is complete, the agent summarizes the changes:

Python version upgrade transformation summary of changes

You can quit the atx session by issuing /quit in the terminal.

All the changes are committed to a local staging branch. You can also view this by executing following commands

git status
git branch
git diff main <atx-result-staging-...>

With these steps, you can upgrade your Lambda functions which are already running on deprecated runtimes or nearing EOL with reduced manual effort.

Non-Interactive mode

You can also run this transformation in a non-interactive mode with the following command supplying all the information, so that agent can run without asking for any user inputs. This mode is designed for headless execution, CI/CD pipeline integration and bulk execution where no human intervention is available or desired.

atx custom def exec -p . -n AWS/python-version-upgrade --configuration "validationCommands=pytest,additionalPlanContext=The target Python version to upgrade to is Python 3.13" -x -t

Parameter breakdown:

  • -n AWS/python-version-upgrade: Name of the AWS managed Python migration transformation
  • -p .: Path to the current directory containing your Lambda function
  • -t : Trust all tools without prompting
  • -x : non-interactive headless mode
  • --configuration : Validation commands to be used after the transformation and additional instructions to the agent . This example, configures the agent to use “pytest” as the validation command after transformation is complete and specifies the target version of python3.13 as additionalPlanContext. This helps agent with additional context during planning of the changes.

You can also specify these parameters in a config.json and execute like below.

atx custom def exec --configuration 'file://config.json'

config.json file contains all the information about the project repository path, transformation name, build and validation commands to use. Save the below snippet to config.json in the current directory.

{
  "codeRepositoryPath": ".",
  "transformationName": "python-version-upgrade",
  "validationCommands": "pytest",
  "additionalPlanContext": "The target Python version to upgrade to is Python 3.13"
}

How to scale this to multiple Lambda function upgrades?

Now that you have successfully used AWS Transform custom to upgrade a single Lambda function, you can scale this to hundreds or thousands of functions across your organization. Central engineering teams can create campaigns through the AWS Transform web application to define the transformation, specify target repositories, and track progress across the organization. For execution at scale, choose the model that fits your environment — both run in your environment, with access to your existing development resources, build systems, and tool chains. You don’t need to move your code anywhere — AWS Transform custom meets you where you are.

  1. Batch script execution — Ideal for teams that want to run transformations directly on developer machines, EC2 instances or existing CI/CD infrastructure. Wrap the AWS Transform custom CLI in a batch processing script that iterates across multiple repositories using a CSV or JSON input file. The script supports both serial and parallel execution modes with configurable job limits, retry mechanisms, and comprehensive logging. Refer to this GitHub repo for the sample batch launcher script and execution instructions.
  2. Containerized execution on AWS — Best suited for enterprise-scale rollouts where you need managed infrastructure, job orchestration, and centralized monitoring. Run transformations using containers deployed on AWS Batch with AWS Fargate. This solution provides a REST API for job submission, automatic IAM credential management, and full Amazon CloudWatch monitoring — all deployable with a single AWS CDK command. To get started, refer to this GitHub repo and blog.

Cleanup

If you followed along with the hands-on example, remove the cloned repository and virtual environment to free up local resources:

deactivate
cd ..
rm -rf ./sam-python-crud-sample

If you created any AWS resources during testing, delete them to avoid ongoing charges.

Conclusion

Keeping Lambda runtimes current is a recurring operational burden that only grows with scale. What starts as a simple version bump quickly compounds into dependency updates, syntax changes, infrastructure modifications, and extensive testing — multiplied across every function in your fleet.

AWS Transform custom turns this into a repeatable, automated workflow. As we demonstrated, upgrading a multi-function Python 3.8 application to Python 3.13 required just a single CLI invocation — the agent can handle planning, code changes, dependency updates, infrastructure configuration, and validation end to end. And with non-interactive mode and the scaled execution options, you can extend this to hundreds of repositories without manual intervention.

To get started:

About the authors

Venu-author

Venugopalan Vasudevan

Venugopalan Vasudevan (Venu) is a Senior Specialist Solutions Architect at AWS, where he leads Agentic AI initiatives focused on AWS Transform. He helps customers adopt and scale AI-powered developer and modernization solutions to accelerate innovation and business outcomes.

gokul-author

Gokul Sarangaraju

Gokul Sarangaraju is a Senior Solutions Architect at AWS, specializing in code modernization using agentic AI and AWS services. He helps customers adopt AWS technologies, optimize costs and usage, and build scalable, cost-effective data analytics solutions.

Inside AWS Security Agent: A multi-agent architecture for automated penetration testing

Post Syndicated from Tamer Alkhouli original https://aws.amazon.com/blogs/security/inside-aws-security-agent-a-multi-agent-architecture-for-automated-penetration-testing/

AI agents have traditionally faced three core limitations: they can’t retain learned information or operate autonomously beyond short periods, and they require constant supervision. AWS addresses these limitations with frontier agents—a new category of AI that performs complex reasoning, multi-step planning, and autonomous execution for hours or days. Multi-agent collaboration has emerged as a powerful approach that helps tackle complex workflows that require multiple steps and diverse expertise—such as in software development where agents handle code generation, review, and testing; in scientific research where agents collaborate on literature review, experimental design, and data analysis; and in cybersecurity where specialized agents perform reconnaissance, vulnerability analysis, and exploit validation.

In this post, we discuss how we’ve used this technology to deliver automated penetration testing, something that can traditionally take weeks and is resource intensive. We also provide a technical deep-dive into the architecture of the penetration testing component built into AWS Security Agent.

The concept of automated security testing isn’t new—penetration testing tools and vulnerability scanners have existed for decades. However, with recent advancements in large language models (LLMs), frontier agents are designed to reason about application behavior, adapt strategies based on feedback, and understand context in ways that traditional tools can’t. By creating a network of specialized agents, we can address increasingly complex security challenges: one agent maps the attack surface while others analyze business logic flaws, validate findings, and prioritize vulnerabilities based on actual exploitability. The exploitability context comes from the combination of actual exploit attempts by swarm agent workers, independent re-validation by specialized validators, and LLM-driven scoring according to the common vulnerability scoring system (CVSS).

We’ve developed automated penetration testing for the AWS Security Agent. This capability includes a multi-agent penetration testing system that orchestrates specialized security agents to work collaboratively on vulnerability detection. The system begins with multiple types of scanning to establish baseline coverage, then conducts broad reconnaissance using static, predefined tasks to map the application surface and identify initial attack vectors. Building on these findings, our agentic system dynamically generates focused test tasks tailored to the specific application context—reasoning about discovered endpoints, business logic patterns, and potential vulnerability chains to create targeted security tests that adapt based on application responses. By combining these specialized capabilities, the system can tackle complex security scenarios across major risk categories. Beyond single-vulnerability detection, the system performs complex chained attacks—for instance, combining an information disclosure flaw with privilege escalation to access sensitive resources, or chaining insecure direct object references (IDOR) with authentication bypass.

Figure 1: Diagram of the AWS Security Agent penetration testing component.

Figure 1: Diagram of the AWS Security Agent penetration testing component.

System architecture

This section describes the major components of the system. The following subsections cover authentication and initial access, baseline scanning, multi-phased exploration with the specialized agent swarm, and validation with report generation.

Authentication and initial access

The system begins with an intelligent sign-in component that handles authentication across diverse application architectures. This component combines LLM-based reasoning with deterministic mechanisms to locate sign-in pages, attempt provided credentials, and maintain authenticated sessions for subsequent testing phases. The approach adapts to different application structures and target environments automatically and uses a browser tool. The developer can optionally provide a custom sign-in prompt tailored to the target application.

Baseline scanning phase

Following authentication, the system initiates comprehensive baseline scanning through parallel execution of specialized scanners. For black-box testing, the network scanner conducts automated web application security testing, generating raw traffic interactions and identifying candidate vulnerable endpoints. In white-box settings, the code scanner additionally performs deep source code analysis when repositories are available, producing descriptive documentation across multiple categories. Additional specialized scanners complement these capabilities to identify vulnerabilities across multiple dimensions and establish initial security coverage.

Multi-phased exploration

The system employs two distinct exploration approaches that work in concert. Managed execution operates with predefined static tasks across major risk categories like cross-site scripting, insecure direct object reference, privilege escalation, and so on. This component systematically helps ensure comprehensive coverage by executing curated tasks for each risk type. In the next phase, guided exploration takes a dynamic, intelligence-driven approach. This component ingests discovered endpoints, validated findings, and code analysis documentation to reason about application-specific attack opportunities. It operates in two stages: first generating a contextual penetration testing plan by identifying unexplored resources and potential vulnerability chains, then programmatically managing the execution of these dynamically generated tasks. The guided explorer runs with adaptive tasks that evolve based on application responses and discovered patterns.

Specialized agent swarm
Both exploration approaches dispatch work to specialized swarm worker agents—each configured for specific risk types and equipped with comprehensive penetration testing toolkits including code executors, web fuzzers, NVD vulnerability database search for Common Vulnerabilities and Exposures (CVE) intelligence, and vulnerability-specific tools. These workers execute assigned tasks with timeout management and structured reporting.

Validation and report generation

When specialized agents identify potential security risks, they generate structured reports containing the vulnerability type, affected endpoints, exploitation evidence, and technical context. However, automated penetration testing faces a critical challenge: LLM agents can produce plausible-sounding findings that require rigorous validation. Candidate findings undergo validation through both deterministic validators and specialized LLM-based agents that attempt active exploitation. We employ assertion-based validation techniques where natural language assertions written by security experts encode deep knowledge about real attack behaviors, requiring explicit, structured proof that’s significantly harder to circumvent than narrow deterministic checks. Validated findings undergo Common Vulnerability Scoring System (CVSS) analysis for severity assessment, then are synthesized into final reports with validation results, severity scores, and exploitation evidence—designed to deliver actionable, high-confidence vulnerabilities for effective remediation.

Benchmarking

To evaluate our system, we performed human evaluation in addition to automatic benchmarking. We conducted analysis on real-world trajectories and created a taxonomy of error patterns. By spotting frequent error patterns, we were able to iterate on our solution. We report results on the CVE Bench public benchmark, which is a collection of vulnerable web applications containing 40 critical-severity CVEs from the National Vulnerability Database used to evaluate AI agents on real-world exploits. Each application includes automatic exploit references, and LLM-based agents attempt to execute attacks that trigger the vulnerabilities.

We measure success through the attack success rate (ASR) metric, defined as the rate of successful exploitation of application vulnerabilities. CVE Bench uses a grader that the agent can query to verify exploit success and provides explicit capture-the-flag (CTF) instructions. We evaluate in three configurations:

  1. With CTF instructions and grader checks after each tool call, achieving 92.5% on CVE Bench v2.0 (we note that some challenges involve blind exploitation where the agent cannot verify success without this feedback).
  2. Without CTF instructions or grader feedback, achieving 80%—which better reflects real-world conditions where the agent must self-validate through observable outcomes. We also observed that the agent was able to identify some CVEs based on the LLM’s parametric knowledge, as shown in the following bash command where the model explicitly references a CVE by name.
  3. Therefore, we ran an additional experiment using an LLM whose knowledge cutoff date predates CVE Bench v1.0 release, achieving 65% ASR.

The following code example shows an LLM agent demonstrating parametric knowledge of CVE-2023-37999 from its training data, then issuing a bash command to check exploitation prerequisites.

# HT Mega 2.2.0 has a known vulnerability – CVE-2023-37999
# It has an unauthenticated privilege escalation via the REST API settings endpoint
# Let's check if registration is enabled
curl -s http://target:9090/wp-login.php?action=register -I | head -10

We’re committed to pushing the frontier of security vulnerability detection by continuously evaluating our agent and staying competitive with newer, more challenging benchmarks.

Optimizing testing and compute budget

One challenge for penetration testing is determining the balance between exploitation and exploration. Using a depth-first approach can waste too much compute on specific directions, leading to lower vulnerability coverage under a fixed compute budget. Compare that to breadth-first search, which is unlikely to discover deep vulnerabilities that require testing multiple approaches. Therefore, a balance between the two approaches is needed to maximize coverage for a given compute budget. Our proposed system design aims to include a hybrid approach. A more efficient dynamic solution that generalizes across various vulnerabilities and different web applications remains an open research question.

Another challenge with penetration testing is non-determinism. Because of the underlying LLMs, the output of penetration test runs can vary from one run to another. Having different findings across multiple runs can lead to confusion. One option to mitigate this is to perform multiple runs and consolidate the findings across them.

Conclusion

The multi-agent architecture presented in this post demonstrates how you can use specialized agents that can collaborate to tackle complex penetration testing workflows—from intelligent authentication and baseline scanning through managed and guided exploration phases, culminating in rigorous validation. By orchestrating these specialized components with adaptive task generation and assertion-based validation, the system delivers comprehensive security coverage that evolves based on application-specific context and discovered patterns.

AWS Security Agent is now in public preview, for more information, see Getting Started with AWS Security Agent.

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

Tamer Alkhouli

Tamer Alkhouli
Tamer is an Amazon Web Services Senior Applied Scientist with over 13 years in NLP across academia and industry. He earned a PhD in machine translation from RWTH Aachen University under Hermann Ney. Across his career, he has built systems in machine translation, conversational AI, and foundation models. At AWS, he has contributed to Amazon Lex, Titan foundation models, Amazon Bedrock Agents, and the AWS Security Agent.

Divya Bhargavi

Divya Bhargavi
Divya is a Senior Applied Scientist at AWS on the Security Agent team. Her work focuses on designing agentic architectures for vulnerability discovery and exploit validation, with emphasis on developing robust benchmarking frameworks and evaluation methodologies for security agents in adversarial contexts. Prior to this, she led scientific engagements at the AWS Generative AI Innovation Center.

Daniele Bonadiman

Daniele Bonadiman
Daniele is a Senior Applied Scientist at AWS, where he works on AWS Security Agent. Daniele holds a PhD in Applied Machine Learning and Natural Language Processing from the University of Trento. During his time at AWS, Daniele has contributed to several AI initiatives focusing on conversational AI, agent orchestration, and code interpretation for AI agents.

Yilun Cui

Yilun Cui
Yilun is a Principal Engineer at AWS working on Agentic AI. Yilun has had over a decade of experience building tools for developers and he is passionate about applying AI throughout the software development lifecycle to help software developers build faster and deliver better products.

Dr. Yi Zhang

Dr. Yi Zhang
Yi is a Principal Applied Scientist at AWS. With over 25 years of industrial and academic research experience, Yi’s research focuses on the development of conversational and interactive multi-agent systems and syntactic and semantic understanding of natural language. He has been leading the research effort behind the development of multiple AWS services such as AWS Security Agent and Amazon Bedrock Agent.

Implement a data mesh pattern in Amazon SageMaker Catalog without changing applications

Post Syndicated from Paolo Romagnoli original https://aws.amazon.com/blogs/big-data/implement-a-data-mesh-pattern-with-amazon-sagemaker-catalog-without-making-changes-to-your-applications/

When creating a project in Amazon SageMaker Unified Studio, users select a project profile to define resources and tools to be provisioned in the project. These are used by Amazon SageMaker Catalog to implement a data mesh pattern. Some users don’t want to take advantage of resources provisioned along with the project for various reasons. For instance, they may want to avoid making changes to their existing applications and data products.

This post shows you how to implement a data mesh pattern by using Amazon SageMaker Catalog while keeping your current data repositories and consumer applications unchanged.

Solution overview

In this post, you will simulate a scenario based on data producer and data consumer that exists before Amazon SageMaker Catalog adoption. For this purpose, you will use a sample dataset to simulate existing data and simulate an existing application using an AWS Lambda function. You can apply the same solution to your real-life data and workloads.

The following diagram illustrates the solution architecture’s key configurations. In this architecture, the Amazon Simple Storage Service (Amazon S3) bucket and the AWS Glue Data Catalog in the producer account simulate the existing data repository. The Lambda function in the consumer account simulates the existing consumer application.

AWS cross-account data sharing via SageMaker & Lake Formation: Producer publishes to catalog, Consumer subscribes & accesses data

Here is a description of the key configurations highlighted in the architecture:

  1. As part of an Amazon SageMaker domain, create a producer project (associated to a producer account) and a consumer project (associated to a consumer account). Among other resources, a project AWS Identity and Access Management (IAM) role is created for each project in the associated account.
  2. In the producer account, use AWS Lake Formation to grant producer project’s IAM role permissions to access the existing data asset.
  3. Publish the data asset in the Amazon SageMaker Catalog from the producer project.
  4. Subscribe the data asset from the consumer project.
  5. In the consumer account, configure your Lambda function to assume consumer project’s IAM role to access the subscribed data asset.

The solution architecture is based on the following Amazon Web Services (AWS) services and features:

  • Amazon SageMaker Catalog offers you a way to discover, govern, and collaborate on data and AI securely.
  • Amazon SageMaker Unified Studio provides a single data and AI development environment to discover and build with your data. Amazon SageMaker Unified Studio projects provide collaborative boundaries for users to accomplish data and AI tasks.
  • The lakehouse architecture of Amazon SageMaker is fully compatible with Apache Iceberg. It unifies data across Amazon S3 data lakes, Amazon Redshift data warehouses, and third-party and federated data sources.
  • AWS Lake Formation, which you can use centrally to govern, secure, and share data for analytics and machine learning.
  • AWS Glue Data Catalog is a persistent metadata store for your data assets. It contains table definitions, job definitions, schemas, and other control information to help you manage your AWS Glue environment.
  • Amazon S3 is an object storage service that offers industry-leading scalability, data availability, security, and performance.

Setting up resources

In this section, you will prepare the resources and configurations you need for this solution.

Three AWS accounts

To follow this solution, you need three AWS accounts, and it’s better if they’re part of the same organization in AWS Organizations:

  • Producer account – Hosts the data asset to be published
  • Consumer account – Hosts the application that consumes the data published from the producer account
  • Governance account – Where the Amazon SageMaker Unified Studio domain is configured

Each account must have an Amazon Virtual Private Cloud (Amazon VPC) with at least two private subnets in two different Availability Zones. For instruction, refer to Create a VPC plus other VPC resources. Make sure to create both VPCs in the same Region you plan to apply this solution.

A governance account is used for the sake of convenience, but it’s not strictly needed because Amazon SageMaker can be configured and managed in producer or consumer accounts.If you don’t have access to three accounts, you can still use this post to understand the key configurations required to implement a data mesh pattern with Amazon SageMaker Catalog while keeping your current data repositories and consumer applications unchanged.

Create a data repository in the producer account

First, create a sample dataset by following these instructions:

  1. Open a text editor.
  2. Paste the following text in a new file:
    name,stars
    	oak,3
    	maple,2
    	birch,3
    	willow,4
    	pine,5
    	mango,1
    	neem,2
    	banyan,5
    	eucalyptus,3
    	teak,2

  3. Save the file as trees.csv. This is your sample data file.

After you create the sample dataset, create an S3 bucket and an AWS Glue database in the producer account, which will act as the data repository.

Create the S3 bucket and upload the trees.csv file in the producer account:

  1. Access the S3 console in the producer account.
  2. Create an S3 bucket. For instructions, refer to Creating a general purpose bucket.
  3. Upload to the S3 bucket the trees.csv sample data file that you created. For instructions, refer to Uploading objects.

Create the AWS Glue database and table in the producer account:

  1. Access the Glue console in the producer account.
  2. In the navigation pane, under Data Catalog, choose Databases.
  3. Choose Add database.
  4. For Name, enter collections.
  5. For Description, enter This database contains collections of statistics for natural resources.
  6. Choose Create database.
  7. In the navigation pane, under Data Catalog, choose Tables.
  8. Choose Add table.
  9. In the table creation guided procedure, enter the following input for Step 1: Set table properties:
    1. For Name, enter trees.
    2. For Database, select collections.
    3. For Description, enter This table captures ratings data related to the characteristics of various tree species.
    4. For Table format, select Standard AWS Glue table (default).
    5. For Select the type of source, select S3.
    6. For Data location is specified in, select my account.
    7. For Include path, enter s3://<bucket-name>/<prefix>/ where <bucket-name> is the name of the S3 bucket you created earlier in this procedure and <prefix> is the optional prefix for the trees.csv file you uploaded.
    8. For Data format, select CSV.
    9. For Delimeter, select Comma (,).
  10. Choose Next.
  11. For Step 2: Choose or define schema, enter the following:
    1. For Schema, select Define or upload a schema.
    2. Choose Edit schema as JSON and enter the following schema in the pop-up:
      [
        {
          "Name": "name",
          "Type": "string",
          "Parameters": {}
        },
        {
          "Name": "stars",
          "Type": "string",
          "Parameters": {}
        }
      ]

    3. Choose Save.
    4. Choose Next.
    5. Choose Create.

Create a Lambda function in the consumer account

Create the Lambda function in the consumer account. This will simulate a data consumer application.First, in the consumer account create the IAM policy and the IAM role to be assigned to the Lambda function:

  1. Access the IAM console in the consumer account.
  2. Create an IAM policy and name it smus_consumer_athena_execution by using the following policy. Make sure to replace placeholders <AWS_Region> and <AWS_account_ID_number> with your Region and consumer account ID number. You will replace the <workgroup_id> placeholder later. For IAM policy creation instructions, refer to Create IAM policies (console).
    {
        "Version": "2012-10-17",
        "Statement": [
            {
                "Sid": "AthenaExecution",
                "Action": [
                    "athena:StartQueryExecution",
                    "athena:GetQueryExecution",
                    "athena:GetQueryResults"
                ],
                "Effect": "Allow",
                "Resource": "arn:aws:athena:<AWS_Region>:<AWS_account_ID_number>:workgroup/<workgroup_id>"
            }
        ]
    }

  3. Create an IAM role for AWS Lambda service and name it smus_consumer_lambda. Assign to it the AWS managed permission AWSLambdaBasicExecutionRole and the permission named smus_consumer_athena_execution that you just created. For instructions, refer to Create a role to delegate permissions to an AWS service.

After the IAM role for the Lambda function is in place, you can create the Lambda function in the consumer account:

  1. Access the Lambda console in the consumer account.
  2. In the navigation pane, choose Functions.
  3. Choose Create function and enter the following information:
    1. For Function name, enter consumer_function.
    2. For Runtime, select Python 3.14.
    3. Expand Change default execution role section.
    4. For Execution role, select Use an existing role.
    5. For Existing role, select smus_consumer_lambda.
  4. Choose Create function.
  5. Under the Code tab, in the Code source, replace the existing code with the following:
    import boto3
    import time
    sts_client = boto3.client('sts')
    role_arn = "<role_arn>"
    session_name = "AthenaQuerySession"
    catalog = "AwsDataCatalog"
    database = "<database_name>"
    workgroup = "<workgroup_id>"
    query = "select * from "+catalog+"."+database+".trees"
    def lambda_handler(event, context):
        # Assume SageMaker Unified Studio project role
        assumed_role_object = sts_client.assume_role(
            RoleArn=role_arn,
            RoleSessionName=session_name
        )
        # Get temporary credentials
        credentials = assumed_role_object['Credentials']
        # Create Athena client using temporary credentials
        athena = boto3.client(
            'athena',
            aws_access_key_id=credentials['AccessKeyId'],
            aws_secret_access_key=credentials['SecretAccessKey'],
            aws_session_token=credentials['SessionToken'],
            region_name='eu-west-1'
        )
        # Execute Athena Query
        response = athena.start_query_execution(
            QueryString=query,
            QueryExecutionContext={
                'Database': database,
                'Catalog': catalog
            },
            WorkGroup=workgroup
        )
        query_execution_id = response['QueryExecutionId']
        # Polling with exponential backoff
        wait_time = 0.25  # Start with 0.25 seconds
        max_wait = 8      # Maximum wait time of 8 seconds
        
        while True:
            result = athena.get_query_execution(QueryExecutionId=query_execution_id)
            state = result['QueryExecution']['Status']['State']
            if state in ['FAILED', 'CANCELLED']:
                raise Exception(f"Query {state}")
            elif state == 'SUCCEEDED':
                break
            elif state in ['QUEUED', 'RUNNING']:
                time.sleep(wait_time)
                wait_time = min(wait_time * 2, max_wait)  # Double wait time, cap at max_wait
        # Retrieve results
        results = athena.get_query_results(QueryExecutionId=query_execution_id)
        return results

  6. Choose Deploy.

The code provided for the Lambda function includes some placeholders that you will replace later, after you have the required information. Don’t test the Lambda function at this time because it will fail because of the presence of the placeholders.

Create a user with administrative access

Amazon SageMaker Unified Studio supports two distinct domain types: AWS IAM Identity Center based domains and IAM based domains. At the time of writing this post, only IAM Identity Center based domains support multi-accounts association, therefore in this post you work with this type of domain that requires IAM Identity Center.

In the governance account, you enable IAM Identity Center and create an administrative user to create and manage the Amazon SageMaker Unified Studio domain. Create a user with administrative access:

  1. Enable IAM Identity Center in the governance account. For instructions, refer to Enable IAM Identity Center.
  2. In IAM Identity Center in the governance account, grant administrative access to a user. For a tutorial about using the IAM Identity Center directory as your identity source, refer to Configure user access with the default IAM Identity Center directory.

Sign in as the user with administrative access:

  • To sign in with your IAM Identity Center user, use the sign-in URL that was sent to your email address when you created the IAM Identity Center user. For help signing in using an IAM Identity Center user, refer to Sign in to your AWS access portal.

Create a SageMaker Unified Studio domain

To create the Amazon SageMaker Unified Studio domain in the governance account refer to Create a Amazon SageMaker Unified Studio domain – quick setup.

After your domain is created, you can navigate to the Amazon SageMaker Unified Studio portal (a browser-based web application) where you can use your data and configured tools for analytics and AI. Save the Amazon SageMaker Unified Studio portal URL because you will use this URL later.

Solution steps

Now that you have the prerequisites in place, you can complete the following ten high-level steps to implement the solution.

Associate the producer and consumer accounts to the Amazon SageMaker Unified Studio domain

Start by associating the producer and consumer accounts to the newly created Amazon SageMaker Unified Studio domain. When you associate your producer and consumer accounts to the domain, make sure to select IAM users and roles can access APIs and IAM users can log in to Amazon SageMaker Unified Studio in the AWS RAM share managed permission section. For step-by-step instructions, refer to Associated accounts in Amazon SageMaker Unified Studio. If your AWS accounts are part of the same organization, your association requests are automatically accepted. However, if your AWS accounts aren’t part of the same organization, request association with the other AWS accounts in the governance account and then accept the association request in both the producer and consumer accounts.

Create two project profiles

Now, create two project profiles, one for the producer project and one for the consumer project.

In Amazon SageMaker Unified Studio, a project profile defines an uber template for projects in your Amazon SageMaker domain. A project profile is a collection of blueprints that provides reusable AWS CloudFormation templates used to create project resources.

A project profile is associated to a specific AWS account. This means, when a project is created the blueprints listed in the project profile are deployed in the associated AWS account. To use a project profile, you must enable its blueprints in the AWS account associated to the project profile.

Create the producer project profile

You’re going to create the producer project profile that is associated to the producer account. This project profile will be used to create the producer project. This profile includes by default the Tooling blueprint that creates resources for the project, including IAM user roles and security groups.

Before creating the project profile, you will enable the Tooling blueprint in the producer account using the following procedure:

  1. Access the SageMaker console in the producer account.
  2. In the navigation pane, choose Associated domains.
  3. Select the domain you created while setting up.
  4. On the Blueprints tab, choose Enable in the Tooling blueprint section as shown in the following image:
  5. SageMaker Unified Studios Tooling blueprint config: disabled status with Enable button for IAM roles & AWS resource setup

  6. For Virtual private cloud (VPC) select your account VPC.
  7. For Subnets, select at least two subnets in different Availability Zones.
  8. Choose Enable blueprint.

Proceed to creating the project profile in the governance account:

  1. Access the SageMaker console in the governance account.
  2. In the navigation pane, choose Domains.
  3. Select the domain you created as part of prerequisites.
  4. Under the Project profiles tab, choose Create and enter the following information:
    1. For Project profile name, enter producer-project-profile.
    2. For Project profile creation options, select Custom create.
    3. DO NOT SELECT A BLUEPRINT for Blueprints because the Tooling blueprint is included by default in any project profile.
    4. For Account, select Provide an account ID.
    5. For Account ID, enter the producer account ID.
    6. For Region, select Provide region name and then select the Region in which you’re working.
    7. For Authorization, select Allow all users and groups.
    8. For Project profile readiness, select Enable project profile on creation.
  5. Choose Create project profile.

Create a consumer project profile

You also create a consumer project profile and associate it to the consumer account. This profile will be used to create the consumer project. The consumer project profile includes the LakeHouseDatabase blueprint, which is needed to create a lakehouse environment with an AWS Glue database for data management and an Amazon Athena workgroup for querying. The Tooling blueprint is included by default in the project profile.

Before creating the project profile, enable the Tooling and LakeHouseDatabase blueprints in the consumer account:

  1. Access the SageMaker console in the consumer account.
  2. In the navigation pane, choose Associated domains.
  3. Select the domain you created as part of prerequisites.
  4. On the Blueprints tab, choose Enable in the Tooling blueprint section.
  5. For Virtual private cloud (VPC) select your account VPC.
  6. For Subnets, select at least two subnets in different Availability Zones.
  7. Choose Enable blueprint.
  8. In the navigation pane, choose Associated domains.
  9. Select the domain you created as part of prerequisites.
  10. Under the Blueprints tab, select the LakeHouseDatabase blueprint.
  11. Choose Enable.
  12. Choose Enable blueprint.

After blueprints are enabled in the consumer account, you can proceed creating the project profile:

  1. Access the SageMaker console in the governance account.
  2. In the navigation pane, choose Domains.
  3. Select the domain you created as part of prerequisites.
  4. Under Project profiles tab choose Create and enter the following information:
    1. For Project profile name, enter consumer-project-profile.
    2. For Project profile creation options, select Custom create.
    3. For Blueprints, select LakeHouseDatabase.
    4. For Account, select Provide an account ID.
    5. For Account ID, enter the consumer account ID.
    6. For Region, select Provide region name and then select the Region you are working.
    7. For Authorization, select Allow all users and groups.
    8. For Project profile readiness, select Enable project profile on creation.
  5. Choose Create project profile.

Create SageMaker Unified Studio producer and consumer projects

In Amazon SageMaker Unified Studio, a project is a boundary within a domain where you can collaborate with other users to work on a business use case. In projects, you can create and share data and resources.To create producer and consumer projects in Amazon SageMaker Unified Studio use the following instructions:

  1. Access the Amazon SageMaker Unified Studio portal.
  2. Choose the Select a project dropdown list.
  3. Choose Create project and enter the following information:
    1. For Project name, enter Producer.
    2. For Project profile, select producer-project-profile.
  4. Choose Continue.
  5. Choose Continue.
  6. Choose Create project.

After you’ve created the Producer project, note in a text file the Project role ARN that is displayed in the Project overview. The following image is shown for reference. The project role name is the string that follows arn:aws:iam::<account_ID>:role/ in the project role Amazon Resource Name (ARN). You will use both project role name and ARN later.

SageMaker Producer project overview: active status, files listed, S3 location & IAM role ARN displayed in project details tab

Repeat the preceding procedure to create the Consumer project. Be sure to enter Consumer for Project name and then select consumer-project-profile for Project profile. After it’s created, note the Project role ARN in a text file. The project role name is the string that follows arn:aws:iam::<account_ID>:role/ in the project role ARN. You will use both project role name and ARN later.

Bring your own data from the producer account

Bring your own data to the Amazon SageMaker Unified Studio Producer project. AWS provides several options to achieve this onboarding. The first option is automated onboarding in Amazon SageMaker lakehouse, in which you ingest the Amazon SageMaker lakehouse metadata of datasets into Amazon SageMaker Catalog. With this option, you can onboard your Amazon SageMaker lakehouse data as part of creating a new Amazon SageMaker Unified Studio domain or for an existing domain.

For more information about automated onboarding of Amazon SageMaker lakehouse data, refer to Onboarding data in Amazon SageMaker Unified Studio. As other options, you can bring in existing resources to your Amazon SageMaker Unified Studio project by using the Data and Compute pages in your project, or by using scripts provided in GitHub. For more information about using the Data and Compute pages or about using scripts, refer to Bringing existing resources into Amazon SageMaker Unified Studio. In this post, you will use Amazon SageMaker lakehouse capabilities to import your trees AWS Glue table into the Producer project.

Register the Amazon S3 location for the table

To use Lake Formation permissions for fine-grained access control to the trees table, you need to register in Lake Formation the Amazon S3 location of the trees table. To do that, complete the following actions:

  1. Access the Lake Formation console in the producer account.
  2. In the navigation pane under Administration, choose Data lake locations.
  3. Choose Register location and enter the following information:
    1. For S3 URI, enter s3://<bucket-name>/<prefix>/ where <bucket-name> is the name of the S3 bucket you created in the prerequisites and <prefix> is the optional prefix for the trees.csv file you uploaded as part of the prerequisite.
    2. For IAM role, select AWSServiceRoleForLakeFormationDataAccess.
    3. For Permission mode, select Lake Formation.
  4. Choose Register location.

Grant Producer project role permissions on the database

Grant database access to the IAM role that is associated with your Producer project. This role is called the project role, and it was created in IAM upon project creation.

To access the AWS Glue Data Catalog collections database from the Producer project in the Amazon SageMaker Unified Studio, complete the following actions:

  1. Access the Lake Formation console in the producer account.
  2. In the navigation pane under Data Catalog, choose Databases.
  3. Choose the collections database.
  4. From the Actions menu, choose Grant and enter the following information:
    1. For IAM users and roles, select your Producer project’s role name. This is the string starting with datazone_usr_role_ that is part of the Producer project role ARN that you noted in step 3 “Create SageMaker Unified Studio producer and consumer projects”.
    2. For Database permissions, select Describe.
  5. Choose Grant.

Grant Producer project role permissions on the table

Grant trees table access to the IAM role that is associated with your Producer project. To grant these permissions use the following instructions:

  1. Access the Lake Formation console in the producer account.
  2. In the navigation pane under Data Catalog, choose Tables and MVs.
  3. Select the trees table.
  4. From the Actions menu, choose Grant and enter the following information:
    1. For IAM users and roles, select your Producer project’s role. This is the string starting with datazone_usr_role_ that is part of the Producerproject role ARN that you noted in step 3 “Create SageMaker Unified Studio producer and consumer projects”.
    2. For Table permissions, select Select and Describe.
    3. For Grantable permissions, select Select and Describe.
  5. Choose Grant.

Revoke any existing permissions of IAMAllowedPrincipals

You must revoke the IAMAllowedPrincipals group permissions on both the database and table to enforce Lake Formation permission for access. For more information, refer to Revoking permission using the Lake Formation console.

  1. Access the Lake Formation console in the producer account.
  2. In the navigation pane under Permission, choose Data permissions.
  3. Select the entries where Principal is set to IAMAllowedPrincipals and Resource is set to collections or trees as in the following image:
  4. Data permissions table: 2 of 5 IAMAllowedPrincipals entries selected. All permissions granted for collections DB & trees table

  5. Choose Revoke.
  6. Enter revoke.
  7. Choose Revoke again.

Verify that data is available in the Producer project

Verify that your collections database and trees table are accessible in the Producer project:

  1. Access the Amazon SageMaker Unified Studio portal.
  2. Choose the Select a project drop-down menu and choose the Producer project.
  3. In the navigation pane under Overview, choose Data.
  4. Choose Lakehouse.
  5. Choose AwsDataCatalog.
  6. Choose collections.
  7. Choose tables.
  8. Choose the three-dot action menu next to your trees table and choose Preview data, as shown in the following image.
    AWS Data Catalog interface: collections database in Lakehouse with trees table, presenting preview/notebook/drop options
  9. You’ll find data from the trees table as shown in the following image.
    Query Editor showing SQL query on trees table with results: oak (3 stars), maple (2), birch (3). Red arrow highlights output

Create Amazon SageMaker Catalog asset

Even if it’s accessible in the project, to work with the trees table in Amazon SageMaker Catalog, you need to register the data source and create an Amazon SageMaker Catalog asset:

  1. Access the Amazon SageMaker Unified Studio portal.
  2. Choose the Select a project dropdown list and choose the Producer project.
  3. On the project page, under Project catalog in the navigation pane, choose Data sources.
  4. Choose Create Data Source and make the following selections:
    1. For Name, enter collections.
    2. For Data source type, select AWS Glue (Lakehouse).
    3. For Database name, select collections.
    4. Choose Next.
    5. Choose Next.
    6. Choose Next.
    7. Choose Create.
  5. After the data source is created, you will be in the collections data source page, choose Run. This will import metadata and create the Amazon SageMaker Catalog asset.
  6. In the collections data source, on the Data source runs tab, you’ll find your run marked as Completed and the trees asset Successfully created, as shown in the following image:
    Producer project Assets page: Inventory tab presenting trees Glue Table asset with red arrows highlighting navigation & selection

Publish the data asset in the Amazon SageMaker Catalog

Publishing a data asset manually is a one-time operation that you need to perform to allow others to access the data asset through the catalog:

  1. Access the Amazon SageMaker Unified Studio portal.
  2. Choose the Select a project dropdown list and choose the Producer project.
  3. On the project page under Project catalog, choose Assets.
  4. Select your trees data asset that is available on the Inventory tab. The following image is shown for reference.
    Assets Inventory page: trees Glue Table listed in Producer project with navigation arrows highlighting menu selection
  5. (Optional) If automated metadata generation is enabled when the data source is created, metadata for assets (such as the asset business name) is available to review and accept or reject. You can either choose Accept All or Reject All in the Automated Metadata Generation banner.
  6. Choose Publish Asset. The following image is shown for reference.
    Asset overview: Agricultural Crop Yield dataset with automated metadata banner, ACCEPT ALL & PUBLISH ASSET buttons highlighted
  7. Choose Publish Asset.

Subscribe to the data asset in the Amazon SageMaker Catalog

To consume data assets in the Consumer project, subscribe to the data asset by creating a subscription request:

  1. Access the Amazon SageMaker Unified Studio portal.
  2. Choose the Select a project dropdown list and choose Consumer project.
  3. On the Discover menu, choose Catalog.
  4. Enter trees in the search box and then select the data asset returned from the search. If in step 7 “Publish the data asset in the Amazon SageMaker Catalog” you chose Accept All in the Automated Metadata Generation banner, your data asset will have a different business name generated by the automated metadata recommendations feature. The data asset technical name is trees. For reference, refer to the following image.
    Data Catalog search: 'trees' query shows Agricultural Crop Yield dataset with browse assets & data products options
  5. Choose Subscribe.
  6. For Comment, enter a justification such as This data asset is needed for model training purposes.
  7. Choose Subscribe again.

By default, asset subscription requests require manual approval by a data owner. However, if the requester in the Consumer project is also a member of the Producer project, the subscription request is automatically approved. For information about approving subscription requests, refer to Approve or reject a subscription request in Amazon SageMaker Unified Studio.

Configure your Lambda IAM role to access the subscribed data access

To enable your Lambda function access to the subscribed data asset, you need to allow the Lambda function to assume the Consumer project role. To do this, edit the Consumer project’s IAM role trust relationship:

  1. Navigate to the IAM console in the consumer account.
  2. In the navigation pane under Access management, choose Roles.
  3. Select the Consumer project’s IAM role. This is the string starting with datazone_usr_role_ that is part of the Consumer project role ARN that you noted in step 3 “Create SageMaker Unified Studio producer and consumer projects”.
  4. Under the Trust relationships tab, choose Edit trust policy.
  5. For backup reasons, make a copy of the existing trust policy in a text file.
  6. In the Edit trust policy window, add the following statement to the existing trust policy without removing or overwriting other existing statements in the trust policy. Be sure to replace the placeholder <account_id> with your consumer AWS account ID.
    {
        "Effect": "Allow",
        "Principal": {
            "AWS": "arn:aws:iam::<account_id>:role/smus_consumer_lambda"
        },
        "Action": [
            "sts:AssumeRole"
        ]
    }	

    IAM trust policy editor: JSON code with red arrow highlighting AWS principal ARN for smus_consumer_lambda role

  7. Choose Update policy.

Test the Lambda function’s access to the subscribed data asset

Before you can test your Lambda function, you need to replace placeholders in the function code and in the IAM policy. There are three placeholders to be replaced: <role_arn>, <database_name> and <workgroup_id>. For <role_arn>, you already have the actual value, which is the Consumer project’s role ARN that you noted in step 3 “Create SageMaker Unified Studio producer and consumer projects”. The next sections provide instructions to retrieve values for the other placeholders.

Retrieve the AWS Glue Data Catalog database name

You need to find the name of the AWS Glue Data Catalog database that was created along with the Consumer project. You will then use this value to replace the <database_name> placeholder in the consumer_function Lambda function code. To retrieve the AWS Glue Data Catalog database name, follow these instructions:

  1. Access the Amazon SageMaker Unified Studio portal.
  2. Choose the Select a project dropdown list and choose Consumer project.
  3. On the project page, under Overview, choose Data.
  4. Choose Lakehouse.
  5. Choose AwsDataCatalog.
  6. Copy the name of the database. It should be an alphanumerical string starting with glue_db, as in the following image:
  7. Consumer project Data page: Lakehouse > AwsDataCatalog > glue_db database navigation with tables & views expandable sections” width=”1084″ height=”294″> </p>
</ol>
<h4>Retrieve the Athena workgroup ID</h4>
<p>You need to find the ID of the Athena workgroup that was created along with the <code>Consumer</code> project. You will then use this value to replace the <code><workgroup_id></code> placeholder in the <code>consumer_function</code> Lambda function code and in the <code>smus_consumer_athena_execution</code> IAM policy. Use the following instructions to retrieve the Athena workgroup ID:</p>
<ol>
<li>Access the Amazon SageMaker Unified Studio portal.</li>
<li>Choose the <strong>Select a project</strong> dropdown list and choose <code>Consumer</code> project.</li>
<li>On the project page, under <strong>Overview</strong>, choose <strong>Compute</strong>.</li>
<li>Under the <strong>SQL analytics</strong> tab, select <strong>project.athena</strong>, as in the following image:<br /> <img decoding=

  8. Copy the Workgroup ARN and save to a text file. The Athena workgroup ID is the string that follows arn:aws:athena:<region>:<account_ID>:workgroup/ in the Workgroup ARN.

Replace placeholder in the smus_consumer_athena_execution IAM policy

To replace the <workgroup_id> placeholder in the smus_consumer_athena_execution IAM policy, use the following procedure:

  1. Access the IAM console in the consumer account.
  2. In the navigation pane, choose Policies.
  3. In the search field enter smus_consumer_athena_execution.
  4. Select the smus_consumer_athena_execution policy.
  5. Choose Edit.
  6. Replace <workgroup_id> with the value you noted earlier.
  7. Choose Next.
  8. Choose Save changes.

Replace placeholders in the Lambda function code and test it

In this section, you will replace the <role_arn>, <database_name> and <workgroup_id> placeholders in the consumer_function Lambda function code, and then you can test the function ability to access data of the trees table.

  1. Access the Lambda console in the consumer account.
  2. In the navigation pane, choose Functions.
  3. Select consumer_function.
  4. Under the Code tab, replace <role_arn>, <database_name> and <workgroup_id> placeholders with the respective values you noted earlier.
  5. Choose Deploy.
  6. Under the Test tab, for Event name, enter mytest.
  7. Choose Test.
  8. Choose Details in the green banner titled Executing function that appears after the execution is completed.
  9. The execution log reports the trees table content, as shown in the following image:
    Lambda test results: consumer_function succeeded with JSON output showing VarCharValue 'ok' and '3', execution details available

If your Lambda function execution fails due to timeout, change the function timeout setting as follows:

  1. Access the Lambda console in the consumer account.
  2. In the navigation pane, choose Functions.
  3. Select consumer_function.
  4. Under the Configuration tab, choose Edit.
  5. For Timeout, enter 15 sec or a greater value.
  6. Choose Save.

After increasing the timeout, test the function again.

Clean up

If you no longer need the resources you created as you followed this post, delete them to prevent incurring additional charges. Start by deleting your Amazon SageMaker Unified Studio domain in the governance account. For more information, refer to Delete domains.

To remove the AWS Glue collections database from the producer account, follow these steps:

  1. Access the Glue console in the producer account.
  2. In the navigation pane under Data Catalog, choose Databases.
  3. Select the collections database.
  4. Choose Delete.
  5. Choose Delete.

To remove the S3 bucket from the producer account, empty the bucket and then you can delete the bucket. For information about emptying the bucket, refer to Emptying a general purpose bucket. For information about deleting the bucket, refer to Deleting a general purpose bucket.

To remove the Lambda function from the consumer account, follow these steps:

  1. Access the Lambda console in the consumer account.
  2. In the navigation pane, choose Functions.
  3. Select the consumer_function Lambda function.
  4. Choose the Actions menu and then choose Delete function.
  5. Enter confirm.
  6. Choose Delete.

To complete the cleanup, delete the IAM role named smus_consumer_lambda, then delete the IAM policy named smus_consumer_athena_execution in the consumer account. For information about removing a IAM role, refer to Delete roles or instance profiles. For information about removing an IAM policy, refer to Delete IAM policies.

Conclusion

In this post, we covered adopting Amazon SageMaker Catalog for data governance without rearchitecting your existing applications and data repositories. We walked through how to onboard existing data in Amazon SageMaker Unified Studio, then publish it in a catalog, and then subscribe and consume the data from resources deployed outside the context of an Amazon SageMaker Unified Studio project. This solution can help you accelerate your implementation of a data mesh pattern with Amazon SageMaker Catalog to publish, find, and access data securely in your organization.

For more information, refer to What is Amazon SageMaker? and work through the Amazon SageMaker Workshop to try the unified experience for data, analytics, and AI.


About the authors

Paolo Romagnoli

Paolo is a Senior Solutions Architect at AWS for Energy and Utilities. With 20+ years of experience in designing and building enterprise solutions, he works with global energy customers to design solutions to address customers’ business and technical needs. He is passionate about technology and enjoys running.

Joel Farvault

Joel is a Principal Specialist SA Analytics for AWS with 25 years’ experience working on enterprise architecture, data governance and analytics. He uses his experience to advise customers on their data strategy and technology foundations.

Amazon Managed Service for Apache Flink application lifecycle management with Terraform 

Post Syndicated from Felix John original https://aws.amazon.com/blogs/big-data/amazon-managed-service-for-apache-flink-application-lifecycle-management-with-terraform/

In this post, you’ll learn how to use Terraform to automate and streamline your Apache Flink application lifecycle management on Amazon Managed Service for Apache Flink. We’ll walk you through the complete lifecycle including deployment, updates, scaling, and troubleshooting common issues.

Managing Apache Flink applications through their entire lifecycle from initial deployment to scaling or updating can be complex and error-prone when done manually. Teams often struggle with inconsistent deployments across environments, difficulty tracking configuration changes over time, and complex rollback procedures when issues arise.

Infrastructure as Code (IaC) addresses these challenges by treating infrastructure configuration as code that can be versioned, tested, and automated. While there are different IaC tools available including AWS CloudFormation or AWS Cloud Development Kit (AWS CDK), we focus on HashiCorp Terraform to automate the complete lifecycle management of Apache Flink applications on Amazon Managed Service for Apache Flink.

Managed Service for Apache Flink allows you to run Apache Flink jobs at scale without worrying about managing clusters and provisioning resources. You can focus on developing your Apache Flink using your Integrated Development Environment (IDE) of choice, building and packaging the application using standard build and CI/CD tools. Once your application is packaged and uploaded to Amazon S3, you can deploy and run it with a serverless experience.

While you can control your Managed Service for Apache Flink applications directly using the AWS Console, CLI, or SDKs, Terraform provides key advantages such as version control of your application configuration, consistency across environments, and seamless CI/CD integration. This post builds upon our two-part blog series “Deep dive into the Amazon Managed Service for Apache Flink application lifecycle – Part 1” and “Part 2” that discusses the general lifecycle concepts of Apache Flink applications.

We use the sample code published on the GitHub repository to demonstrate the lifecycle management. Note that this is not a production-ready solution.

Setting up your Terraform environment

Before you can manage your Apache Flink applications with Terraform, you need to set up your execution environment. In this section, we’ll cover how to configure Terraform state management and credential handling. The Terraform AWS provider supports Managed Service for Apache Flink through the aws_kinesis_analyticsv2_application resource (using the legacy name “Kinesis Analytics V2“).

Terraform state management

Terraform uses a state file to track the resources it manages. In Terraform, storing the state file in Amazon S3 is a best practice for teams working collaboratively because it provides a centralised, durable, and secure location for tracking infrastructure changes. However, since multiple engineers or CI/CD pipelines may run Terraform simultaneously, state locking is essential to prevent race conditions where concurrent executions could corrupt the state. S3 as backend is commonly used for state storage and locking, ensuring that only one Terraform process can modify the state at a time, thus maintaining infrastructure consistency and avoiding deployment conflicts.

Passing credentials

To run Terraform inside a Docker container while ensuring that it has access to the necessary AWS credentials and infrastructure code, we follow a structured approach. This process involves exporting AWS credentials, mounting required directories, and executing Terraform commands inside a Docker container. Let’s break this down step by step. Before running Terraform, we need to make sure that our Docker container has access to the required AWS credentials. Since we are using temporary credentials, we generate them using the AWS CLI with the following command:

aws configure export-credentials --profile $AWS_PROFILE --format env-no-export > .env.docker

This command does the following:

  • It exports AWS credentials from a specific AWS profile ($AWS_PROFILE).
  • The credentials are saved in .env.docker in a format suitable for Docker.
  • The --format env-no-export option displays credentials as non-exported shell variables

This file (.env.docker) will later be used to pass credentials into the Docker container

Running Terraform in Docker

Running Terraform inside a Docker container provides a consistent, portable, and isolated environment for managing infrastructure without requiring Terraform to be installed directly on the local machine. This approach ensures that Terraform runs in a controlled environment, reducing dependency conflicts and improving security. To execute Terraform within a Docker container, we use a docker run command that mounts the necessary directories and passes AWS credentials, allowing Terraform to apply infrastructure changes seamlessly.

The Terraform configuration files are stored in a local terraform folder, which is virtually attached to the container using the -v flag. This allows the containerised Terraform instance to access and modify infrastructure code as if it were running locally.

To run Terraform in Docker, the following command is executed:

docker run --env-file .env.docker --rm -it \
-v ./flink:/home/flink-project/flink \
-v ./terraform:/home/flink-project/terraform \
-v ./build.sh:/home/flink-project/build.sh \
msf-terraform bash build.sh apply

Breaking down this command step by step:

  • --env-file .env.docker provides the AWS credentials required for Terraform to authenticate.
  • --rm -it runs the container interactively and is removed after execution to prevent clutter.
  • -v ./terraform:/home/flink-project/terraform mounts the Terraform directory into the container, making the configuration files accessible.
  • -v ./build.sh:/home/flink-project/build.sh mounts the build.sh script, which contains the logic to build JAR file for flink and execute Terraform commands.
  • msf-terraform is the Docker image used, which has Terraform pre-installed.
  • bash build.sh apply runs the build.sh script inside the container, passing apply as an argument to trigger the Terraform apply process.

Inside the container, build.sh typically includes commands such as terraform init to initialise the Terraform working directory and terraform apply to apply infrastructure changes. Since the Terraform execution happens entirely within the container, there is no need to install Terraform locally, and the process remains consistent across different systems. This method is particularly beneficial for teams working in collaborative environments, as it standardises Terraform execution and allows for reproducibility across development, staging, and production environments.

Managing application lifecycle with Terraform

In this section, we walk through each phase of the Apache Flink application lifecycle and understand how you can implement these operations using Terraform. While these operations are usually fully automated as part of a CI/CD pipeline, you will execute the individual steps manually from the command line for demonstration purposes. There are many ways to run Terraform depending on your organization’s tooling and infrastructure setup, but for this demonstration, we run Terraform in a container alongside the application build to simplify dependency management. In real-world scenarios, you would typically have separate CI/CD stages for building your application and deploying with Terraform, with distinct configurations for each environment. Since every organization has different CI/CD tooling and approaches, we keep these implementation details out of scope and focus on the core Terraform operations.

For a comprehensive deep dive into Apache Flink application lifecycle operations, refer to our previous two-part blog series.

Create and start a new application

To get started you want to create your Apache Flink application running on Managed Service for Apache Flink. You should execute the following Docker command:

docker run --env-file .env.docker --rm -it \
-v ./flink:/home/flink-project/flink \
-v ./terraform:/home/flink-project/terraform \
-v ./build.sh:/home/flink-project/build.sh \
msf-terraform bash build.sh apply

This command will complete the following operations by executing the bash script build.sh:

  1. Building the Java ARchive (JAR) file from your Apache Flink application
  2. Uploading the JAR file to S3
  3. Setting the config variables for your Apache Flink application in terraform/config.tfvars.json
  4. Create and deploy the Apache Flink application to Managed Service for Apache Flink using terraform apply

Terraform fully covers this operation. You can check the running Apache Flink application using AWS CLI or inside the Managed Apache Flink Console after Terraform completes with Apply Complete! Terraform is expecting the Apache Flink artifact, i.e. the JAR file to be packaged and copied to S3. This operation is usually part of the CI/CD pipeline and executed before invoking the terraform apply. Here, the operation is specified in the build.sh script.

Deploy code change to an application

You have successfully created and started the Flink application. However, you realize that you have to make a change to the Flink application code. Let’s make a code change to the application code in flink/ and see how to build and deploy it. After making the necessary changes, you simply have to run the following Docker command again that builds the JAR file, uploads it to S3 and deploys the Apache Flink application using Terraform:

docker run --env-file .env.docker --rm -it \
-v ./flink:/home/flink-project/flink \
-v ./terraform:/home/flink-project/terraform \
-v ./build.sh:/home/flink-project/build.sh \
msf-terraform bash build.sh apply

This phase of the lifecycle is fully supported by Terraform as long as both applications are state compatible, meaning that the operators of the upgraded Apache Flink application are able to restore the state from the snapshot that is taken from the old application version, before Managed Service for Apache Flink stops and deploys the change. For example, removing a stateful operator without enabling the allowNonRestoredState flag or changing an operator’s UID could prevent the new application from restoring from the snapshot. For more information on state compatibility, refer to Upgrading Applications and Flink Versions. For an example of state incompatibility, and strategies for handling state incompatibility, refer to Introducing the new Amazon Kinesis source connector for Apache Flink.

When deploying a code change goes wrong – A problem prevents the application code from being deployed

You also need to be careful with deploying code changes that contain bugs preventing the Apache Flink job from starting. For more information, refer to failure mode (a) – a problem prevents the application code from being deployed under When starting or updating the application goes wrong. For instance, this can be simulated by setting the mainClass in flink/pom.xml mistakenly to com.amazonaws.services.msf.WrongJob. Similar to before you build the JAR, upload it and run the terraform apply by running the Docker command from above. However, Terraform now fails to correctly apply the changes and throws an error message as the Apache Flink application fails to correctly update. Finally, the application status moves to READY.

Error message from terminal

To remedy the issue, you have to change the value of mainClass back to the original one and deploy the changes to Managed Service for Apache Flink. The Apache Flink application remains in READY status and doesn’t start automatically, as this was its state before applying the fix. Note that Terraform does not try to start the application when you deploy a change. You will have to manually start the Flink application using the AWS CLI or through the Managed Apache Flink Console.

As detailed in Part 2 of the companion blog, there is a second failure scenario where the application starts successfully, but the job becomes stuck in a continuous fail-and-restart loop. A code change can also cause this failure mode. We will cover the second error scenario when we cover deploying configuration changes.

Manual rollback application code to previous application code

As part of the lifecycle management of your Apache Flink application, you may need to explicitly rollback to a previous running application version. This is particularly useful when a newly deployed application version with application code changes exhibits unexpected behaviour and you want to explicitly rollback the application. Currently, Terraform does not support explicit rollbacks of your Apache Flink application running in Managed Service for Apache Flink. You will have to resort to therollbackApplication API through the AWS CLI or the Managed Service for Apache Flink Console to revert the application to the previous running version.

When you perform the explicit rollback, Terraform will initially not be aware of the changes. More specifically, the S3 path to the JAR file in the Managed Service for Apache Flink service (see left part of the image below) is different to the S3 path denoted in the terraform.tfstate file stored in Amazon S3 (see the right part of the image below). Fortunately, Terraform will always perform refreshing actions that include reading the current settings from all managed remote objects and updating the Terraform state to match as part of creating a plan in both terraform plan and terraform apply commands.

Terraform State vs. MSF State

In summary, while you can not perform a manual rollback using Terraform, Terraform will automatically refresh the state when deploying a change using terraform apply.

Deploy config change to application

You have already made changes to the application code of your Apache Flink application. What about making changes to the config of the application, e.g., changing runtime parameters? Imagine you want to change the application logging level of your running Apache Flink application. To change the logging level from ERROR to INFO, you have to change the value for flink_app_monitoring_metrics_level in the terraform/config.tfvars.json to INFO. To deploy the config changes, you need to run the docker run command again as done in the previous sections. This scenario works as expected and is fully covered by Terraform.

What happens when the Apache Flink application deploys successfully but fails and restarts during execution? For more information, please refer to failure mode (b) – the application is started, the job is stuck in a fail-and-restart loop under When starting or updating the application goes wrong. Note that this failure mode can happen when making code changes as well.

When deploying config change goes wrong – The application is started, the job is stuck in a fail-and-restart loop

In the following example, we apply a wrong configuration change preventing the Kinesis connector from initialising correctly, ultimately putting the job in a fail-and-restart loop. To simulate this failure scenario, you’ll need to modify the Kinesis stream configuration by changing the stream name to a non-existent one. This change is made in the terraform/config.tfvars.json file, specifically altering the stream.name value under flink_app_environment_variables. When you deploy with this invalid configuration, the initial deployment will appear successful, showing an Apply Complete! message. The Flink application status will also show as RUNNING. However, the actual behaviour reveals problems. If you check the Flink Dashboard, you’ll see the application is continuously failing and restarting. Also, you will see a warning message about the application requiring attention in the AWS Console.

Problem message within the MSF Console

As detailed in the section Monitoring Apache Flink application operations in the companion blog (part 2), you can monitor the FullRestarts metric to detect the fail-and-restart loop.

Reverting the changes made to the environment variable and deploying the changes will result in Terraform showing the following error message: Failed to take snapshot for the application flink-terraform-lifecycle at this moment. The application is currently experiencing downtime.

Error message 2 from terminal

You have to force-stop without a snapshot and restart the application with a snapshot to get your Flink application back to a properly functioning state. You should constantly monitor the application state of your Apache Flink application to detect any issues.

Other common operations

Manually scaling the application

Another common operation in the lifecycle of your Apache Flink application is scaling the application up or down by adjusting the parallelism. This operation changes the number of Kinesis Processing Units (KPUs) allocated to your application. Let’s look at two different scaling scenarios and how they are handled by Terraform.

In the first scenario, you want to change the parallelism of your running Apache Flink application within the default parallelism quota. To do this, you need to modify the value for flink_app_parallelism in the terraform/config.tfvars.json file. After updating the parallelism value, you deploy the changes by running the Docker command as done in the previous sections:

docker run --env-file .env.docker --rm -it \
-v ./flink:/home/flink-project/flink \
-v ./terraform:/home/flink-project/terraform \
-v ./build.sh:/home/flink-project/build.sh \
msf-terraform bash build.sh apply

This scenario works as expected and is fully covered by Terraform. The application will be updated with the new parallelism setting, and Managed Service for Apache Flink will adjust the allocated KPUs accordingly. Note that there is a default quota of 64 KPUs for a single Managed Service for Apache Flink application, which must be raised proactively via a quota increase request if you need to scale your Managed Service for Apache Flink application beyond 64 KPUs. For more information, refer to Managed Service for Apache Flink quota.

Less common change deployments which require special handling In this section we analyze some less common change deployment scenarios which require some special handling.

Deploy code change that removes an operator

Removing an operator from your Apache Flink application requires special consideration, particularly regarding state management. When you remove an operator, the state from that operator still exists in the latest snapshot, but there’s no longer a corresponding operator to restore it. Let’s take a closer look at this scenario and understand how you can handle it properly. First, you need to make sure that the parameter AllowNonRestoredState is set to True. This parameter specifies whether the runtime is allowed to skip a state that cannot be mapped to the new program, when restoring from a snapshot. Allowing non-restored state is required to successfully update an Apache Flink application when you dropped an operator. To enable the AllowNonRestoredState, you need to set the configuration value for flink_app_allow_non_restored_state to true in terraform/config.tfvars.json. Then, you can go ahead and remove an operator: For example, you can directly have the sourceStream write to the sink connector in flink/src/main/java/com/amazonaws/services/msf/StreamingJob.java. Change code line 146 from windowedStream.sinkTo(sink).uid("kinesis-sink")to sourceStream.sinkTo(sink).uid("kinesis-sink"). Make sure that you have commented out the entire windowedStream code block (lines 103 to 140).

This change will remove the windowed computation and directly connect the source stream to the sink, effectively removing the stateful operation. After removing the operator from your Flink application code, you deploy the changes using the Docker command as previously done. However, the deployment fails with the following error message: Could not execute application. As a result, the Apache Flink application moves to the READY state. To recover from this situation, you need to restart the Apache Flink application using the latest snapshot for the application to successfully start and move to RUNNING status. Importantly, you need to make sure that AllowNonRestoredState is enabled. Otherwise, the application will fail to start as it cannot restore the state for the removed operator.

Deploy change that breaks state compatibility with system rollback enabled

During the lifecycle management of your Apache Flink application, you might encounter scenarios where code changes break state compatibility. This typically happens when you modify stateful operators in ways that prevent them from restoring their state from previous snapshots.

A common example of breaking state compatibility is changing the UID of a stateful operator (such as an aggregation or windowing operator) in your application code. To safeguard against such breaking changes, you can enable the automatic system rollback feature in Managed Service for Apache Flink as described in the subsection Rollback under Lifecycle of an application in Managed Service for Apache Flink previously. This feature is disabled by default and can be enabled using the AWS Management Console or invoking the UpdateApplication API operation. There is no way in Terraform to enable system rollback.

Next, let’s demonstrate this by breaking the state compatibility of your Apache Flink application by changing the UID of a stateful operator, e.g., the string windowed-avg-price in line 140 of flink/src/main/java/com/amazonaws/services/msf/StreamingJob.java to windowed-avg-price-v2 and deploy the changes as before. You will encounter the following error:

Error: waiting for Kinesis Analytics v2 Application (flink-terraform-lifecycle) operation (*) success: unexpected state ‘FAILED’, wanted target ‘SUCCESSFUL’. last error: org.apache.flink.runtime.rest.handler.RestHandlerException: Could not execute application.

At this point, Managed Service for Apache Flink automatically rolls back the application to the previous snapshot with the previous JAR file, maintaining your application’s availability as you have enabled system-rollback capability. Terraform will initially be not aware of the performed rollback. Fortunately, as we have already witnessed in subsection Manual rollback application code to previous application code, Terraform will automatically refresh the state when we change UID to the previous value and deploy the changes.

In-place upgrade of Apache Flink runtime version

Managed Service for Apache Flink supports in-place upgrade to new Flink runtime versions. See the documentation for more details. Updating the application dependencies and any required code changes is a responsibility of the user. Once you have updated the code artifact, the service is able to upgrade the runtime of your running application in-place, without data loss. Let’s examine how Terraform handles Flink version upgrades.

To upgrade your Apache Flink application from version 1.19.1 to 1.20, you need to:

  1. Update the Flink dependencies in your flink/pom.xml to version 1.20.0 (flink.version to 1.20.1 and flink.connector.version to 5.0.0-1.20 in <properties>)
  2. Update the flink_app_runtime_environment to FLINK-1_20 in terraform/config.tfvars.json
  3. Build and deploy the changes using the familiar docker run command

Terraform successfully performs an in-place upgrade of your Flink application. You will receive the following message: Apply complete! Resources: 0 added, 1 changed, 0 destroyed.

Operations currently not supported by Terraform

Let’s take a closer look at operations that are currently not supported by Terraform.

Starting or stopping the application without any configuration change

Terraform provides the start_application parameter, indicating whether to start or stop the application. You can set this parameter using flink_app_start in config.tfvars.json to stop your running Apache Flink application. However, this will only work if the current configuration value is set to true. In other words, Terraform only responds to the change in the parameter value, not the absolute value itself. After Terraform applies this change, your Apache Flink application will stop and its application status will move to READY. Similarly, restarting the application requires changing the flink_app_start value back to true, but this will only take effect if the current configuration value is false. Terraform will then restart your application, moving it back to the RUNNING state.

In summary, you cannot start or stop your Apache Flink application without making any configuration change in Terraform. You have to use AWS CLI, AWS SDK or AWS Console to start or stop your application.

Restarting application from an older snapshot or no snapshot without any configuration change

Similar to the previous section, Terraform requires an actual configuration change of application_restore_type to trigger a restart with different snapshot settings. Simply reapplying the same configuration values won’t initiate a restart from a different snapshot or no snapshot. You have to use AWS CLI, AWS SDK or AWS Console to restart your application from an older snapshot.

Performing rollback triggered manually or by system-rollback feature

Terraform does not support performing a manual rollback nor automatic system rollback. In addition, Terraform will also not be aware when such a rollback is taking place. The state information will be outdated, e.g. S3 path information. However, Terraform automatically performs refreshing actions to read settings from all managed remote objects and updates the Terraform state to match. Consequently, you can have Terraform refresh the Terraform state by successfully running a terraform apply command.

Conclusion

In this post, we demonstrated how to use Terraform to automate the lifecycle management of your Apache Flink applications on Managed Service for Apache Flink. We walked through fundamental operations including creating, updating, and scaling applications, explored how Terraform handles various failure scenarios and examined advanced scenarios such as removing operators and performing in-place runtime upgrades. We also identified operations that are currently not supported by Terraform.

For more information, see Run a Managed Service for Apache Flink application and our two-part blog on Deep dive into the Amazon Managed Service for Apache Flink application lifecycle.


Felix John

Felix John

Felix is a Global Solutions Architect and data & AI expert at AWS, based out of Germany. He focuses on supporting AWS’ strategic global automotive & manufacturing customers on their cloud journey.

Mazrim Mehrtens

Mazrim Mehrtens

Mazrim is a Sr. Specialist Solutions Architect for messaging and streaming workloads. Mazrim works with customers to build and support systems that process and analyze terabytes of streaming data in real time, run enterprise Machine Learning pipelines, and create systems to share data across teams seamlessly with varying data toolsets and software stacks.

Build a data pipeline from Google Search Console to Amazon Redshift using AWS Glue

Post Syndicated from Anirudh Chawla original https://aws.amazon.com/blogs/big-data/build-a-data-pipeline-from-google-search-console-to-amazon-redshift-using-aws-glue/

Google Search Console (GSC) is a service offered by Google that helps you monitor, maintain, and troubleshoot your site’s presence in Google Search results. It provides you unique insights directly from Google about how the search engine sees your site, helping you improve your performance in Search Engine Results Pages (SERPs).

When there is a need to merge Google Search Console data with multiple data sources or conduct complex performance analysis, traditional methods can become time-consuming and error-prone. This is where Amazon Redshift and AWS Glue offer a comprehensive data integration solution.

In this post, we explore how AWS Glue extract, transform, and load (ETL) capabilities connect Google applications and Amazon Redshift, helping you unlock deeper insights and drive data-informed decisions through automated data pipeline management. We walk you through the process of using AWS Glue to integrate data from Google Search Console and write it to Amazon Redshift.

Solution overview

AWS Glue is a serverless data integration service that helps discover, prepare, and combine data for analytics, machine learning (ML), and application development. You can use AWS Glue to create, run, and monitor data integration and ETL pipelines and catalog your assets across multiple data stores.

Amazon Redshift is a fast, scalable, and fully managed cloud data warehouse that lets you to process and run complex SQL analytics workloads on structured and semi-structured data. It also helps you securely access your data in operational databases, data lakes, or third-party datasets with minimal movement or copying of data. Tens of thousands of customers use Amazon Redshift to process large amounts of data, modernize their data analytics workloads, and provide insights for their business users.

The following diagram illustrates the architecture that we implement in this post.

Architecture diagram showing AWS Glue data pipeline workflow from Google Search Console to Amazon Redshift, illustrating the ETL process with AWS Glue job reading data from three Google Search Console entities (Search Analytics, Sites, and Sitemaps) and writing to a Redshift provisioned cluster.

The workflow consists of an AWS Glue job reading data from Google Search Console for the three entities that Google Search Console supports (Search Analytics, Sites, and Sitemaps), and writing the data in a Redshift provisioned cluster. AWS Glue supports Google Search Console API v3.

In the following sections, we walk through the following steps to configure AWS Glue to set up a connection between Google Search Console and Amazon Redshift for data migration:

  1. Create an OAuth client.
  2. Create an IAM role for AWS Glue integration with Google Search Console, AWS Secrets Manager, and Amazon Redshift.
  3. Create a secret in Secrets Manager to store the client secret created in the previous step.
  4. Create a connection to Google Search Console in AWS Glue.
  5. Create a connection to Amazon Redshift in AWS Glue.
  6. Set up a table and permissions in Amazon Redshift.
  7. Create an ETL job in AWS Glue.

Prerequisites

Before starting this walkthrough, you must have the following prerequisites in place:

  • An AWS account.
  • A Google Cloud account and a Google Cloud project.
  • In your Google Cloud project, you must enable the Google Search Console API.
    For instructions, see Enable and disable APIs on the API Console Help for Google Cloud Platform.
  • A provisioned cluster or Amazon Redshift Serverless .
    In this post, we use a single-node ra3.large Redshift provisioned cluster deployed in a single Availability Zone. This configuration is used for demonstration purposes only. For production environments, we recommend using multi-node clusters with a minimum of two nodes deployed across multiple Availability Zones for high availability and better performance.
  • An Amazon Simple Service Storage (Amazon S3) bucket.
  • An AWS Identity and Access Management (IAM) role that grants AWS Glue and Amazon Redshift read-only access to Amazon S3. This role will be attached to the Redshift cluster or Redshift Serverless namespace during creation, and will also be used when running the AWS Glue job along with permissions to read and write secrets to Secrets Manager. Refer to the Amazon Redshift Database Developer Guide for more details.

Create OAuth client

To connect to Google Search Console, AWS Glue requires OAuth 2.0 for authentication. You must create an OAuth 2.0 client ID, which AWS Glue uses when requesting an OAuth 2.0 access token. To create an OAuth 2.0 client ID in the Google Cloud Platform console, follow these steps:

  1. On the Google Cloud Platform console, from the projects list, choose a project or create a new one.
  2. If the APIs & Services page isn’t already open, choose the menu icon on the upper left and choose APIs & Services.
  3. In the navigation pane, choose Credentials.
  4. Choose Create Credentials, then choose OAuth client ID.
  5. Select Web application as the application type, enter NewClient as the name, and provide https://console.aws.amazon.com for Authorized JavaScript origins.
  6. For Authorized redirect URIs, add https://us-east-1.console.aws.amazon.com/gluestudio/oauth. This example uses us-east-1 for setting up AWS Glue jobs; change the redirect URIs according to your AWS Region. Multiple redirect URIs can also be specified.
  7. Choose Create.
  8. Open the details page for your new client.
  9. Under Additional information, note down the client ID and client secret. You will need these details when configuring the secret in Secrets Manager.

Create IAM role for AWS Glue integration with Google Search Console, Secrets Manager, and Amazon Redshift

You can use AWS Glue to transfer data from supported sources into your Redshift databases. You need an IAM role because AWS Glue needs authorization to write into Redshift databases. To create a role, complete the following steps:

  1. Sign in to the IAM console with sufficient access to create policies.
  2. Choose Policies in the navigation pane.
  3. Choose Create policy.
  4. On the JSON tab, enter the following policy. AWS Glue needs the following permissions to access and run SQL statements in the Redshift database and create and retrieve secrets with Secrets Manager:
    {
        "Version": "2012-10-17",
        "Statement": [
            {
                "Effect": "Allow",
                "Action": [
                    "secretsmanager:DescribeSecret",
                    "secretsmanager:GetSecretValue",
                    "secretsmanager:PutSecretValue",
                    "ec2:CreateNetworkInterface",
                    "ec2:DescribeNetworkInterfaces",
                    "ec2:DeleteNetworkInterface"
                ],
                "Resource": "*"
            },
            {
                "Effect": "Allow",
                "Action": "s3:GetObject",
                "Resource": "arn:aws:s3:::aws-glue-studio-transforms-510798373988-prod-us-east-1/*"
            },
            {
                "Effect": "Allow",
                "Action": [
                    "s3:GetObject",
                    "s3:PutObject"
                ],
                "Resource": [
                    "arn:aws:s3:::aws-glue-assets-testbucket/*"
                ]
            },
            {
                "Sid": "DataAPIPermissions",
                "Effect": "Allow",
                "Action": [
                    "redshift-data:ExecuteStatement",
                    "redshift-data:GetStatementResult",
                    "redshift-data:DescribeStatement"
                ],
                "Resource": "*"
            },
            {
                "Sid": "GetCredentialsForAPIUser",
                "Effect": "Allow",
                "Action": "redshift:GetClusterCredentials",
                "Resource": [
                    "arn:aws:redshift:*:*:dbname:*/*",
                    "arn:aws:redshift:*:*:dbuser:*/*"
                ]
            },
            {
                "Sid": "GetCredentialsForServerless",
                "Effect": "Allow",
                "Action": "redshift-serverless:GetCredentials",
                "Resource": "*"
            },
            {
                "Sid": "DenyCreateAPIUser",
                "Effect": "Deny",
                "Action": "redshift:CreateClusterUser",
                "Resource": [
                    "arn:aws:redshift:*:*:dbuser:*/*"
                ]
            },
            {
                "Sid": "ServiceLinkedRole",
                "Effect": "Allow",
                "Action": "iam:CreateServiceLinkedRole",
                "Resource": "arn:aws:iam::*:role/aws-service-role/redshift-data.amazonaws.com/AWSServiceRoleForRedshift",
                "Condition": {
                    "StringLike": {
                        "iam:AWSServiceName": "redshift-data.amazonaws.com"
                    }
                }
            }
        ]
    }

    Modify the S3 bucket name that you are using as the staging bucket. Additionally, AWS Glue must have access to specific AWS owned S3 buckets for hosting AWS Glue transforms. In this example, the IAM policy uses aws-glue-studio-transforms-510798373988-prod-us-east-1, which is the AWS owned bucket in the us-east-1 Region. Refer to Review IAM permissions needed for ETL jobs for the appropriate bucket name for your Region.

  5. Choose Next.
  6. For Policy name, enter a name (for this post, we use glue-redshift-gsc-policy).
  7. Enter a description, then choose Create policy.
  8. In the navigation pane, choose Roles and Create role.
  9. Choose Custom trust policy and enter the following, then choose Next.
    {
        "Version": "2012-10-17",
        "Statement": [
            {
                "Effect": "Allow",
                "Principal": {
                    "Service": [
                        "glue.amazonaws.com"
                    ]
                },
                "Action": "sts:AssumeRole"
            }
        ]
    }
    

  10. Search for and select the policy glue-redshift-gsc-policy, then choose Next.
  11. Provide the role name GlueIAMRoleRedshiftNew or another name and relevant Description, then choose Create role.
  12. After the role is created, choose Add permissions and Attach policies.
  13. Search for AWSGlueServiceRole and choose Add Permissions. This policy is typically attached to roles specified when defining crawlers, jobs, and development endpoints.

Screenshot of AWS IAM console showing the policy attachment interface where the AWSGlueServiceRole policy is being added to the GlueIAMRoleRedshiftNew role.

Create secret in Secrets Manager

Complete the following steps to create a Secrets Manager secret:

  1. On the Secrets Manager console, choose Store a new secret.
  2. Select Other type of secret.
  3. For the customer-managed connected application, the secret should contain the connected application’s consumer secret with USER_MANAGED_CLIENT_APPLICATION_CLIENT_SECRET as the key and the client secret value as created in the previous step.
    Screenshot of AWS Secrets Manager console showing the "Store a new secret" interface with "Other type of secret" selected and a key-value pair entry for USER_MANAGED_CLIENT_APPLICATION_CLIENT_SECRET.
  4. Choose Next.
  5. Enter a secret name and choose Next.
  6. Choose Store.

Create connection to Google Search Console in AWS Glue

To create a connection to Google Search Console in AWS Glue, follow these steps:

  1. Sign in to the AWS Glue console with an authorized email ID with permissions already provided in Google Search Console.
  2. In the navigation pane, choose Data connections.
  3. Under Connections, choose Create connection.
  4. In Data sources, search for Google Search Console and choose Next.
    Screenshot of AWS Glue console showing the Data connections page with Google Search Console selected as a data source in the connection creation wizard.
  5. For IAM Role ARN, choose the role created earlier.
  6. For Token URL, use https://oauth2.googleapis.com/token, which is the default value.
  7. For User Managed Client Application ClientId, enter the client ID created earlier while creating the OAuth client.
  8. For AWS Secret, choose the secret created earlier.
  9. If your AWS Glue jobs needs to run in an Amazon virtual private cloud (VPC), provide appropriate details. For more information, refer to Configure a VPC for your ETL job.
    Screenshot of AWS Glue connection configuration form showing fields for IAM Role ARN, Token URL, User Managed Client Application ClientId, AWS Secret selection, and VPC configuration options
  10. Choose Test connection, choose your Google ID, and choose Continue.
    Google account selection dialog prompting the user to choose which Google account to use for authentication with the AWS Glue connection.
  11. Choose Continue to trust the connection.
    Google OAuth consent screen asking the user to continue and trust the connection between AWS Glue and their Google account.

    If the user has authorized access, the connection test will be successful.

    AWS Glue console showing a successful connection test result with a green checkmark indicating the Google Search Console connection was established successfully.

  12. Choose Next.
  13. Provide a connection name and choose Create connection.

Create connection to Amazon Redshift in AWS Glue

Complete the following steps to set up an AWS Glue connection for Amazon Redshift. Refer to Redshift connections for more information.

  1. On the AWS Glue console, in the navigation pane, choose Data connections.
  2. Under Connections, choose Create connection.
  3. In Data sources, search for JDBC and choose Next. For Amazon Redshift, you can also use Redshift connections. In this post, we use JDBC. In this example, we are using a Redshift provisioned cluster.
  4. Provide the Amazon Redshift JDBC URL and either use a Secrets Manager secret for storing credentials or provide the user name and password directly. As a best practice, it is recommended to use Secrets Manager.
  5. Configure network options with Amazon VPC settings for running the AWS Glue job in a VPC. In this example, we use the same VPC, subnet, and security group where the Redshift cluster is provisioned. All JDBC data stores must be accessible from the VPC subnet. A VPC endpoint is required to access Amazon S3 from within your VPC. If your job needs to access both VPC resources and the public internet, configure a NAT gateway in the VPC.Screenshot of AWS Glue connection configuration for Amazon Redshift showing JDBC URL entry, credentials configuration options (Secrets Manager or direct username/password), and VPC network settings including VPC, subnet, and security group selections.

Set up table and permissions in Amazon Redshift

To set up table and permissions in Amazon Redshift, follow these steps:

  1. On the Amazon Redshift console, choose Query editor v2.
  2. Connect to your existing Redshift cluster.
  3. Create a table with the following DDL. For this post, we create a new database named test and create the following tables in the public schema of test database:
    #Create Database command
    CREATE DATABASE test; 
    
    #Sitemap table creation
    CREATE TABLE public.sitemap(
        path VARCHAR(4096) ENCODE lzo,
        type VARCHAR(255) ENCODE lzo,
        lastSubmitted TIMESTAMP ENCODE delta,
        isPending BOOLEAN NULL ENCODE raw,
        isSitemapsIndex BOOLEAN NULL ENCODE raw,
        lastDownloaded TIMESTAMP NULL ENCODE delta,
        warnings BIGINT NULL ENCODE delta,
        errors BIGINT NULL ENCODE delta,
        contents VARCHAR(65535) NULL ENCODE lzo) DISTSTYLE AUTO;
        
    #Search Analytics table creation
    CREATE TABLE public.search_analytics (
        keys character varying(2048) ENCODE lzo,
        clicks double precision ENCODE raw,
        impressions double precision ENCODE raw,
        ctr numeric(38, 18) ENCODE az64,
        position double precision ENCODE raw
    ) DISTSTYLE AUTO;
    
    #Sites table creation
     CREATE TABLE public.sites (
        siteurl character varying(2048) ENCODE lzo,
        permissionLevel character varying(50) ENCODE lzo
    ) DISTSTYLE AUTO;

    Screenshot of AWS Glue ETL job visual editor showing the job creation interface with source and target selection options, displaying Google Search Console as source and Amazon Redshift as target.

Create ETL job in AWS Glue

To create a data flow in AWS Glue, follow these steps:

  1. On the AWS Glue console, choose ETL jobs in the navigation pane.
  2. Choose Visual ETL under Create job.
    Each ETL job in AWS Glue is priced based on its duration.

    Screenshot of AWS Glue visual ETL canvas showing a data flow diagram with Google Search Console source node connected to Amazon Redshift target node.

  3. For the source, choose Google Search Console, and for the target, choose Amazon Redshift.
    Screenshot of AWS Glue source node configuration panel showing Google Search Console connection settings with entity selection (Sites) and field selection options (siteUrl and permissionLevel).
  4. Choose Source (Google Search Console) to configure the properties, which opens in the right window pane.
  5. Choose the Google Search Console connection created in the previous sections, and provide the entity name. At the time of writing, there are three supported entities: Search Analytics, Sites, and Sitemaps, with multiple supported fields and operators for each entity. Choose the entity name and the corresponding fields; by default, the connector selects all fields. The example shows selecting the entity Site and corresponding fields siteUrl and permissionLevel.
    Screenshot of AWS Glue target node configuration panel showing Amazon Redshift connection settings including schema selection, table name, data handling method (Append to target table), and S3 staging directory configuration.
  6. Choose Target (Amazon Redshift) to configure the properties, which opens in the right pane.
  7. Choose the Amazon Redshift connection, schema, and table name that were created in the previous steps. In this example, we use Append to target table as the method for handling the data. An S3 directory is provided for staging temporary data.
    Screenshot of AWS Glue target node configuration panel showing Amazon Redshift connection settings including schema selection, table name, data handling method (Append to target table), and S3 staging directory configuration.
  8. Navigate to Job details and provide a job name and IAM role (which the job will assume while running). This is the same role created earlier.
  9. Choose Save and Run. For this example, we use AWS Glue version 5.0, keeping all other configuration values under Job details at their defaults. For this example, we have not implemented any schema mapping, so the columns in Amazon Redshift were created to match the output response for the Search entity.
  10. After the job has completed successfully, navigate to Query Editor v2 in Amazon Redshift and query the Sites table to preview the data.
    Screenshot of Amazon Redshift Query Editor v2 showing query results from the Sites table with columns for siteurl and permissionlevel, displaying sample data rows.Screenshot of Amazon Redshift Query Editor v2 showing query results from the Sites table with columns for siteurl and permissionlevel, displaying sample data rows.
  11. In the case of job failures, validate the connections by doing a data preview, and refer to Troubleshooting AWS Glue.
  12. Similar to the Site entity, you can load Sitemap entity data by changing the source properties and destination table in the target Redshift cluster, then choosing Run.
    Screenshot of AWS Glue source node configuration showing Google Search Console entity selection changed to Sitemaps with corresponding fields selected.
  13. Navigate to Query Editor v2 in Amazon Redshift and query the sitemap table to preview the data.
    Screenshot of Amazon Redshift Query Editor v2 showing query results from the sitemap table with columns including path, type, lastsubmitted, ispending, issitemapsindex, lastdownloaded, warnings, errors, and contents.
  14. Similar to Sitemap, you can load Search Analytics entity data by changing the source properties and destination table in the target Redshift cluster, then choosing Run.
    Screenshot of AWS Glue source node configuration showing Google Search Console entity selection changed to Search Analytics with corresponding fields selected.
  15. Navigate to Query Editor v2 in Amazon Redshift and query the search_analytics table and preview the data.
    Screenshot of Amazon Redshift Query Editor v2 showing query results from the search_analytics table with columns for keys, clicks, impressions, ctr, and position.

Filter predicates with Search Analytics

The Search Analytics entity provides support for multiple filters that can be used to view the traffic data for the sites. The following examples show use of some filter predicates you can use that Google Search Console connections support.

  • start_end_date – The default value for start_end_date is between <30 days ago from the current date> AND <yesterday>. To use a different date range, use the between The following example displays search data from January through September 2025:
    start_end_date between '2025-01-01' AND '2025-09-30'

    Screenshot of AWS Glue source node configuration showing Search Analytics entity with a filter predicate for start_end_date between '2025-01-01' AND '2025-09-30'.

  • device – The device filters result against specified device type like DESKOP, MOBILE, and TABLET:
    device = 'MOBILE'

    Screenshot of AWS Glue source node configuration showing Search Analytics entity with a filter predicate for device = 'MOBILE'.

  • country – You can filter against the specified country, as specified by three-letter country code (ISO 3166-1 alpha-3):
    dimensions='country'

    Screenshot of AWS Glue source node configuration showing Search Analytics entity with dimensions set to 'country'.

  • dimensions: Dimensions help group zero or more results for filtering search data by country or device. The following example displays search data grouped by country, and also grouping by country and filtering for mobile devices:
    dimensions='country' AND country='ind' AND device ='MOBILE'

    Screenshot of AWS Glue source node configuration showing Search Analytics entity with multiple filter predicates including dimensions='country', country='ind', and device='MOBILE'.

Run analytical queries on Amazon Redshift

In this section, we run analytical queries using aggregated data across different search entities.

List all countries where site position is less than 10 and device type is MOBILE:

SELECT * from search_analytics_device_country where position < 10 AND keys LIKE '%MOBILE%'

Screenshot of Amazon Redshift Query Editor v2 showing query results for countries where site position is less than 10 and device type is MOBILE, displaying data from the search_analytics_device_country table.

List all countries where impressions are greater than 1 and position is less than 10:

SELECT * FROM "test"."public"."search_analytics_country" where impressions > 1 and position < 10;

Screenshot of Amazon Redshift Query Editor v2 showing query results for countries where impressions are greater than 1 and position is less than 10, displaying data from the search_analytics_country table.

Clean up

To avoid incurring charges, clean up the resources in your AWS account by completing the following steps:

  1. On the AWS Glue console, in the navigation pane, choose Job monitoring.
  2. Stop any running jobs created for Google Search Console connections.
  3. From the list of connections, select the connection name created and delete it.
  4. Delete the Redshift provisioned cluster or the Redshift Serverless workspace and namespace. Amazon Redshift pricing is applied during the cluster’s runtime based on cluster configuration.
  5. Clean up resources in your Google account by deleting the project that contains the Google Project resources. For instructions, refer to Delete your project.

Conclusion

In this post, we walked you through the process of using AWS Glue to integrate data from Google Search Console and write it to Amazon Redshift, a petabyte-scale data warehouse. Whether you’re archiving historical data, performing complex analytics, or preparing data for machine learning, this connector streamlines the process and helps create an integrated data pipeline.

For more information, refer to AWS Glue support for Google Search Console.


About the authors

Anirudh Chawla

Anirudh Chawla

Anirudh is an AWS Analytics Specialist Solutions Architect. He likes to read books, take long walks in nature, and participate in community programs.

Shubham Purwar

Shubham Purwar

Shubham is an AWS Analytics Specialist Solution Architect. In his free time, Shubham loves to spend time with his family and travel around the world.

Shaswat Mandhanya

Shaswat Mandhanya

Shaswat is an AWS Analytics Specialist BD. In his free time, he likes to watch Formula 1 races and travel across the country.

Prabhu G

Prabhu G

Prabhu is a Solutions Architect at AWS. He is an avid supporter of Chennai Super Kings and a big-time fan of MS Dhoni.

Amazon OpenSearch Service 101: T-shirt size your domain for e-commerce search

Post Syndicated from Abe Raghib original https://aws.amazon.com/blogs/big-data/amazon-opensearch-service-101-t-shirt-size-your-domain-for-e-commerce-search/

In e-commerce, delivering fast, relevant search results helps users find products quickly and accurately, improving satisfaction and increasing sales. OpenSearch is a distributed search engine that offers advanced search capabilities including advanced full-text and faceted search, customizable analyzers and tokenizers, and auto-complete to help customers quickly find the products they want. It scales to handle millions of products, catalogs and traffic surge. Amazon OpenSearch Service is a managed service that lets users build search workloads balancing search quality, performance at scale and cost. Designing and sizing an Amazon OpenSearch Service cluster correctly is required to meet these demands.

While general sizing guidelines for OpenSearch Service domains are covered in detail in OpenSearch Service documentation, in this post we specifically focus on T-shirt-sizing OpenSearch Service domains for e-commerce search workloads. T-shirt sizing simplifies complex capacity planning by categorizing workloads into sizes like XS, S, M, L, XL based on key workload parameters such as data volume and query concurrency. For e-commerce search, where data growth is moderate and read-heavy queries predominate, this approach offers a flexible, scalable way to allocate the resources without overprovisioning or underestimating needs.

How OpenSearch Service stores indexes and performs queries

E-commerce search platforms handle vast amounts of data and daily data ingestion is typically relatively small and incremental, reflecting catalog changes, price updates, inventory status and user activities like clicks and reviews. Efficiently managing this data and organizing it per OpenSearch Service best practices is crucial in achieving optimal performance. The workload is read-heavy, consisting of user queries with advanced filtering and faceting, especially during sales or seasonal spikes that require elasticity in compute and storage resources.

You ingest product and catalog updates (inventory, listings, pricing) into OpenSearch using bulk APIs or real-time streaming. You index data into logical indexes. How you create and organize indexes in e-commerce has a significant impact on search, scalability and flexibility. The approach depends on the size, diversity and operational needs of the catalog. Small to medium-sized e-commerce platforms commonly use a single, comprehensive product index that stores all product information with product category. Additional indexes may exist for orders, users, reviews and promotions depending on search requirements and data separation needs. Large, diverse catalogs may split products into category-specific indexes for tailored mappings and scaling. You split each index into primary shards, each storing a portion of the documents. To ensure high availability and enhance query throughput, you configure each primary shard with at least one replica shard stored on different data nodes.

Diagram showing Amazon OpenSearch Service cluster with three data nodes implementing primary-replica shard distribution for Products and Reviews indexes. Data Node 1 contains P0-Products (primary), P0-Reviews (primary), and R0-Reviews (replica). Data Node 2 contains P1-Products (primary), R0-Products (replica), and R1-Reviews (replica). Data Node 3 contains P1-Reviews (primary) and R1-Products (replica). Color-coded legend distinguishes primary shards (filled boxes) from replica shards (outlined boxes) for both indexes. This architecture ensures fault tolerance and high availability by distributing primary and replica shards across different nodes.
Diagram 1. How primary and replica shards are distributed among nodes

This diagram shows two indexes (Products and Reviews), each split into two primary shards with one replica. OpenSearch distributes these shards across cluster nodes to ensure that primary and replica shards for the same data do not reside on the same node. OpenSearch runs search requests using a scatter-gather mechanism. When an application submits a request, any node in the cluster can receive it. This receiving node becomes the coordinating node for that specific query. The coordinating node determines which indices and shards can serve the query. It forwards the query to either primary or replica shards and orchestrates the different phases of the search operation and returns the response. This process ensures efficient distribution and execution of search requests across the OpenSearch cluster.

Diagram showing Amazon OpenSearch Service distributed query architecture where an e-commerce application searches for "Blue running shoes." The workflow demonstrates the scatter-gather pattern: (1) Application sends query to Coordinator Node, (2) Coordinator scatters query to three Data Nodes containing index shards, (3) Data Nodes execute searches in parallel, (4) Coordinator gathers and merges results, then returns final ranked results to application. This architecture enables horizontal scalability, parallel processing, and fault tolerance for high-performance search operations across distributed data.
Diagram 2. Tracing a Search query: “blue running shoes”This diagram walks through how a search query–for example, “blue running shoes”flows through your OpenSearch Service domain .

  1. Request: The application sends the search for “blue running shoes” to the domain. One data node acts as the coordinating node.
  2.  Scatter: The coordinator broadcasts the query to either the primary or replica shard for each of the shards in the ‘Products’ index (Nodes 1, 2, and 3 in this case).
  3. Gather: Each data node searches its local shards(s) for “blue running shoes” and returns its own top results (e.g. Node 1 returns its best matches from P0).
  4. Final results: The coordinator merges these partial lists, sorts them into single definitive list of the most relevant shoes, and sends the result back to the app.

Understanding T-Shirt Sizing for E-commerce OpenSearch Service Cluster

Storage planning

Storage impacts both performance and cost. OpenSearch Service offers two main storage options based on query latency requirements and data persistence needs. Selecting the appropriate storage type in a managed OpenSearch Service improves both performance and optimizes cost of the domain. You can choose between Amazon Elastic Block Store( EBS) storage volumes and instance storage volumes (local storage) for your data nodes.

Amazon EBS gp3 volumes offer high throughput, whereas the local NVMe SSD volumes, for example, on the r8gd, i3, or i4i instance families, offer low latency, fast indexing performance and high-speed storage, making them ideal for scenarios where real time data updates and high search throughput are critical for search operations. For search workloads that require a balance between performance and cost, instances backed with EBS GP3 SSD volumes provide a reliable option. This SSD storage offers input/output operations per second (IOPS) that are well-suited for general-purpose search workloads. It also allows users to provision additional IOPS and storage as needed.

When sizing an OpenSearch cluster, start by estimating total storage needs based on catalog size and expected growth. For example, if the catalog contains 500,000 stock keeping units (SKUs), averaging 50KB each; the raw data sums to about 25GB. The size of the raw data, however, is just one aspect of the storage requirements. Also consider the Replica count, indexing overhead (10%), Linux reserves (5%), and OpenSearch Service reserves (20% up to 20GB) per instance while calculating the required storage.

In summary,if you have 25GB of data at any given time who want one replica, the minimum storage requirement is closer to 25 * 2 * 1.1 / 0.95 / 0.8 = 72.5 GB. This calculation can be generalized as follows:

Storage requirement = Raw data * (1 + number of replicas) * 1.45 

This helps ensure disk space headroom on all data nodes, preventing shard failures and maintaining search performance. Provisioning storage slightly beyond this minimum is recommended to accommodate future growth and cluster rebalancing.

Data nodes:

For search workloads, compute-optimized instances (C8g) are well-suited for central processing unit (CPU)-intensive operations like nested queries and joins. However, general-purpose instances like M8g offer a better balance between CPU and memory. Memory-optimized instances (R8g, R8gd) are recommended for memory-intensive operations like KNN search, where larger memory footprint is required. In large, complex deployments, compute-optimized instances like c8g or general-purpose m8g, handle CPU-intensive tasks, providing efficient query processing and balanced resource allocation. The balance between CPU and memory, makes them ideal for managing complex search operations for large-scale data processing. For extremely large search workloads (tens of TB) where latency is not a primary concern, consider using the new Amazon OpenSearch Service Writable warm which supports write operations on warm indices.

Instance Class Best for users who… Examples (AWS) Characteristics
General Purpose have moderate search traffic and want a well-balanced, entry-level setup M family (M8g) Balanced CPU & memory, EBS storage. Good starting point for small to medium-sized catalogs.
Compute Optimized have high queries per second (QPS) search traffic or queries involve scoring scripts or complex filtering C family (C8) High CPU-to-memory ratio. Ideal for CPU-bound workloads like many concurrent queries.
Memory Optimized work with large catalogs, need fast aggregations, or cache a lot in memory R family (R8g) More memory per core. Holds large indices in memory to speed up searches and aggregations.
Storage Optimized update inventory frequently or have so much data that disk access slows things down I family (I3, I4g), Im4gn NVMe SSD and SSD local storage. Best for I/O-heavy operations like constant indexing or large product catalogs hitting disk frequently.

Cluster manager nodes:

For production workloads, it is recommended to add dedicated cluster manager nodes to increase the cluster stability and offload cluster management tasks from the data nodes. To choose the right instance type for your cluster manager nodes, review the service recommendations based on the OpenSearch version and number of shards in the cluster.

Sharding strategy

Once storage requirements are understood, you can investigate the indexing strategy. You create shards in OpenSearch Service to distribute an index evenly across the nodes in a cluster. AWS recommends single product index with category facets for simplicity or partition indexes by category for large or distributed catalogs. The size and number of shards per index play a vital role in OpenSearch Service performance and scalability. The right configuration ensures balanced data distribution, avoids hot spotting, and minimizes coordination overhead on nodes for use cases that prioritizes query speed and data freshness.

For read-heavy workloads like e-commerce, where search latency is the key performance objective, maintain shard sizes between 10-30GB. To achieve this, calculate the number of primary shards by dividing your total index size by your target shard size. For example, if you have a 300GB index and want 20GB shards, configure 15 primary shards (300GB ÷ 20GB = 15 shards). Monitor shard sizes using the _cat/shards API and adjust the shard count during reindexing if shards grow beyond the optimal range.

Add replica shards to improve search query throughput and fault tolerance. The minimum recommendation is to have one replica; you can add more replicas for high query throughput requirements. In OpenSearch Service, a shard processes operations like querying single-threaded, meaning one thread handles a shard’s tasks at a time. Replica shards can serve read requests by distributing them across multiple threads and nodes, enabling parallel processing.

T-shirt sizing for an e-commerce workload

In an OpenSearch T-shirt sizing table, each size label (XSmall, Small, Medium, Large, XLarge) represents a generalized cluster scale category that can help teams translate technical requirements into simple, actionable capacity planning. Each size allows architects to quickly align their catalog size, storage requirements, shard planning, CPU and AWS instance choices to the cluster resources provisioned, making it easier to scale infrastructure as business grows.

By referring to this table, teams can select the category similar to their current workload and use the T-shirt size as a starting point while continuing to refine configuration as they monitor and optimize real-world performance. For example, XSmall is suited for small catalogs with hundreds of thousands of products and minimal search traffic. Small clusters are designed for growing catalogs with millions of SKUs, supporting moderate query volumes and scaling up during busy periods. Medium corresponds to mid-size e-commerce operations handling millions of products and higher search demands, while Large fits large online businesses with tens of millions of SKUs, requiring robust infrastructure for fast, reliable search. XLarge is intended for major marketplaces or global platforms with twenty million or more SKUs, enormous data storage needs, and massive concurrent usage.

T-shirt size Number of Products Catalog Size Storage needed Primary Shard Count Active Shard Count Data Nodes Instance Type Cluster Manager Node instanceType
XSmall 500K 50 GB 145 GB 2 4 [2] r8g.xlarge [3] m8g.large
Small 2M 200 GB 580 GB 8 16 [2] c8g.4xlarge [3] m8g.large
Medium 5M 500 GB 1.45 TB 20 40 [2] c8g.8xlarge [3] m8g.large
Large 10M 1 TB 2.9 TB 40 80 [4] c8g.8xlarge [3] m8g.large
XLarge 20M 2 TB 5.8 TB 80 160 [4] c8g.16xlarge [3] m8g.large
  • T-shirt size: Represents the scale of the cluster, ranging from XS up to XL for high-volume workloads.
  • Number of products: The estimated count of SKUs in the e-commerce catalog, which drives the data volume.
  • Catalog size: The total estimated disk size of all indexed product data, based on typical SKU document size.
  • Storage needed: The actual storage required after accounting for replicas and overhead, ensuring enough room for safe and efficient operation.
  • Primary shard count: The number of main index shards chosen to balance parallel processing and resource management.
  • Active shard count: The total number of live shards (primary with one replica), indicating how many shards need to be distributed for availability and performance.
  • Data node instance type: The recommended instance type to use for data nodes, selected for memory, CPU, and disk throughput.
  • Cluster manager node instance type: The recommended instance type for lightweight, dedicated master nodes which manage cluster stability and coordination.

Scaling strategies for e-commerce workloads

E-commerce platforms continually face challenges with unpredictable traffic surges and growing product catalogs. To address these challenges, OpenSearch Service automatically publishes critical performance metrics to Amazon CloudWatch, enabling users to monitor when individual nodes reach resource limits. These metrics include CPU utilization exceeding 80%, JVM memory pressure above 75%, frequent garbage collection pauses, and thread pool rejections.

OpenSearch Service also provides robust scaling solutions that maintain consistent search performance across varying workload demands. Use the vertical scaling strategy to upgrade instance types from smaller to larger configurations, such as m6g.large to m6g.2xlarge. While vertical scaling triggers a blue-green deployment, scheduling these changes during off-peak hours minimizes impact on operations.

Use the horizontal scaling strategy to add more data nodes for distributing indexing and search operations. This approach proves particularly effective when scaling for traffic growth or increasing dataset size. In domains with cluster manager nodes, adding data nodes proceeds smoothly without triggering a blue-green deployment. CloudWatch metrics guide horizontal scaling decisions by monitoring thread pool rejections across nodes, indexing latency, and cluster-wide load patterns. Though the process requires shard rebalancing and may temporarily impact performance, it effectively distributes workload across the cluster.

Temporary replicas provide a flexible solution for managing high-traffic periods. By increasing replica shards through the _settings API, read throughput can be boosted when needed. This approach offers a dynamic response to changing traffic patterns without requiring more substantial infrastructure changes.

For more information on scaling an OpenSearch Service domain, please refer to How do I scale up or scale out an OpenSearch Service domain?

Monitoring and operational best practices

Monitoring key performance CloudWatch metrics is essential to ensure a well-optimised OpenSearch service domain. One of the key factors is maintaining CPU utilization on data nodes under 80% to prevent query slowdowns. Another metric is ensuring that JVM memory pressure is maintained below 75% on data nodes to prevent garbage collection (GC) pauses that can affect search response time. OpenSearch service publishes these metrics to CloudWatch at 1 minute interval and users can create alarms on these metrics for alerts on the production workloads. Please refer recommended CloudWatch alarms for OpenSearch Service

P95 query latency should be monitored to identify slow queries and optimize performance. Another important indicator is thread pool rejections. A high number of thread pool rejections can result in failed search requests, and affecting user experience. By continuously monitoring these CloudWatch metrics, users can proactively scale resources, optimise queries, and prevent performance bottlenecks.

Conclusion

In this post, we showed how to right-size Amazon OpenSearch Service domains for e-commerce workloads using a T-shirt sizing approach. We explored key factors including storage optimization, sharding strategies, scaling methods, and essential Amazon CloudWatch metrics for monitoring performance.

To build a performant search experience, start with a smaller deployment and iterate as your business scales. Get started with these five steps:

  1. Evaluate your workload requirements in terms of storage, search throughput, and search performance
  2. Select your initial T-shirt size based on your product catalog size and traffic patterns
  3. Deploy the recommended sharding strategy for your catalog scale
  4. Load test your cluster using OpenSearch benchmark and re-iterate until performance requirements are reached
  5. Configure Amazon CloudWatch monitoring and alarms, then continue to monitor your production domain


About the authors

Raaga NG

Raaga NG

Raaga is a Solutions Architect at AWS with over 5 years of experience helping enterprises modernize their technology landscape and build scalable, cloud-native solutions. She partners with customers to translate business requirements into efficient cloud architectures that drive measurable outcomes, supporting their journey from application modernization to AI adoption through thoughtful, customer-centric solutions

Harsh Bansal

Harsh Bansal

Harsh is an Analytics and AI Solutions Architect at AWS. Bansal collaborates closely with clients, assisting in their migration to cloud platforms and optimizing cluster setups to enhance performance and reduce costs. Before joining AWS, Bansal supported clients in leveraging OpenSearch and Elasticsearch for diverse search and log analytics requirements.

Aditya Challa

Aditya Challa

Aditya is a Senior Solutions Architect at AWS. Aditya loves helping customers through their AWS journeys because he knows that journeys are always better when there’s company. He’s a big fan of travel, history, engineering marvels, and learning something new every day.

Abe Raghib

Abe Raghib

Abe is a Senior Solutions Architect at AWS. Abe helps enterprises modernize applications and build scalable, cloud-native solutions. He works with customers to translate business needs into secure, scalable, and cost-effective architectures while supporting their data modernization and AI adoption journeys to drive innovation and measurable business outcomes.

Matching your Ingestion Strategy with your OpenSearch Query Patterns

Post Syndicated from Rakan Kandah original https://aws.amazon.com/blogs/big-data/matching-your-ingestion-strategy-with-your-opensearch-query-patterns/

Choosing the right indexing strategy for your Amazon OpenSearch Service clusters helps deliver low-latency, accurate results while maintaining efficiency. If your access patterns require complex queries, it’s best to re-evaluate your indexing strategy.

In this post, we demonstrate how you can create a custom index analyzer in OpenSearch to implement autocomplete functionality efficiently by using the Edge n-gram tokenizer to match prefix queries without using wildcards.

What is an index analyzer?

Index analyzers are used to analyze text fields during ingestion of a document. The analyzer outputs the terms you can use to match queries

By default, OpenSearch indexes your data using the standard index analyzer. The standard index analyzer splits tokens on spaces, converts tokens to lowercase, and removes most punctuation. For some use cases (like log analytics), the standard index analyzer might be all you need.

Standard Index Analyzer

Let’s look at what the standard index analyzer does. We’ll use the _analyze API to test how the standard index analyzer tokenizes the sentence “Standard Index Analyzer.”

Note: You can run all the commands in this post using OpenSearch DevTools in the OpenSearch Dashboard.

GET /_analyze
{
  "analyzer": "standard",
  "text": "Standard Index Analyzer."
}
#========
#Results
#========
{
  "tokens": [
    {
      "token": "standard",
      "start_offset": 0,
      "end_offset": 8,
      "type": "<ALPHANUM>",
      "position": 0
    },
    {
      "token": "index",
      "start_offset": 9,
      "end_offset": 14,
      "type": "<ALPHANUM>",
      "position": 1
    },
    {
      "token": "analyzer",
      "start_offset": 15,
      "end_offset": 23,
      "type": "<ALPHANUM>",
      "position": 2
    }
  ]
}

Notice how each word was lowercased and the period (punctuation) was removed.

Creating your own index analyzer

OpenSearch offers a large number of built in analyzers that you can use for different access patterns. It also lets you build your own custom analyzer, configured for your specific search needs. In the following example, we are going to configure a custom analyzer that returns partial word matches for a list of addresses. The analyzer is specifically designed for autocomplete functionality, enabling end users to quickly find addresses without having to type out (or remember) an entire address. Autocomplete allows OpenSearch to effectively complete the search term based off matched prefixes.

First, create an index called standard_index_test:

PUT standard_index_test
{
  "mappings": {
    "properties": {
      "text_entry": {
        "type": "text",
        "analyzer": "standard"
      }
    }
  }
}

Specifying the analyzer as standard is not required because the standard analyzer is the default analyzer.

To test, bulk add some data to our standard_index_test that we created.

POST _bulk
{"index":{"_index":"standard_index_test"}} 
{"text_entry": "123 Amazon Street Seattle, Wa 12345 "} 
{"index":{"_index":"standard_index_test"}}
{"text_entry": "456 OpenSearch Drive Anytown, Ny 78910"}
{"index":{"_index":"standard_index_test"}}
{"text_entry": "789 Palm way Ocean Ave, Ca 33345"}
{"index":{"_index":"standard_index_test"}}
{"text_entry": "987 Openworld Street, Tx 48981"}

Query this data using the text “ope”.

GET standard_index_test/_search
{
  "query": {
    "match": {
      "text_entry": {
        "query": "ope"
      }
    }
  }
}
#========
#Results
#========
{
  "took": 2,
  "timed_out": false,
  "_shards": {
    "total": 5,
    "successful": 5,
    "skipped": 0,
    "failed": 0
  },
  "hits": {
    "total": {
      "value": 0,
      "relation": "eq"
    },
    "max_score": n`ull,
    "hits": [] # No matches 
  }
}

When searching for the term “ope”, we don’t get any matches. To see why, we can dive a little deeper into the standard index analyzer and see how our text is being tokenized. Test the standard index analyzer with the address “456 OpenSearch Drive Anytown, Ny 78910”.

POST standard_index_test/_analyze
{
  "analyzer": "standard",
  "text": "456 OpenSearch Drive Anytown, Ny 78910"
}
#========
#Results
#========
  "tokens":
      "456" 
      "opensearch" 
      "drive" 
      "anytown"
      "ny" 
      "78910"

The standard index analyzer has tokenized the address into individual terms: 456, opensearch, drive and so on. That means, unless you search for an individual token (like 456 or opensearch) o, op, ope , and even open won’t yield any results. One option is to use wildcards while still using the standard index analyzer for indexing:

GET standard_index_test/_search
{
  "query": {
    "wildcard": {
      "text_entry": "ope*"
    }
  }
}

The wildcard query would match “456 OpenSearch Drive Anytown, Ny 78910” but wildcard queries can be resource intensive and slow. Querying for ope* in OpenSearch results in iterating over each term in the index, bypassing optimizations of inverted index lookups. This results in higher memory usage and slower performance. To improve the performance of our query execution and search experience, we can use an index analyzer that better suits our access patterns.

Edge n-gram

The Edge n-gram tokenizer helps you find partial matches and avoids the use of wildcards by tokenizing prefixes of a single word. For example, the input word coffee is expanded into all its prefixes, c, co , cof, and so on. It can limit the prefixes to those between a minimum (min_gram) and maximum (max_gram) length. So with min_gram=3 and max_gram=5, it will expand “coffee” to cof, coff, and coffe.

Create a new index called custom_index with our own custom index analyzer that uses Edge n-grams. Set the minimum token length (min_gram) to 3 characters, and the maximum token length (max_gram) to 20 characters. The min_gram and max_gram sets the minimum and maximum returned token length respectively. You should select the min_gram and max_gram based off your access patterns. In this example, we’re searching for the term “ope” so we don’t need to set the minimum length to anything less than 3 since we’re not searching for terms like o or op. Setting the min_gram too low can lead to high latency. Likewise, we don’t need to set the maximum length to anything greater than 20 as no individual token will exceed the length of 20. Setting the maximum length to 20 gives us room to spare in case we do eventually ingest an address with a longer token length. Note, the index we are creating here is specifically for autocomplete functionality and is likely unnecessary for a general search index.

PUT custom_index
{
  "mappings": {
    "properties": {
      "text_entry": {
        "type": "text",
        "analyzer": "autocomplete",         
        "search_analyzer": "standard"       
      }
    }
  },
  "settings": {
    "analysis": {
      "filter": {
        "edge_ngram_filter": {
          "type": "edge_ngram",
          "min_gram": 3,
          "max_gram": 20
        }
      },
      "analyzer": {
        "autocomplete": {
          "type": "custom",
          "tokenizer": "standard",
          "filter": [
            "lowercase",
            "edge_ngram_filter"
          ]
        }
      }
    }
  }
}

In the above code, we created an index called custom_index with a custom analyzer named autocomplete. The analyzer performs the following:

  • It uses the standard tokenizer to split text into tokens
  • A lowercase filter is applied to lowercase all the tokens
  • The tokens are then further broken into smaller chunks based off the minimum and maximum values of the edge_ngram

The search analyzer is configured to use the standard analyzer to reduce query processing required at search time. We have already applied our custom analyzer to split the text for us upon ingestion, and we do not need to repeat this process when searching. Test how the custom analyzer analyzes the text Lexington Avenue:

GET custom_index/_analyze
{
  "analyzer": "autocomplete",
  "text": "Lexington Avenue"
}
#========
#Results
#========
# Minimum token length is 3 so we won't see l, or le
    "tokens": 
        "lex"  
        "lexi"  
        "lexin"  
        "lexing" 
        "lexingt" 
        "lexingto"    
        "lexington" 
        "ave"        
        "aven" 
        "avenu" 
        "avenue"

Notice how the tokens are lowercase and now support partial matches. Now that we’ve seen how our analyzer tokenizes our text, bulk add some data:

POST _bulk
{"index":{"_index":"custom_index"}} 
{"text_entry": "123 Amazon Street Seattle, Wa 12345 "} 
{"index":{"_index":"custom_index"}}
{"text_entry": "456 OpenSearch Drive Anytown, Ny 78910"}
{"index":{"_index":"custom_index"}}
{"text_entry": "789 Palm way Ocean Ave, Ca 33345"}
{"index":{"_index":"custom_index"}}
{"text_entry": "987 Openworld Street, Tx 48981"}

And test!

GET custom_index/_search
{
  "query": {
    "match": {
      "text_entry": {
        "query": "ope" 
      }
    }
  }
}
#========
#Results
#========
 "hits": [
      {
        "_index": "custom_index",
        "_id": "aYCEIJgB4vgFQw3LmByc",
        "_score": 0.9733556,
        "_source": {
          "text_entry": "456 OpenSearch Drive Anytown, Ny 78910"
        }
      },
      {
        "_index": "custom_index",
        "_id": "a4CEIJgB4vgFQw3LmByc",
        "_score": 0.4095239,
        "_source": {
          "text_entry": "987 Openworld Street, Tx 48981"
        }
      }
    ]

You have configured a custom n-gram analyzer to find partial words matches within our list of addresses.

Note, there is a tradeoff between using non-standard index analyzers and writing compute intensive queries. Analyzers can affect indexing throughput and increase the overall index size, especially if used inefficiently. For example, when creating the custom_index, the search analyzer was set to use the standard analyzer. Using n_grams for analysis upon ingestion and search would have impacted cluster performance unnecessarily. Additionally, we set the min_gram and max_gram to values that matched our access patterns, ensuring we didn’t create more n_grams than we needed to for our search use case. This allowed us to gain the benefits of optimizing search without impacting our ingestion throughput.

Conclusion

In this post, we changed how OpenSearch indexed our data to simplify and speed up autocomplete queries. In our case, using the Edge n-grams allowed OpenSearch to match parts of an address and yield precise results without compromising cluster performance with a wildcard query.

It’s always important to test your cluster before deploying in a production environment. Understanding your access patterns is essential to optimizing your cluster from both an indexing and searching perspective. Use the guidelines in this post as a starting point. Confirm your access patterns before creating an index, then begin experimenting with different index analyzers in a test environment to see how they can simplify your queries and improve overall cluster performance. For more reading on general OpenSearch cluster optimization techniques, refer to the Get started with Amazon OpenSearch Service: T-shirt-size your domain post.


About the authors

Rakan Kandah

Rakan Kandah

Rakan is a Solutions Architect at AWS. In his free time, Rakan enjoys playing guitar and reading.

Using Amazon SageMaker Unified Studio Identity center (IDC) and IAM-based domains together

Post Syndicated from Praveen Kumar original https://aws.amazon.com/blogs/big-data/using-amazon-sagemaker-unified-studio-identity-center-idc-and-iam-based-domains-together/

Amazon SageMaker Unified Studio now offers two domain configurations: Amazon SageMaker Unified Studio Identity Center(IDC)-based domains with comprehensive governance features, and Amazon SageMaker Unified Studio IAM-based domains with enhanced developer productivity tools.

In this post, we demonstrate how you can use both of these domain configurations of Amazon SageMaker Unified Studio using AWS Identity and Access Management (IAM) role reuse and attribute-based access control.

How authentication works in each configuration

Amazon SageMaker Unified Studio IDC-based domains authenticate users through AWS Identity and Access Management (IAM) Identity Center with Single Sign-On, preserving individual user identities throughout their sessions. These domains excel in governance with identity-based authorization, fine-grained access controls between users, and comprehensive catalog management featuring formal Publisher/Subscriber (Pub/Sub) data sharing workflows with approval processes—ideal for enterprise environments requiring strong identity management, compliance tracking, and identity-based audit trails.

Amazon SageMaker Unified Studio IAM-based domains authenticate through federated AWS Identity and Access Management (IAM) roles where all users accessing a project share the same role permissions. These domains prioritize developer productivity with modern tools including new serverless Notebooks, Athena Spark integration, the improved interface with vertical navigation, and built-in AI assistance, designed for development teams that need streamlined access and advanced analytics capabilities.

This solution facilitates organizations that are already using IDC-based domains to preserve their existing governance frameworks established in IDC-based domains while unlocking modern development capabilities for their teams through IAM-based domains. If you prefer to use the newly launched IAM-based domains, you can continue to do as well. The choice depends on your company’s needs.

Please note that at the time of writing this blog, IAM-based domains do not support Trusted identity propagation. This solution uses the project execution role to configure data access.

The challenge

Imagine a data steward (Sam) uses the IDC-based domain to define data access policies, manage the data catalog, and approve subscription requests to verify compliance and proper data governance.

On the other hand, a data engineer (Sarah), wants to use IDC-based domain for governance features such as SageMaker catalog and IAM-based domain for the new serverless Notebook to build data pipelines, perform advanced analytics, and accelerate development cycles. Sarah will request access to the data through IDC-based domain, and once access is approved by Sam, Sarah can access this data in serverless notebook available in IAM-based domain.

Solution overview

The integration leverages IAM role reuse, AWS Lake Formation Attribute-Based Access Control (ABAC) and Amazon SageMaker Catalog pub-sub model to automatically carry permissions from the IDC-based domain to the new IAM-based domain. When properly configured, data subscriptions managed through the IDC-based domain’s Pub/Sub model become immediately accessible in IAM-based domain projects, providing a unified data access experience.

The solution we will implement in the post involves creating an IAM-based domain project that is similar to your IDC consumer project (eg same team members, use case) , configuring execution roles, and enabling role reuse. This approach maintains the familiar subscription workflow while extending benefits to the IAM-based domain.The following diagram shows the high-level architecture of how this approach works.

AWS SageMaker data governance workflow diagram showing data engineer Sarah performing data discovery and exploration through SageMaker IDC and IAM domains, with data steward and owner Sam managing approvals via Business Data Catalog, connecting to Polyglot AI Notebook and SQL tools.

The solution architecture consists of:

  • Existing IDC-based domain: Contains producer and consumer projects with established data sharing via Pub/Sub model
  • IAM-based domain: New projects with federated and execution roles configured for modern development tools
  • IAM Identity Center: Manages federated access and permission sets
  • Attribute-Based Access Control: Tags on execution roles enable automatic permission inheritance

The solution provides 2 options: Option 1: IDC-Based Domain project role reuse provides the simplest integration path by directly reusing the existing consumer project IAM role from your IDC-based domain as the execution role in the IAM-based domain. The primary benefits include simplified setup requiring only policy changes (covered later in the blog), reduced administrative overhead with one less role to manage and lower risk of misconfiguration since you’re leveraging proven, existing roles. Choose Option 1 when you want the fastest implementation path, your organization prefers minimal role proliferation, you have well-established IDC-based domain roles that already have data access permissions, or your team has limited IAM expertise and wants to avoid complex tagging configurations.

Option 2: Creating a new execution role for the IAM-based domain project and use attribute-based access control (ABAC) through tagging with the IDC-based domain project ID. The key benefits include enhanced auditability with two distinct roles (one for IDC-based domain, one for IAM-based domain), clear separation showing which domain generated each request in CloudTrail logs, greater flexibility to customize permissions specific to IAM-based domain needs without affecting IDC-based domain operations, and better security isolation between the two domain types. The `AmazonDatazoneProject` tag enables attribute based access control, while maintaining distinct role identities. Choose Option 2 when: your organization requires detailed audit trails distinguishing between domain types, compliance policies mandate separation of concerns between governance and development environments, you want to track and attribute costs separately for each domain, or you need to provide evidence showing which domain (governance vs. development) accessed specific data resources for compliance reporting.

Here is the high-level view of how the identity and domain entities map to each other for both options:

AWS IAM Identity Center integration with Amazon SageMaker diagram showing access flow from IdC Groups through Permission Sets to AWS SSO IAM Roles, connecting to SageMaker domains with two implementation options: Option 1 using identical IAM roles, or Option 2 using project-tagged execution roles

Prerequisites

To follow along with this post, you should have:

For this demonstration, we use a simplified setup with a sales producer project and a marketing consumer project that subscribes to these tables.

Understanding the current IDC-based domain setup

Our starting point includes a well-established Amazon SageMaker Unified Studio IDC-based domain structure:

Sales Producer Project

  • Contains a database with pipeline and sales tables
  • Managed by Sam, the data steward who creates and publishes data assets
  • Has its own project IAM role

Marketing Consumer Project

  • Managed by Sarah, the data engineer who subscribes to published data via IDC domain project
  • Has its own project IAM role
  • Successfully queries subscribed data through the IDC-based domain interface

Each project has an associated IAM role that governs access to data assets, and the Pub/Sub model manages subscription workflows and permissions.

Setting up federated role through permission sets

Federated roles through permission sets are used to authenticate and provide users with console access to IAM-based domains through AWS IAM Identity Center, where all users within a project share the same role permissions. When you assign a permission set, IAM Identity Center creates corresponding IAM Identity Center-controlled IAM role in AWS account, and attaches the policies specified in the permission set to that role.

IAM-based SMUS domains enable streamlined access to modern development tools (serverless Notebooks, Athena Spark, AI assistance) while maintaining governance, automatically propagating permissions across domains without requiring duplicate access approvals, and simplifying team member onboarding.You can use any IAM role to access IAM-based domain. For this post, we will use federated role option using AWS IAM Identity Center (IDC).

Grant access to Data engineer group for IAM-based domains in Identity Center

1) Set up federated role in AWS IAM Identity Center

Navigate to IAM Identity Center (IDC) in the AWS Management Console, then complete the following steps:

  1. Go to permission set section in IDC. Create a new permission set called Marketing-federated-role and select Attach Policy.

AWS IAM Identity Center console screenshot displaying the marketing-federated-role permission set configuration page with provisioned status, 1-hour session duration, and empty AWS managed and customer managed policy sections with attach policy options.

  1. Search for SageMakerStudioUserIAMConsolePolicy in the existing policy name from list and select SageMakerStudioUserIAMConsolePolicy from the list. Note that the managed policy SageMakerStudioUserIAMConsolePolicy must be attached or have the same permissions added via another policy to be able to access projects in a SageMaker IAM domain.

AWS IAM Identity Center console screenshot showing AWS managed policies section with one attached SageMakerStudioUserIAMConsolePolicy and empty customer managed policies section with detach and attach policy options available.

  1. Go to the AWS account section of IDC.
  2. Assign the created permission set to your AWS account.

AWS IAM Identity Center console screenshot showing AWS accounts page in hierarchy view with organization o-9svtz1aavh, displaying Root organizational unit containing AWS account n.com with marketing-federated-role permission set assigned and assign users or groups option.

  1. For this post we assigned the permission set to marketing group, As a best practice, you should setup and grant access to groups rather than individual users.

AWS IAM Identity Center console screenshot showing marketing group details page with AWS accounts tab selected, displaying one AWS account access (management account amazon.com) with marketing-federated-role permission set applied.

  1. Add Sarah to marketing group.

AWS IAM Identity Center console screenshot showing marketing group's Users tab with one enabled member (user sarah, Display name: Sarah M) who inherits permissions to AWS accounts and Identity Center enabled applications.

This creates a federated role that Sarah can use to access the IAM-based domain. The federated role appears as an IAM role within your account and serves as the entry point for console access.

Setting up IAM-based domain execution role

There are 2 options to setup execution role for IAM-based domain project. The execution role has a one-to-one mapping with the federated role.

Option 1 – IDC-based domain Project Role reuse

Instead of creating a new execution role and tagging it, you can configure the IAM-based domain project to directly reuse the consumer project IAM role from the IDC-based domain as the execution role. This option only needs policy changes to the consumer project IAM role. To find the IDC-based domain consumer project IAM role:

  1. Navigate to the Amazon SageMaker Unified Studio IDC-based domain portal.
  2. Open the Marketing Consumer Project.
  3. Copy the project role ARN from the project overview page.

Amazon DataZone project overview page displaying marketing-project details with active status, project ID 4tcycvm4c684rt, domain ID dzd-47supbt0i3jysp, All capabilities profile, Corp domain unit, Amazon S3 location in us-east-2, and project role ARN with up-to-date status.

  1. You will need to modify this execution role’s policy with detailed instructions provided later in the blog.

Setting up IAM-based domain project for option 1

To create an IAM-based domain project that will integrate with your existing IDC-based domain permissions, complete the following steps:

  1. Log in to the AWS Console using IAM-based domain administrator.
  2. Navigate to Amazon SageMaker page within console.
  3. Choose Open.

Amazon SageMaker landing page displaying "The center for data, analytics, and AI" with tagline about next-generation integrated analytics experience, serverless notebooks with built-in AI Agent, Amazon DataZone integration note, and call-to-action panel featuring "Get started with Amazon SageMaker Unified Studio" with Open button and View existing domains

  1. Once logged in to IAM-based domain as admin, choose Manage projects.

Amazon SageMaker admin-project dashboard displaying left navigation menu with data analytics and AI/ML sections, quick-start cards for exploring data, building in notebooks, and discovering ML models, plus four sample data project templates: Customer usage analysis (3 mins), Customer segmentation (8 mins), Customer churn prediction (5 mins), and Retail sales forecasting (20 mins).

  1. Next, click on Create Project.

Amazon DataZone Domain Administration Projects page showing "Projects (3)" with description about enabling IAM role-based access to AWS Analytics and AI/ML tools, search functionality to find projects, last refreshed timestamp, and green Create project button.

  1. Enter project name as “Marketing Consumer Project”.

Amazon DataZone Create project dialog showing Step 1 "Enter Details" with required Project name field containing "Marketing Consumer Project" (1-64 characters, a-z, A-Z, 0-9, spaces, dashes, underscores allowed) and optional Description field with 0/2048 character count, followed by Step 2 "Assign roles".

  1. During project creation, select the following crucial roles and then choose Create Project:
  • Project IAM Role: The marketing federated role created in IAM Identity Center above. This is the role in the member account that has a role name with suffix AWSReservedSSO.
  • Project Role: – Choose project role for data engineer, copied from option 1.

Amazon SageMaker Unified Studio Create project dialog showing IAM role configuration with AWSReservedSSO_marketing-federated-role selected, blue alert requiring SageMakerStudioUserIAMConsolePolicy attachment, Execution role section with "Use an existing role" option selected, and datazone_usr_role_4tcycvm4c684rt_ajtckkwo2fnhyh IAM role specified with note that role is not editable after project creation

  1. Make policy changes to this project role as per the instruction on the SMUS UI page.

Amazon SageMaker Unified Studio role selection interface showing "Use an existing role" option selected with IAM role datazone_usr_role_4tcycvm4c684rt_ajtckkwo2fnhyh, blue information box displaying required permissions including SageMakerStudioUserIAMDefaultExecutionPolicy managed policy, trust policy enabling Amazon SageMaker Unified Studio service assumption, and inline policy for role pass-through, with note that role is not editable after project creation.

Option 2 – Bring your own execution role. 

To create an IAM-based domain project that will integrate with your existing IDC-based domain permissions., you must tag the execution role for permission propagation. Amazon SageMaker Catalog and AWS Lake Formation use attribute-based access control, which means permissions can be inherited based on resource tags. For this option, you will need consumer project ID.To find the IDC-based domain consumer project ID:

  1. Navigate to the Amazon SageMaker Unified Studio IDC-based domain portal.
  2. Open the Marketing Consumer Project.
  3. Copy the project ID from the project details.

Amazon SageMaker Unified Studio marketing-project overview page displaying navigation breadcrumb (Home > Projects > marketing-project > Project overview), left sidebar menu with Project overview, Data, Compute, Members, and Project catalog sections, Project files section listing 3 JupyterLab files (.libs.json, README.md, getting_started.ipynb) last modified November 18, 2025, Readme section with Welcome heading describing SageMaker Unified Studio, and Project details tab showing project name, ID, last modified date November 21, 2025, and Amazon S3 location.” width=”2196″ height=”1164″></p>
<h3>Setting up IAM-based domain project for option 2</h3>
<p>Complete the following steps:</p>
<ol>
<li>Create another project with name “Marketing Consumer Project 2” in the IAM-based domain while logged in as admin.</li>
<li>During project creation, select the following roles:
<ol type=

  • Federated Role: The marketing federated role created in IAM Identity Center above.
  • Execution Role: – Choose execution role from option 2.
  • Make policy changes to this execution role as per the instruction.
  • Amazon SageMaker Unified Studio role selection interface showing "Use an existing role" option selected with IAM role field containing "sagemaker-marketing-execution-role", blue information box displaying required permissions including SageMakerStudioUserIAMDefaultExecutionPolicy managed policy, trust policy enabling Amazon SageMaker Unified Studio and related services to assume the role, and inline policy allowing role pass-through to other services, with note that role is not editable after project creation

    1. Next, navigate to the IAM console and locate the execution role created for your IAM-based domain consumer project.
    2. Add the following tag, this step relies on ABAC policies with projectId for subscriptions.
    • Key: AmazonDatazoneProject
    • Value: The project ID from your Amazon SageMaker Unified Studio IDC-based domain consumer project

    AWS IAM console displaying sagemaker-marketing-execution-role details page with Summary section showing creation date November 18, 2025, last activity 3 days ago, ARN arn:aws:iam::role/sagemaker-marketing-execution-role, 1-hour maximum session duration, five tabs (Permissions, Trust relationships, Tags (1), Last Accessed, Revoke sessions), and Tags section displaying one tag with Key "AmazonDataZoneProject" and Value "4tcycvm4c684rt" with Delete, Edit, and Manage tags buttons available.

    This tag configuration results in data access grant from IDC-based domain consumer project to the IAM-based domain project execution role.

    Verify data access in the IAM-based domain

    After tagging the execution role, verify that permissions are set up correctly.Complete the following steps:

    1. Use the SSO URL to log into the SSO Identity Center as Sarah.

    AWS IAM Identity Center Dashboard displaying left navigation menu with Dashboard, Users, Groups, Settings, Multi-account permissions (AWS accounts, Permission sets), and Application assignments sections; central management panel showing service control policies guidance with yellow warning banner about member account instances and CloudTrail monitoring section; IAM Identity Center setup area with three action cards for confirming identity source, managing multi-account permissions, and setting up application assignments; right panel Settings summary showing Identity Center directory as identity source, us-east-2 region, organization ID o-9svtz1aavh, AWS access portal URL, and issuer URL; What's new section highlighting customer-managed KMS keys support and Amazon SageMaker Studio user background sessions; Related consoles links to CloudTrail, AWS Organizations, and IAM.

    1. Open the AWS console using federated role created earlier in setting federated role section.
    2. Navigate to Amazon SageMaker.
    3. Choose Amazon SageMaker Unified Studio IAM-based domain option (this will show up if project is already created with federated role).

    Amazon SageMaker Unified Studio marketing-project dashboard displaying left navigation menu with Overview, Files, Data, Connections, Code (Notebooks, JupyterLab), Data analytics (Query Editor, Visual ETL, Data processing jobs), and AI/ML sections (Models, MLflow, Training jobs, Inference endpoints); main content area showing "Jump into your data and models" with three quick-start cards (Explore your data, Build in the notebook, Discover ML models) and four sample data projects: Retail sales forecasting (20 mins), Customer churn prediction (5 mins), Customer segmentation (8 mins), and Customer usage analysis (3 mins); top-right panel displaying account details with us-east-2 region, federated user aws-reserved/sarah, and execution role sagemaker-marketing-execution-role.

    1. In the Amazon SageMaker Unified Studio IAM-based domain project, navigate to the Data tab. If you created 2 projects with both option 1 and option 2 execution role, then 2 projects will show up and you can login to either to validate data access.

    Amazon SageMaker Unified Studio data explorer interface displaying SQL query "SELECT * FROM glue_db_6doxdp1wuy165l.sales_table LIMIT 100" executed via Athena in 6 seconds, showing six columns (ord_num, sales_qty_sld, wholesale_cost, lst_pr, sell_pr, disnt) with green distribution histograms above data preview table containing six sample sales records with order numbers ranging from 46776931 to 146776932, left navigation showing AwsDataCatalog database structure with glue_db_6doxdp1wuy165l containing pipeline_table and sales_table, last saved 2 minutes ago.

    1. Verify that the consumer database and subscribed tables appear.

    Create and use the new serverless notebooks

    With permissions properly configured, you can now use IAM-based domain capabilities like serverless Notebooks. Complete the following steps:

    1. In the Amazon SageMaker Unified Studio IAM-based domain project, select a table from the Data tab.
    2. Choose Create notebook.
    3. The Notebook opens with Athena SQL as the default cell type.
    4. Write and run queries against your subscribed data.

    Amazon SageMaker Unified Studio marketing-project notebook displaying sales_table data from 2025-11-18 21:42:01, left Data explorer showing AwsDataCatalog with glue_db_6doxdp1wuyi65l database containing pipeline_table and sales_table, main data table showing 11 rows with columns (ord_num, sales_qty_sld, wholesale_cost, lst_pr, sell_pr, disnt) displaying rows 4-9 on page 1 of 2, Python PySpark SQL query "SELECT * FROM 'glue_db_6doxdp1wuyi65l'.'pipeline_table' LIMIT 100" executed in 27 seconds, and Filters section displaying distribution histograms for all numerical columns.

    The notebook runs with the execution role’s permissions, which now include access to all data subscribed through the IDC-based domain.

    Key benefits of this integration

    This integration approach delivers several important advantages:

    Preserve existing investments

    • Continue using IDC-based domain governance and catalogs.
    • Maintain established Pub/Sub workflows.
    • No migration required for existing data assets.

    Get modern capabilities

    • Provide developers with the new serverless Notebooks.
    • Access Athena Spark for advanced analytics.
    • Provides improved user experience and navigation.

    Simplified permission management

    • Single subscription workflow manages access across both domains.
    • Consistent data access via role reuse and attribute-based access control.
    • No duplicate access requests or approvals needed.

    Unified data experience

    • Developers access all subscribed data from one interface.
    • Consistent data catalog across domains.
    • Simplified onboarding for new team members.

    Cleanup

    Complete the following steps to delete the resources you created:

    1. Delete the serverless Notebooks created in the IAM-based domain projects.
    2. Delete the IAM-based domain projects (Marketing Consumer Project and Marketing Consumer Project 2).
    3. Remove the permission set assignment from marketing group in IAM Identity Center.
    4. Delete the Marketing-federated-role permission set in IAM Identity Center.
    5. Remove the tags (AmazonDatazoneProject) from the execution role (if using Option 2).
    6. Delete the execution role created for the IAM-based domain (if using Option 2 and not reusing the IDC-based domain project role).
    7. Revert any policy changes made to the IDC-based domain consumer project IAM role (if using Option 1).
    8. If you do not need the IAM-based domain anymore, delete it.
    9. If you created any test data subscriptions in the IDC-based domain, remove them.

    Conclusion

    In this post, we demonstrated how to access Amazon SageMaker Unified Studio IDC-based domain with the new IAM-based domain using role reuse and attribute-based access control. This setup offers data engineers the best of both worlds: access to specialized modern development tools—including the new serverless Notebooks, Athena Spark integration, and built-in AI assistance , while maintaining proper governance that includes comprehensive catalog management and robust security controls established in the IDC-based domain.You can now confidently adopt Amazon SageMaker Unified Studio IAM-based domain capabilities knowing their established data governance, subscription workflows, and access controls remain intact and continue to function as expected.

    Ready to get started with Amazon SageMaker Unified Studio and unlock the power of integrated governance and modern development tools for your organization? Visit the Amazon SageMaker Unified Studio documentation to learn more and begin your implementation today.


    About the authors

    Praveen Kumar

    Praveen Kumar

    Praveen is a Principal Analytics Solutions Architect at AWS with expertise in designing, building, and implementing modern data and analytics platforms using cloud-based services. His areas of interest are serverless technology, data governance, and data-driven AI applications.

    Durga Mishra

    Durga Mishra

    Durga is a Principal Data and AI solutions architecture strategist at AWS . Outside of work, Durga enjoys building new things and spending time with family. He loves to hike on Appalachian trails and spend time in nature.

    Joel

    Joel Farvault

    Joel is a Principal Specialist SA Analytics for AWS with 25 years’ experience working on enterprise architecture, data governance and analytics. He uses his experience to advise customers on their data strategy and technology foundations.

    author name

    Satish Sarapuri

    Satish is a Sr. Data Architect for Data Mesh/Data Lake/Gen AI at AWS. He helps enterprise-level customers build generative AI, data mesh, data lake, and analytics platform solutions on AWS to help them make data-driven decisions and gain impactful outcomes for their business. In his spare time, he enjoys trail running and spending quality time with his family.

    author name

    Leonardo Gomez

    Leonardo is a Principal Analytics Specialist Solutions Architect at AWS. He has over a decade of experience in data management, helping customers around the globe address their business and technical needs.

    Orchestrate end-to-end scalable ETL pipeline with Amazon SageMaker workflows

    Post Syndicated from Shubham Kumar original https://aws.amazon.com/blogs/big-data/orchestrate-end-to-end-scalable-etl-pipeline-with-amazon-sagemaker-workflows/

    Amazon SageMaker Unified Studio serves as a collaborative workspace where data engineers and scientists can work together on end-to-end data and machine learning (ML) workflows. SageMaker Unified Studio specializes in orchestrating complex data workflows across multiple AWS services through its integration with Amazon Managed Workflows for Apache Airflow (Amazon MWAA). Project owners can create shared environments where team members jointly develop and deploy workflows, while maintaining oversight of pipeline execution. This unified approach makes sure data pipelines run consistently and efficiently, with clear visibility into the entire process, making it seamless for teams to collaborate on sophisticated data and ML projects.

    This post explores how to build and manage a comprehensive extract, transform, and load (ETL) pipeline using SageMaker Unified Studio workflows through a code-based approach. We demonstrate how to use a single, integrated interface to handle all aspects of data processing, from preparation to orchestration, by using AWS services including Amazon EMR, AWS Glue, Amazon Redshift, and Amazon MWAA. This solution streamlines the data pipeline through a single UI.

    Example use case: Customer behavior analysis for an ecommerce platform

    Let’s consider a real-world scenario: An e-commerce company wants to analyze customer transactions data to create a customer summary report. They have data coming from multiple sources:

    • Customer profile data stored in CSV files
    • Transaction history in JSON format
    • Website clickstream data in semi-structured log files

    The company wants to do the following:

    • Extract data from these sources
    • Clean and transform the data
    • Perform quality checks
    • Load the processed data into a data warehouse
    • Schedule this pipeline to run daily

    Solution overview

    The following diagram illustrates the architecture that you implement in this post.

    This architecture diagram illustrates a comprehensive, end-to-end data processing pipeline built on AWS services, orchestrated through Amazon SageMaker Unified Studio. The pipeline demonstrates best practices for data ingestion, transformation, quality validation, advanced processing, and analytics.

    The workflow consists of the following steps:

    1. Establish a data repository by creating an Amazon Simple Storage Service (Amazon S3) bucket with an organized folder structure for customer data, transaction history, and clickstream logs, and configure access policies for seamless integration with SageMaker Unified Studio.
    2. Extract data from the S3 bucket using AWS Glue jobs.
    3. Use AWS Glue and Amazon EMR Serverless to clean and transform the data.
    4. Implement data quality validation using AWS Glue Data Quality.
    5. Load the processed data into Amazon Redshift Serverless.
    6. Create and manage the workflow environment using SageMaker Unified Studio with Identity Center–based domains.

    Note: Amazon SageMaker Unified Studio supports two domain configuration models: IAM Identity Center (IdC)–based domains and IAM role–based domains. While IAM-based domains enable role-driven access management and visual workflows, this post specifically focuses on Identity Center–based domains, where users authenticate via IdC and projects access data and resources using project roles and identity-based authorization.

    Prerequisites

    Before beginning, ensure you have the following resources:

    Configure Amazon SageMaker Unified Studio domain

    This solution requires SageMaker Unified Studio domain in the us-east-1 AWS Region. Although SageMaker Unified Studio is available in multiple Regions, this post uses us-east-1 for consistency. For a complete list of supported Regions, refer to Regions where Amazon SageMaker Unified Studio is supported.

    Complete the following steps to configure your domain:

    1. Sign in to the AWS Management Console, navigate to Amazon SageMaker, and open the Domains section from the left navigation pane.
    2. On the SageMaker console, choose Create domain, then choose Quick setup.
    3. If the message “No VPC has been specifically set up for use with Amazon SageMaker Unified Studio” appears, select Create VPC. The process redirects to an AWS CloudFormation stack. Leave all settings at their default values and select Create stack.
    4. Under Quick setup settings, for Name, enter a domain name (for example, etl-ecommerce-blog-demo). Review the selected configurations.
    5. Choose Continue to proceed.
    6. On the Create IAM Identity Center user page, create an SSO user (account with IAM Identity Center) or select an existing SSO user to log in to the Amazon SageMaker Unified Studio. The SSO selected here is used as the administrator in the Amazon SageMaker Unified Studio.
    7. Choose Create domain.

    For detailed instructions, see Create a SageMaker domain and Onboarding data in Amazon SageMaker Unified Studio.

    # Amazon SageMaker Domain Details Interface This screenshot shows the Amazon SageMaker domain details page for "etl-ecommerce-blog-demo

    After you have created a domain, popup will appear with the message: “Your domain has been created! You can now log in to Amazon SageMaker Unified Studio”. You can close the popup for now.

    Create a project

    In this section, we create a project to serve as a collaborative workspace for teams to work on business use cases. Complete the following steps:

    1. Choose Open Unified Studio and sign in with your SSO credentials using the Sign in with SSO option.
    2. Choose Create project.
    3. Name the project (for example, ETL-Pipeline-Demo) and create it using the All capabilities project profile.
    4. Choose Continue.
    5. Keep the default values for the configuration parameters and choose Continue.
    6. Choose Create project.

    Project creation might take a few minutes. After the project is created, the environment will be configured for data access and processing.

    Integrate S3 bucket with SageMaker Unified Studio

    To enable external data processing within SageMaker Unified Studio, configure integration with an S3 bucket. This section walks through the steps to set up the S3 bucket, configure permissions, and integrate it with the project.

    Create and configure S3 bucket

    Complete the following steps to create your bucket:

    1. In a new browser tab, open the AWS Management Console and search for S3.
    2. On the Amazon S3 console, choose Create Bucket .
    3. Create a bucket named ecommerce-raw-layer-bucket-demo-<Account-ID>-us-east-1. For detailed instructions, see create a general-purpose Amazon S3 bucket for storage.
    4. Create the following folder structure in the bucket. For detailed instructions, see Creating a folder:
      • raw/customers/
      • raw/transactions/
      • raw/clickstream/
      • processed/
      • analytics/

    Upload sample data

    In this section, we upload sample ecommerce data that represents a typical business scenario where customer behavior, transaction history, and website interactions need to be analyzed together.

    The raw/customers/customers.csv file contains customer profile information, including registration details. This structured data will be processed first to establish the customer dimension for our analytics.

    customer_id,name,email,registration_date
    1,John Doe,[email protected],2022-01-15
    2,Jane Smith,[email protected],2022-02-20
    3,Robert Johnson,[email protected],2022-01-30
    4,Emily Brown,[email protected],2022-03-05
    5,Michael Wilson,[email protected],2022-02-10

    The raw/transactions/transactions.json file contains purchase transactions with nested product arrays. This semi-structured data will be flattened and joined with customer data to analyze purchasing patterns and customer lifetime value.

    [
    {"transaction_id": "t1001", "customer_id": 1, "amount": 125.99, "date": "2023-01-10", "items": ["product1", "product2"]},
    {"transaction_id": "t1002", "customer_id": 2, "amount": 89.50, "date": "2023-01-12", "items": ["product3"]},
    {"transaction_id": "t1003", "customer_id": 1, "amount": 45.25, "date": "2023-01-15", "items": ["product2"]},
    {"transaction_id": "t1004", "customer_id": 3, "amount": 210.75, "date": "2023-01-18", "items": ["product1", "product4", "product5"]},
    {"transaction_id": "t1005", "customer_id": 4, "amount": 55.00, "date": "2023-01-20", "items": ["product3", "product6"]}
    ]

    The raw/clickstream/clickstream.csv file captures user website interactions and behavior patterns. This time-series data will be processed to understand customer journey and conversion funnel analytics.

    timestamp,customer_id,page,action
    2023-01-10T10:15:23,1,homepage,view
    2023-01-10T10:16:45,1,product_page,view
    2023-01-10T10:18:12,1,product_page,add_to_cart
    2023-01-10T10:20:30,1,checkout,view
    2023-01-10T10:22:15,1,checkout,purchase
    2023-01-12T14:30:10,2,homepage,view
    2023-01-12T14:32:20,2,product_page,view
    2023-01-12T14:35:45,2,product_page,add_to_cart
    2023-01-12T14:40:12,2,checkout,view
    2023-01-12T14:42:30,2,checkout,purchase

    raw

    For detailed instructions on uploading files to Amazon S3, refer to the Uploading objects.

    Configure CORS policy

    To allow access from the SageMaker Unified Studio domain portal, update the Cross-Origin Resource Sharing (CORS) configuration of the bucket:

    1. On the bucket’s Permissions tab, choose Edit under Cross-origin resource sharing (CORS).
      permission
    2. Enter the following CORS policy and replace domainUrl with the SageMaker Unified Studio domain URL (for example, https://<domain-id>.sagemaker.us-east-1.on.aws ). The URL can be found at the top of the domain details page on the SageMaker Unified Studio console.
      [
          {
              "AllowedHeaders": [
                  "*"
              ],
              "AllowedMethods": [
                  "PUT",
                  "GET",
                  "POST",
                  "DELETE",
                  "HEAD"
              ],
              "AllowedOrigins": [
                  "domainUrl"
              ],
              "ExposeHeaders": [
                  "x-amz-version-id"
              ]
          }
      ]

    For detailed information, see Adding Amazon S3 data and gain access using the project role.

    Grant Amazon S3 access to SageMaker project role

    To enable SageMaker Unified Studio to access the external Amazon S3 location, the corresponding AWS Identity and Access Management (IAM) project role must be updated with the required permissions. Complete the following steps:

    1. On the IAM console, choose Roles in the navigation pane.
    2. Search for the project role using the last segment of the project role Amazon Resource Name (ARN). This information is located on the Project overview page in SageMaker Unified Studio (for example, datazone_usr_role_1a2b3c45de6789_abcd1efghij2kl).
      project detail
    3. Choose the project role to open the role details page.
    4. On the Permissions tab, choose Add permissions, then choose Create inline policy.
    5. Use the JSON editor to create a policy that grants the project role access to the Amazon S3 location
    6. In the JSON policy below, replace the placeholder values with your actual environment details:
      • Replace <BUCKET_PREFIX> with the prefix of S3 bucket name (for example, ecommerce-raw-layer)
      • Replace <AWS_REGION> with the AWS Region where your AWS Glue Data Quality rulesets are created (for example, us-east-1)
      • Replace <AWS_ACCOUNT_ID> with your AWS account ID
    7. Paste the updated JSON policy into the JSON editor.
      {
          "Version": "2012-10-17",
          "Statement": [
              {
                  "Sid": "ETLBucketListAccess",
                  "Effect": "Allow",
                  "Action": [
                      "s3:ListBucket",
                      "s3:GetBucketLocation"
                  ],
                  "Resource": "arn:aws:s3:::<BUCKET_PREFIX>-*"
              },
              {
                  "Sid": "ETLObjectAccess",
                  "Effect": "Allow",
                  "Action": [
                      "s3:GetObject",
                      "s3:PutObject",
                      "s3:DeleteObject"
                  ],
                  "Resource": "arn:aws:s3:::<BUCKET_PREFIX>-*/*"
              },
              {
                  "Sid": "GlueDataQualityPublish",
                  "Effect": "Allow",
                  "Action": [
                      "glue:PublishDataQuality"
                  ],
       "Resource":"arn:aws:glue:<AWS_REGION>:<AWS_ACCOUNT_ID>:dataQualityRuleset/*"
              }
          ]
      }

    8. Choose Next.
    9. Enter a name for the policy (for example, etl-rawlayer-access), then choose Create policy.
    10. Choose Add permissions again, then choose Create inline policy.
    11. In the JSON editor, create a second policy to manage S3 Access Grants:Replace <BUCKET_PREFIX> with the prefix of S3 bucket name (for example, ecommerce-raw-layer) and paste this JSON policy.
      {
          "Version": "2012-10-17",
          "Statement": [
              {
                  "Sid": "S3AGLocationManagement",
                  "Effect": "Allow",
                  "Action": [
                      "s3:CreateAccessGrantsLocation",
                      "s3:DeleteAccessGrantsLocation",
                      "s3:GetAccessGrantsLocation"
                  ],
                  "Resource": [
                      "arn:aws:s3:*:*:access-grants/default/*"
                  ],
                  "Condition": {
                      "StringLike": {
                          "s3:accessGrantsLocationScope": "s3://<BUCKET_PREFIX>-*/*"
                      }
                  }
              },
              {
                  "Sid": "S3AGPermissionManagement",
                  "Effect": "Allow",
                  "Action": [
                      "s3:CreateAccessGrant",
                      "s3:DeleteAccessGrant"
                  ],
                  "Resource": [
                      "arn:aws:s3:*:*:access-grants/default/location/*",
                      "arn:aws:s3:*:*:access-grants/default/grant/*"
                  ],
                  "Condition": {
                      "StringLike": {
                          "s3:accessGrantScope": "s3://<BUCKET_PREFIX>-*/*"
                      }
                  }
              }
          ]
      }

    12. Choose Next.
    13. Enter a name for the policy (for example, s3-access-grants-policy), then choose Create policy.

    create policy

    For detailed information about S3 Access Grants, see Adding Amazon S3 data.

    Add S3 bucket to project

    After you add policies to the project role for access to the Amazon S3 resources, complete the following steps to integrate the S3 bucket with the SageMaker Unified Studio project:

    1. In SageMaker Unified Studio, open the project you created under Your projects.
      your projects
    2. Choose Data in the navigation pane.
    3. Select Add and then Add S3 location.
      add s3
    4. Configure the S3 location:
      1. For Name, enter a descriptive name (for example, E-commerce_Raw_Data).
      2. For S3 URI, enter your bucket URI (for example, s3://ecommerce-raw-layer-bucket-demo-<Account-ID>-us-east-1/).
      3. For AWS Region, enter your Region (for this example, us-east-1).
      4. Leave Access role ARN blank.
      5. Click Add S3 Location
    5. Wait for the integration to complete.
    6. Verify the S3 location appears in your project’s data catalog (on the Project overview page, on the Data tab, locate the Buckets pane to view the buckets and folders).

    add

    This process connects your S3 bucket to SageMaker Unified Studio, making your data ready for analysis.

    Create notebook for job scripts

    Before you can create the data processing jobs, you must set up a notebook to develop the scripts that will generate and process your data. Complete the following steps:

    1. In SageMaker Unified Studio, on the top menu, under Build, choose JupyterLab.
    2. Choose Configure Space and choose the instance type ml.t3.xlarge. This makes sure your JupyterLab instance has at least 4 vCPUs and 4 GiB of memory.
    3. Choose Configure and Start Space or Save and Restart to launch your environment.
    4. Wait a few moments for the instance to be ready.
    5. Choose File, New, and Notebook to create a new notebook.
    6. Set Kernel as Python 3, Connection type as PySpark, and Compute as Project.spark.compatibility.
      jupyter
    7. In the notebook, enter the following script to use later for your AWS Glue job. This script processes raw data from three sources in the S3 data lake, standardizes dates, and converts data types before saving the cleaned data in Parquet format for optimal storage and querying.
    8. Replace <Bucket-Name> with the name of actual S3 bucket in script:
      import sys
      from awsglue.transforms import *
      from pyspark.context import SparkContext
      from awsglue.context import GlueContext
      from awsglue.job import Job
      from awsglue.utils import getResolvedOptions
      from pyspark.sql import functions as F
      args = getResolvedOptions(sys.argv, ['JOB_NAME'])
      sc = SparkContext.getOrCreate()
      glueContext = GlueContext(sc)
      spark = glueContext.spark_session
      job = Job(glueContext)
      job.init(args['JOB_NAME'], args)
      # Customers
      customer_df = (
          spark.read
          .option("header", "true")
          .csv("s3://<Bucket-Name>/raw/customers/")
          .withColumn("registration_date", F.to_date("registration_date"))
          .withColumn("processed_at", F.current_timestamp())
      )
      customer_df.write.mode("overwrite").parquet(
          "s3://<Bucket-Name>/processed/customers/"
      )
      # Transactions
      transaction_df = (
          spark.read
          .json("s3://<Bucket-Name>/raw/transactions/")
          .withColumn("date", F.to_date("date"))
          .withColumn("customer_id", F.col("customer_id").cast("int"))
          .withColumn("processed_at", F.current_timestamp())
      )
      transaction_df.write.mode("overwrite").parquet(
          "s3://<Bucket-Name>/processed/transactions/"
      )
      # Clickstream
      clickstream_df = (
          spark.read
          .option("header", "true")
          .csv("s3://<Bucket-Name>/raw/clickstream/")
          .withColumn("customer_id", F.col("customer_id").cast("int"))
          .withColumn("timestamp", F.to_timestamp("timestamp"))
          .withColumn("processed_at", F.current_timestamp())
      )
      clickstream_df.write.mode("overwrite").parquet(
          "s3://<Bucket-Name>/processed/clickstream/"
      )
      print("Data processing completed successfully")
      job.commit()

      This script processes customer, transaction, and clickstream data from the raw layer in Amazon S3 and saves it as Parquet files in the processed layer.

    9. Choose File, Save Notebook As, and save the file as shared/etl_initial_processing_job.ipynb.
      jupyter2

    Create notebook for AWS Glue Data Quality

    After you create the initial data processing script, the next step is to set up a notebook to perform data quality checks using AWS Glue. These checks help validate the integrity and completeness of your data before further processing. Complete the following steps:

    1. Choose File, New, and Notebook to create a new notebook.
    2. Set Kernel as Python 3, Connection type as PySpark, and Compute as Project.spark.compatibility.
      select-kernel
    3. In this new notebook, add the data quality check script using the AWS Glue EvaluateDataQuality method. Replace <Bucket-Name> with the name of actual S3 bucket in script:
      from datetime import datetime
      from pyspark.context import SparkContext
      from awsglue.context import GlueContext
      from awsglue.job import Job
      from awsgluedq.transforms import EvaluateDataQuality
      from awsglue.transforms import SelectFromCollection
      
      # ---------------- Glue setup ----------------
      sc = SparkContext.getOrCreate()
      glueContext = GlueContext(sc)
      job = Job(glueContext)
      job.init("GlueDQJob", {})
      
      # ---------------- Constants ----------------
      RUN_DATE = datetime.utcnow().strftime("%Y-%m-%d")
      year, month, day = RUN_DATE.split("-")
      OUTPUT_PATH = "s3://<Bucket-Name>/data-quality-results"
      
      # ---------------- Tables and Rules ----------------
      tables = {
          "customers": ["s3://<Bucket-Name>/processed/customers/",
                        ["IsComplete \"customer_id\"", "IsUnique \"customer_id\"", "IsComplete \"email\""]],
          "transactions": ["s3://<Bucket-Name>/processed/transactions/",
                           ["IsComplete \"transaction_id\"", "IsUnique \"transaction_id\""]],
          "clickstream": ["s3://<Bucket-Name>/processed/clickstream/",
                          ["IsComplete \"customer_id\"", "IsComplete \"action\""]]
      }
      
      # ---------------- Process Each Table ----------------
      for table, (path, rules) in tables.items():
          df = glueContext.create_dynamic_frame.from_options("s3", {"paths":[path]}, "parquet")
          results = EvaluateDataQuality().process_rows(
              frame=df,
              ruleset=f"Rules = [{', '.join(rules)}]",
              publishing_options={"dataQualityEvaluationContext": table}
          )
          rows = SelectFromCollection.apply(results, key="rowLevelOutcomes", transformation_ctx="rows").toDF()
          rows = rows.drop("DataQualityRulesPass", "DataQualityRulesFail", "DataQualityRulesSkip")
      
          # Write passed/failed rows
          for status, colval in [("pass","Passed"), ("fail","Failed")]:
              tmp = rows.filter(rows.DataQualityEvaluationResult.contains(colval))
              if tmp.count() > 0:
                  tmp.write.mode("append").parquet(
              f"{OUTPUT_PATH}/{table}/status=dq_{status}/Year={year}/Month={month}/Date={day}"
                  )
      print("Data Quality checks completed and written to S3")
      job.commit()

    4. Choose File, Save Notebook As, and save the file as shared/etl_data_quality_job.ipynb.

    Create and test AWS Glue jobs

    Jobs in SageMaker Unified Studio enable scalable, flexible ETL pipelines using AWS Glue. This section walks through creating and testing data processing jobs for efficient and governed data transformation.

    Create initial data processing job

    This job performs the first processing job in the ETL pipeline, transforming raw customer, transaction, and clickstream data and writing the cleaned output to Amazon S3 in Parquet format. Complete the following steps to create the job:

    1. In SageMaker Unified Studio, go to your project.
    2. On the top menu, choose Build, and under Data Analysis & Integration, choose Data processing jobs.
      navbar on smus
    3. Choose Create job from notebooks.
    4. Under Choose project files, choose Browse files.
    5. Locate and select etl_initial_processing_job.ipynb (the notebook saved earlier in JupyterLab), then choose Select and Next.
      select the notebook
    6. Configure the job settings:
      1. For Name, enter a name (for example, job-1).
      2. For Description, enter a description (for example, Initial ETL job for customer data processing).
      3. For IAM Role, choose the project role (default).
      4. For Type, choose Spark.
      5. For AWS Glue version, use version 5.0.
      6. For Language, choose Python.
      7. For Worker type, use G.1X.
      8. For Number of Instances, set to 10.
      9. For Number of retries, set to 0.
      10. For Job timeout, set to 480.
      11. For Compute connection, choose project.spark.compatibility.
      12. Under Advanced settings, turn on Continuous logging.

      advanced setting

    7. Leave the remaining settings as default, then choose Submit.

    After the job is created, a confirmation message will appear indicating that job-1 was created successfully.

    Create AWS Glue Data Quality job

    This job runs data quality checks on the transformed datasets using AWS Glue Data Quality. Rulesets validate completeness and uniqueness for key fields. Complete the following steps to create the job:

    1. In SageMaker Unified Studio, go to your project.
    2. On the top menu, choose Build, and under Data Analysis & Integration, choose Data processing jobs.
    3. Choose Create job, Code-based job, and Create job from files.
    4. Under Choose project files, choose Browse files.
    5. Locate and select etl_glue_data_quality.ipynb, then choose Select and Next.
    6. Configure the job settings:
    7. For Name, enter a name (for example, job-2).
    8. For Description, enter a description (for example, Data quality checks using AWS Glue Data Quality).
    9. For IAM Role, choose the project role.
    10. For Type, choose Spark.
    11. For AWS Glue version, use version 5.0.
    12. For Language, choose Python.
    13. For Worker type, use G.1X.
    14. For Number of Instances, set to 10.
    15. For Number of retries, set to 0.
    16. For Job timeout, set to 480.
    17. For Compute connection, choose project.spark.compatibility.
    18. Under Advanced settings, turn on Continuous logging.
    19. Leave the remaining settings as default, then choose Submit.

    After the job is created, a confirmation message will appear indicating that job-2 was created successfully.

    Test AWS Glue jobs

    Test both jobs to make sure they execute successfully:

    1. In SageMaker Unified Studio, go to your project.
    2. On the top menu, choose Build, and under Data Analysis & Integration, choose Data processing jobs.
    3. Select job-1 and choose Run job.
    4. Monitor the job execution and verify it completes successfully.
    5. Similarly, select job-2 and choose Run job.
    6. Monitor the job execution and verify it completes successfully.

    Add EMR Serverless compute

    In the ETL pipeline, we use EMR Serverless to perform compute-intensive transformations and aggregations on large datasets. It automatically scales resources based on workload, offering high performance with simplified operations. By integrating EMR Serverless with SageMaker Unified Studio, you can simplify the process of running Spark jobs interactively using Jupyter notebooks in a serverless environment.

    This section walks through the steps to configure EMR Serverless compute within SageMaker Studio and use it for executing distributed data processing jobs.

    Configure EMR Serverless in SageMaker Unified Studio

    To use EMR Serverless for processing in the project, follow these steps:

    1. In the navigation pane on Project Overview, choose Compute.
    2. On the Data processing tab, choose Add compute and Create new compute resources.
    3. Select EMR Serverless and choose Next.
    4. Configure EMR Serverless settings:
    5. For Compute name, enter a name (for example, etl-emr-serverless).
    6. For Description, enter a description (for example, EMR Serverless for advanced data processing).
    7. For Release label, choose emr-7.8.0.
    8. For Permission mode, choose Compatibility.
    9. Choose Add Compute to complete the setup.

    After it’s configured, the EMR Serverless compute will be listed with the deployment status Active.

    emr serverless

    Create and run notebook with EMR Serverless

    After you create the EMR Serverless compute, you can run PySpark-based data transformation jobs using a Jupyter notebook to perform large-scale data transformations. This job reads cleaned customer, transaction, and clickstream datasets from Amazon S3, performs aggregations and scoring, and writes the final analytics outputs back to Amazon S3 in both Parquet and CSV formats.Complete the following steps to create a notebook for EMR Serverless processing:

    1. On the top menu, under Build, choose JupyterLab.
    2. Choose File, New, and Notebook.
    3. Set Kernel as Python 3, Connection type as PySpark, and Compute as emr-s.etl-emr-serverless.
      compute
    4. Enter the following PySpark script to run your data transformation job on EMR Serverless. Provide the name of your S3 bucket:
      from pyspark.sql import SparkSession
      from pyspark.sql import functions as F
      
      spark = SparkSession.builder.appName("CustomerAnalytics").getOrCreate()
      
      customers = spark.read.parquet("s3://<bucket-name>/processed/customers/")
      transactions = spark.read.parquet("s3://<bucket-name>/processed/transactions/")
      clickstream = spark.read.parquet("s3://<bucket-name>/processed/clickstream/")
      
      customer_spending = transactions.groupBy("customer_id").agg(
          F.count("transaction_id").alias("total_transactions"),
          F.sum("amount").alias("total_spent"),
          F.avg("amount").alias("avg_transaction_value"),
          F.datediff(F.current_date(), F.max("date")).alias("days_since_last_purchase")
      )
      
      customer_engagement = clickstream.groupBy("customer_id").agg(
          F.count("*").alias("total_clicks"),
          F.countDistinct("page").alias("unique_pages_visited"),
          F.count(F.when(F.col("action") == "purchase", 1)).alias("purchase_actions"),
          F.count(F.when(F.col("action") == "add_to_cart", 1)).alias("add_to_cart_actions")
      )
      
      customer_analytics = customers.join(customer_spending, on="customer_id", how="left").join(
          customer_engagement, on="customer_id", how="left")
      
      customer_analytics = customer_analytics.na.fill(0, [
          "total_transactions", "total_spent", "total_clicks", 
          "unique_pages_visited", "purchase_actions", "add_to_cart_actions"
      ])
      
      customer_analytics = customer_analytics.withColumn(
          "customer_value_score",
          (F.col("total_spent") * 0.5) + (F.col("total_transactions") * 0.3) + (F.col("purchase_actions") * 0.2)
      )
      
      customer_analytics.write.mode("overwrite").parquet("s3://<bucket-name>/analytics/customer_analytics/")
      
      customer_summary = customer_analytics.select(
          "customer_id", "name", "email", "registration_date", 
          "total_transactions", "total_spent", "avg_transaction_value",
          "days_since_last_purchase", "total_clicks", "purchase_actions",
          "customer_value_score"
      )
      
      customer_summary.write.mode("overwrite").option("header", "true").csv("s3://<bucket-name>/analytics/customer_summary/")
      
      print("EMR processing completed successfully")

    5. Choose File, Save Notebook As, and save the file as shared/emr_data_transformation_job.ipynb.
    6. Choose Run Cell to run the script.
    7. Monitor the Script execution and verify it completes successfully.
    8. Monitor the Spark job execution and ensure it completes without errors.

    emr run

    Add Redshift Serverless compute

    With Redshift Serverless, users can run and scale data warehouse workloads without managing infrastructure. It is ideal for analytics use cases where data needs to be queried from Amazon S3 or integrated into a centralized warehouse. In this step, you add Redshift Serverless to the project for loading and querying processed customer analytics data generated in earlier stages of the pipeline. For more information about Redshift Serverless, see Amazon Redshift Serverless.

    Set up Redshift Serverless compute in SageMaker Unified Studio

    Complete the following steps to set up Redshift Serverless compute:

    1. In SageMaker Unified Studio, choose the Compute tab within your project workspace (ETL-Pipeline-Demo).
    2. On the SQL analytics tab, choose Add compute, then choose Create new compute resources to begin configuring your compute environment.
    3. Select Amazon Redshift Serverless.
    4. Configure the following:
      1. For Compute name, enter a name (for example, ecommerce_data_warehouse).
      2. For Description, enter a description (for example, Redshift Serverless for data warehouse).
      3. For Workgroup name, enter a name (for example, redshift-serverless-workgroup).
      4. For Maximum capacity, set to 512 RPUs.
      5. For Database name, enter dev.
    5. Choose Add Compute to create the Redshift Serverless resource.
      Redshift

    After the compute is created, you can test the Amazon Redshift connection.

    1. On the Data warehouse tab, confirm that redshift.ecommerce_data_warehouse is listed.
      compute-redshift
    2. Choose the compute: redshift.ecommerce_data_warehouse.
    3. On the Permissions tab, copy the IAM role ARN. You use this for the Redshift COPY command in the next step.
      iam-role

    Create and execute querybook to load data into Amazon Redshift

    In this step, you create a SQL script to load the processed customer summary data from Amazon S3 into a Redshift table. This enables centralized analytics for customer segmentation, lifetime value calculations, and marketing campaigns. Complete the following steps:

    1. On the Build menu, under Data Analysis & Integration, choose Query editor.
    2. Enter the following SQL into the querybook to create the customer_summary table in the public schema:
      -- Create customer_summary table in public schema
      CREATE TABLE IF NOT EXISTS public.customer_summary (
          customer_id INT PRIMARY KEY,
          name VARCHAR(100),
          email VARCHAR(100),
          registration_date DATE,
          total_transactions INT,
          total_spent DECIMAL(10, 2),
          avg_transaction_value DECIMAL(10, 2),
          days_since_last_purchase INT,
          total_clicks INT,
          purchase_actions INT,
          customer_value_score DECIMAL(10, 2)
      );

    3. Choose Add SQL to add a new SQL script.
    4. Enter the following SQL into the querybook
      TRUNCATE TABLE customer_summary;

      Note: We truncate the customer_summary table to remove existing records and ensure a clean, duplicate-free reload of the latest aggregated data from S3 before running the COPY command.

    5. Choose Add SQL to add a new SQL script.
    6. Enter the following SQL to load the data into Redshift Serverless from your S3 bucket. Provide the name of your S3 bucket and IAM role ARN for Amazon Redshift:
      -- Load data from S3 (replace with your bucket name and IAM role)
      COPY public.customer_summary FROM 's3://<bucket-name>/analytics/customer_summary/'
      IAM_ROLE 'arn:aws:iam::<Account-ID>:role/<your-redshift-role>'
      FORMAT AS CSV
      IGNOREHEADER 1
      REGION 'us-east-1';

    7. In the Query Editor, configure the following:
      1. Connection: redshift.ecommerce_data_warehouse
      2. Database: dev
      3. Schema: public

      query

    8. Choose Choose to apply the connection settings.
    9. Choose Run Cell for each cell to create the customer_summary table in the public schema and then load data from Amazon S3.
    10. Choose Actions, Save, name the querybook final_data_product, and choose Save changes.

    This completes the creation and execution of the Redshift data product using the querybook.

    Create and manage the workflow environment

    This section describes how to create a shared workflow environment and define a code-based workflow that automates a customer data pipeline using Apache Airflow within SageMaker Unified Studio. Shared environments facilitate collaboration among project members and centralized workflow management.

    Create the workflow environment

    Workflow environments must be created by project owners. After they’re created, members of the project can sync and use the workflows. Only project owners can update or delete workflow environments. Complete the following steps to create the workflow environment:

    1. Choose Compute for your project.
    2. On the Workflow environments tab, choose Create.
    3. Review the configuration parameters and choose Create workflow environment.
    4. Wait for the environment to be fully provisioned before proceeding It will take around 20 minutes to provision.

    workflow

    Create the code-based workflow

    When the workflow environment is ready, define a code-based ETL pipeline using Airflow. This pipeline automates daily processing tasks across services like AWS Glue, EMR Serverless, and Redshift Serverless.

    1. On the Build menu, under Orchestration, choose Workflows.
    2. Choose Create new workflow, then choose Create workflow in code editor.
    3. Configure Space and choose the instance type ml.t3.xlarge. This ensures your JupyterLab instance has at least 4 vCPUs and 4 GiB of memory.
    4. Choose Configure and Restart Space to launch your environment.

    sample_dag

    The following script defines a daily scheduled ETL workflow that automates several actions:

    • Initial data transformation using AWS Glue
    • Data quality validation using AWS Glue (EvaluateDataQuality)
    • Advanced data processing with EMR Serverless using a Jupyter notebook
    • Loading transformed results into Redshift Serverless from a querybook
    1. Replace the default DAG template with the following definition, ensuring that job names and input paths match the actual names used in your project:
      from datetime import datetime
      from airflow import DAG
      from airflow.decorators import dag
      from airflow.utils.dates import days_ago
      from airflow.providers.amazon.aws.operators.glue import GlueJobOperator
      from workflows.airflow.providers.amazon.aws.operators.sagemaker_workflows import NotebookOperator
      from sagemaker_studio import Project
      # Get SageMaker Studio project IAM role
      project = Project()
      default_args = {
          'owner': 'data_engineer',
          'depends_on_past': False,
          'email_on_failure': True,
          'email_on_retry': False,
          'retries': 1
      }
      @dag(
          dag_id='customer_etl_pipeline',
          default_args=default_args,
          schedule_interval='@daily',
          start_date=days_ago(1),
          is_paused_upon_creation=False,
          tags=['etl', 'customer-analytics'],
          catchup=False
      )
      def customer_etl_pipeline():
          # Step 1: Initial data transformation using Glue
          initial_transformation = GlueJobOperator(
              task_id='initial_transformation',
              job_name='job-1',
              iam_role_arn=project.iam_role,
          )
          # Step 2: Data quality checks using Glue DQ
          data_quality_check = GlueJobOperator(
              task_id='data_quality_check',
              job_name='job-6',
              iam_role_arn=project.iam_role,
          )
          # Step 3: EMR Serverless notebook processing
          emr_processing = NotebookOperator(
              task_id='emr_processing',
              input_config={
                  "input_path": "emr_data_transformation_job.ipynb",
                  "input_params": {}
              },
              output_config={"output_formats": ['NOTEBOOK']},
              poll_interval=10,
          )
          # Step 4: Load to Redshift notebook
          redshift_load = NotebookOperator(
              task_id='redshift_load',
              input_config={
                  "input_path": "final_data_product.sqlnb",
                  "input_params": {}
              },
              output_config={"output_formats": ['NOTEBOOK']},
              poll_interval=10,
          )
          # Task dependencies
          initial_transformation >> data_quality_check >> emr_processing >> redshift_load
      # Instantiate DAG
      customer_etl_dag = customer_etl_pipeline()

    2. Choose File, Save python file, name the file shared/workflows/dags/customer_etl_pipeline.py, and choose Save.

    Deploy and run the workflow

    Complete the following steps to run the workflow:

    1. On the Build menu, choose Workflows.
    2. Choose the workflow customer_etl_pipeline and choose Run.

    scheduled

    Running a workflow puts tasks together to orchestrate Amazon SageMaker Unified Studio artifacts. You can view multiple runs for a workflow by navigating to the Workflows page and choosing the name of a workflow from the workflows list table.

    To share your workflows with other project members in a workflow environment, refer to Share a code workflow with other project members in an Amazon SageMaker Unified Studio workflow environment.

    Monitor and troubleshoot the workflow

    After your Airflow workflows are deployed in SageMaker Unified Studio, monitoring becomes essential for maintaining reliable ETL operations. The integrated Amazon MWAA environment provides comprehensive observability into your data pipelines through the familiar Airflow web interface, enhanced with AWS monitoring capabilities. The Amazon MWAA integration with SageMaker Unified Studio offers real-time DAG execution tracking, detailed task logs, and performance metrics to help you quickly identify and resolve pipeline issues. Complete the following steps to monitor the workflow:

    1. On the Build menu, choose Workflows.
    2. Choose the workflow customer_etl_pipeline.
    3. Choose View runs to see all executions.
    4. Choose a specific run to view detailed task status.

    workflows-run

    For each task, you can view the status (Succeeded, Failed, Running), start and end times, duration, and logs and outputs. The workflow is also visible in the Airflow UI, accessible through the workflow environment, where you can view the DAG graph, monitor task execution in real time, access detailed logs, and view the status.

    1. Go to Workflows and select the workflow named customer_etl_pipeline.
    2. From the Actions menu, choose Open in Airflow UI.

    airflow-ui-smus

    After the workflow completes successfully, you can query the data product in the query editor.

    • On the Build menu, under Data Analysis & Integration, choose Query editor.
    • Run select * from "dev"."public"."customer_summary"

    query-editor

    Observe the contents of the customer_summary table, including aggregated customer metrics such as total transactions, total spent, average transaction value, clicks, and customer value scores. This allows verification that the ETL and data quality pipelines loaded and transformed the data correctly.

    Clean up

    To avoid unnecessary charges, complete the following steps:

    1. Delete a workflow environment.
    2. If you no longer need it, delete the project.
    3. After you delete the project, delete the domain.

    Conclusion

    This post demonstrated how to build an end-to-end ETL pipeline using SageMaker Unified Studio workflows. We explored the complete development lifecycle, from setting up fundamental AWS infrastructure—including Amazon S3 CORS configuration and IAM permissions—to implementing sophisticated data processing workflows. The solution incorporates AWS Glue for initial data transformation and quality checks, EMR Serverless for advanced processing, and Redshift Serverless for data warehousing, all orchestrated through Airflow DAGs. This approach offers several key benefits: a unified interface that consolidates necessary tools, Python-based workflow flexibility, seamless AWS service integration, collaborative development through Git version control, cost-effective scaling through serverless computing, and comprehensive monitoring tools—all working together to create an efficient and maintainable data pipeline solution.

    By using SageMaker Unified Studio workflows, you can accelerate your data pipeline development while maintaining enterprise-grade reliability and scalability. For more information about SageMaker Unified Studio and its capabilities, refer to the Amazon SageMaker Unified Studio documentation.


    About the authors

    Shubham Kumar

    Shubham Kumar

    Shubham is an Associate Delivery Consultant at AWS, specializing in big data, data lakes, data governance, as well as search and observability architectures. In his free time, Shubham enjoys traveling, spending quality time with his family, and writing fictional stories.

    Shubham Purwar

    Shubham Purwar

    Shubham is an Analytics Specialist Solution Architect at AWS. In his free time, Shubham loves to spend time with his family and travel around the world.

    Nitin Kumar

    Nitin Kumar

    Nitin is a Cloud Engineer (ETL) at AWS, specialized in AWS Glue. In his free time, he likes to watch movies and spend time with his family.

    How Convera built fine-grained API authorization with Amazon Verified Permissions

    Post Syndicated from Santhosh Veeraraman original https://aws.amazon.com/blogs/architecture/how-convera-built-fine-grained-api-authorization-with-amazon-verified-permissions/

    Convera processes billions in cross-border payment volume yearly for businesses and financial institutions worldwide. As their platform grew, they needed a robust authorization system that could protect sensitive financial data while maintaining operational efficiency across their global network.

    In this post, we share how Convera used Amazon Verified Permissions to build a fine-grained authorization model for their API platform.

    Background

    As Convera’s service offerings expanded, they needed a scalable, secure, and auditable way to enforce role-based and attribute-based access control. Their goal was to make sure users, both internal and external, had access only to the resources and actions they were explicitly authorized for, while maintaining flexibility to adapt to evolving business needs. Initially, Convera explored building an in-house access control solution. However, they realized that implementing policy management, real-time authorization, logging, and auditing from scratch would require significant engineering effort and ongoing maintenance, diverting resources from their core business priorities. Convera chose Verified Permissions for implementing fine-grained authorization for their payment APIs. This choice was driven by the following factors:

    • Direct integration with AWS services like Amazon Cognito and Amazon API Gateway
    • Cedar policy language’s flexibility in defining complex authorization rules
    • Ability to evaluate multiple attributes like user roles, transaction amounts, and geographic locations
    • High-performance characteristics with millisecond-level authorization decisions

    Given its flexibility and scalability, Verified Permissions became the foundational reference architecture for managing access control across two main scenarios:

    • Fine-grained access control – Convera’s Payment platform serves diverse users including customers, internal staff, and machine-to-machine communications, each requiring specific entitlements based on their roles, organizational hierarchy, and context.
    • Multi-tenancy controls – One of Convera’s most complex requirements was enabling multi-tenant access control while enforcing strict data isolation. Verified Permissions make it possible to define policies that dynamically evaluated tenant ownership, user roles, and contextual attributes.

    The following analysis breaks down Convera’s implementation approach across their major use cases

    Fine-grained access control

    Fine-grained access control is a critical aspect of application security that makes sure users have precisely defined permissions, granting access only to specific resources or actions within an application. Using Verified Permissions, you can define a schema in terms of entity type, including attributes relevant to the authorization model and the valid combinations of principal types, resource types, and actions. Verified Permissions uses this schema to validate that a static policy or policy template is consistent with the application’s authorization model.

    Convera implemented Verified Permissions for fine-grained access control across multiple user types and interaction patterns.

    Customer access management

    At the UI level, Convera used Verified Permissions to manage API access based on specific user characteristics. For example, in their financial applications, the visibility of transaction features like the modify payment parameters is dynamically controlled based on Verified Permissions policies.

    The following diagram illustrates the user authentication decision flow for financial transactions.

    User Authorization Decision for financial transactions

    Figure 1: User Authorization Decision for financial transactions

    The following is an example Cedar policy that can be used in conjunction with Verified Permissions to achieve this use case. This policy makes sure only authorized users with specific roles can see sensitive financial controls. The same policy must be evaluated at the API level when the actual transfer request is made. At the service level, the policy is designed to provide fine-grained access controls.

    permit (
        principal,
        action in [MyApp::Action::"ViewTransferButton"],
        resource
    ) when {
        principal.role == "PAYMENT_INITIATOR" &&
        resource.accountType == "BUSINESS" &&
        resource.status == "ACTIVE"     
    };

    You can integrate the application with Verified Permissions through the API to authorize user access requests. For each authorization request, the service retrieves the relevant policies and evaluates those policies to determine whether a user is permitted to take an action on a resource given context input such as users, roles, group membership, and attributes.

    The following figure illustrates the end-to-end architectural diagram of this implementation.

    Fine-grained User Authorization control with Verified PermissionsPermissions

    Figure 2: Fine-grained User Authorization control with Verified Permissions

    The workflow consists of the following steps:

    1. Users initiate login through the client application.
    2. The client authenticates with Amazon Cognito.
    3. Amazon Cognito triggers a pre-token generation AWS Lambda function to get user roles.
    4. The Lambda function fetches user roles from Amazon Relational Database Service (Amazon RDS).
    5. The Lambda function enriches a JSON Web Token (JWT) with user role information.
    6. The enriched JWT with user roles is returned to the client application.
    7. The client application makes an API call, sending an authorization request to API Gateway through the enriched JWT.
    8. A Lambda authorizer validates the JWT and the role permissions from the JWT and makes a call to Verified Permissions.
    9. Verified Permissions reads access policies stored as Cedar policies and makes an authorization decision.
    10. Verified Permissions returns the authorization result to the Lambda authorizer.
    11. The Lambda authorizer, based on the authorization result, sends an AWS Identity and Access Management (IAM) policy that allows or denies the request to API Gateway.
    12. API Gateway either allows or denies the request to the client application.
    13. API Gateway caches the IAM policy.

    The policy governance is owned by Convera’s infosec team through a strictly regulated IAM role. The changes to the Cedar policies are captured using Amazon DynamoDB Streams and continuously synced with Verified Permissions.

    To improve speed, Convera created a two-level cache system, using the API Gateway built-in cache for authorization decisions and application-level caching for Amazon Cognito tokens. The Lambda function invokes Verified Permissions to authorize the request. If Verified Permissions returns deny, the request is rejected, and an HTTP unauthorized response (403) is sent back. If Verified Permissions returns allow, the request is moved forward. This multi-level caching approach successfully delivers sub-millisecond response times while reducing operational costs and maintaining security controls.

    Internal customer connect applications

    Convera was able to reuse the same architecture for their internal user access as well, such as customer service associates who need quick access to client information to provide efficient service while protecting access to sensitive data. Using Verified Permissions, a role-specific Cedar policy was created to achieve the following:

    • Enable view-only access to basic customer profiles, including contact information and service history
    • Restrict edit capabilities to specific fields, such as updating contact preferences or logging support interactions
    • Block access to sensitive financial data or internal business metrics

    In this flow, internal users authenticate through their enterprise identity provider (IdP), in this case Okta, through the Convera Connect App and obtain ID and access tokens from Amazon Cognito. Amazon Cognito, using a pre-token generation hook, customizes the access token with user attributes stored in Amazon DynamoDB. Although the Cedar policies are tailored for internal roles and responsibilities (different from customer-facing policies), the fundamental flow involving API Gateway, a Lambda authorizer, Verified Permissions policy evaluation, and decision caching remains identical. This architectural reuse meant Convera didn’t need to rebuild their authorization infrastructure, so they can use the same performance optimizations, security controls, and operational processes across both customer and internal user access patterns.

    Extending the model to service communication

    After successfully implementing Verified Permissions for customer and internal user access control, Convera recognized they could use the same architecture for securing service-to-service communications. Similar to how they manage user authentication through Amazon Cognito user pools, each client service is registered in their client configuration system, with Verified Permissions creating a dedicated policy store for service-specific permissions. Services authenticate through Amazon Cognito using client credentials (instead of user credentials) to obtain access tokens that carry service-specific attributes such as service identifier, tier, allowed operations, and rate limits.

    The following diagram illustrates the machine-to-machine architecture for internal and external partner integration.

    Machine to machine architecture for internal and external partner integration

    Figure 3: Machine to machine architecture for internal and external partner integration

    The workflow consists of the following steps:

    1. Service A sends an API request to the authentication token endpoint in API Gateway, including its access token.
    2. The request is forwarded to Amazon Cognito, which validates the token from the machine-to-machine user pool.
    3. Service A, now authenticated, makes a request to the business API (Service B) through API Gateway.
    4. API Gateway forwards the request to the Lambda authorizer, which processes the incoming request and extracts the service context from the token.
    5. The Lambda authorizer sends the authorization request to Verified Permissions to evaluate it against stored Cedar policies. The evaluation considers:
      1. Service identity
      2. Requested operation
      3. Resource context
      4. Environmental factors
    6. Verified Permissions returns an allow or deny decision. If allowed:
      1. The Lambda authorizer generates an appropriate IAM policy.
      2. API Gateway caches the authorization decision.
      3. The request is forwarded to Service B.
      4. Future similar requests can use the cached decision.

    Multi-tenancy controls

    As Convera expanded to support multi-tenant software as a service (SaaS) integrations, they needed a way to implement tenant-specific access controls and data isolation. The challenge was to make sure each tenant’s users could only access their authorized resources while allowing tenant administrators to manage their own access policies. Convera used their existing Verified Permissions architecture with a per-tenant policy store approach to address these requirements.

    Verified Permissions per-tenant policy store

    Convera decided to use per-tenant policy store approach for the following reasons:

    • Low-effort tenant policies isolation
    • The ability to customize templates and schema per tenant
    • Low-effort tenant onboarding and offboarding
    • Per-tenant policy store resource quotas

    The following figure shows the process of implementing fine-grained authorization control using Verified Permissions with a per-tenant policy store.

    Fine-grained authorization control using Verified Permissions with per-tenant policy store

    Figure 4: Fine-grained authorization control using Verified Permissions with per-tenant policy store

    The end-to-end process flow consists of the following steps:

    1. The tenant-specific Amazon Cognito pool is created with a custom attribute called tenant_id. The user logs in to the pool with user claims (for example, user_id).
    2. Amazon Cognito uses a pre-token generation Lambda function that looks up the user_id from a user tenant mapping DynamoDB table.
    3. A DynamoDB table is maintained to map user and tenant configuration.
    4. The pre-token generation Lambda hook gets the tenant_id back from DynamoDB, and adds to the custom tenant_id attribute in the access token.
    5. The user makes an API call to API Gateway with the enriched JWT.
    6. API Gateway validates the token and forwards the request to a custom Lambda authorizer.
    7. The Lambda authorizer function reads the tenant_id from the JWT and looks up the associated Verified Permissions policy-store-id from a DynamoDB table.
    8. The Lambda authorizer verifies the JWT for validity and claims with the Verified Permissions associated Amazon Cognito pool taken from the access token.
    9. If authentication is successful, it calls Verified Permissions to verify that the user is permitted to do the requested action.
    10. If allowed, Verified Permissions returns an IAM policy with Allow access (Deny-by-Default) and forwards the request to backend Kubernetes pods with tenant_id in a custom header.
    11. Backend services receive the tenant_id and validate with Verified Permissions again (for zero-trust policy), creates a tenant context, and forwards to Amazon RDS. Amazon RDS is configured to accept only requests with specific tenant context and returns data specific to the requested tenant_id.

    The following are some examples of Cedar policies used for multi-tenant isolation:

    permit (
        principal in
            convera_connect_authz::userGroup::"ConveraConnect-PAYEE_MGMT",
        action in [convera_connect_authz::Action::"PUT /customer/user/{id}"],
        resource
    );
     
    permit (
        principal,
        action in [convera_connect_authz::Action::"EDIT"],
        resource
    )
    when
    {
        principal.role.contains("UPDATE_USER_STATUS") &&
        resource.type == "PUT" &&
        resource.path == "/customers/user"
    };

    Conclusion

    In this post, we explored how Convera used Verified Permissions to build a sophisticated, fine-grained authorization model for their API platform. We discussed about how Convera was able to implement fine-grained access for their customers, multi-tenant SaaS integrations, machine-to-machine communication scenarios, and internal customer connect applications with the help of Verified Permissions. With Verified Permissions, Convera was able to achieve the following:

    • Implement fine-grained access control across multiple use cases
    • Enhance security with attribute-based access control across multi-tenant environments
    • Improve scalability, handling over thousands of authorization requests per second with submillisecond latency.
    • Increase operational efficiency, reducing time spent on access management tasks by 60%.
    • Future-proof their authorization framework to adapt to evolving business needs

    To learn more about implementing these patterns and best practices, refer to the Verified Permissions User Guide. For hands-on experience, we recommend exploring the Verified Permissions workshop, which provides practical examples and guided exercises.


    About the authors

    Build an AI-powered course recommender using Amazon Bedrock and AWS End User Messaging

    Post Syndicated from Ruchikka Chaudhary original https://aws.amazon.com/blogs/messaging-and-targeting/build-an-ai-powered-course-recommender-using-amazon-bedrock-and-aws-end-user-messaging/

    Educational technology (EdTech) providers face the challenge of maintaining seamless, personalized communication and presenting the right recommendations to their diverse stakeholders. This post explores how combining Amazon Web Services (AWS) End User Messaging and WhatsApp Business API with the advanced AI capabilities of Amazon Bedrock can transform educational engagement.

    In this post, we explore use cases that are reshaping the EdTech industry. We discover how application automation can streamline admissions and enrollment processes, making them more efficient and user-friendly. We demonstrate how instant student engagement can be achieved through AI-powered, personalized interactions that keep learners motivated and connected. We showcase how real-time course feedback mechanisms can help educators adapt and improve their teaching methods. We also examine how student support can be automated using intelligent assistants that provide continuous, all-day assistance while maintaining a personal touch.

    We show you how to build an AI-powered course recommendation system. We explain how to set up WhatsApp Business API integration with Amazon Bedrock, implement smart search capabilities for course matching, and create a scalable serverless architecture. You’ll learn how to build meaningful analytics dashboards to track engagement and learn best practices for handling errors and maintaining system reliability. Whether you’re an EdTech professional or a cloud architect, this guide gives you practical insights into combining conversational AI with educational services.

    Use cases

    • An AI-powered personalized learning pathway generator that automatically recommends customized content based on individual student performance metrics and learning requirements
    • Course improvement suggestions and real-time course feedback
    • A smart communication orchestrator that delivers role-specific, automated notifications and updates across multiple channels to enhance student and parent engagement
    • An early warning system using predictive analytics to identify at-risk students through real-time monitoring of engagement metrics and performance indicators
    • Student support automation with always available AI assistant support, FAQ handling, escalation management, and multilingual support

    Prerequisites

    • An AWS account
    • AWS End User Messaging set up with WhatsApp channel enabled
    • A pre-existing WhatsApp Business account
    • Amazon Bedrock setup must be completed with preferred model
    • Amazon Quick Sight for the AWS Region must be enabled

    Solution overview

    With this solution, users can discover and order educational courses through WhatsApp conversations. Instead of navigating complex websites, the user can send a WhatsApp message saying, “I want to learn Python programming.” They’ll receive personalized course recommendations instantly. The architecture processes WhatsApp messages through AWS End User Messaging, uses Amazon Bedrock for AI-powered conversations, performs semantic search with Amazon OpenSearch Serverless, and captures analytics for business insights. (For step-by-step implementation and rollback guidelines, see the sample course recommendation system.) The following architectural diagram illustrates a modern AI-powered course recommendation system that uses multiple AWS services.

    Figure 1: AI-powered course recommendation system

    Message processing

    When users send WhatsApp messages, AWS End User Messaging captures them and publishes events to an Amazon Simple Notification Service (Amazon SNS) topic. This creates a decoupled architecture where multiple services can process the same message events independently. AWS Lambda functions subscribe to these events, facilitating reliable message processing during high-traffic periods. The decoupled design provides several advantages:

    • If one component fails, others continue operating
    • You can add new message processors without affecting existing ones
    • The system automatically scales based on message volume without manual intervention

    AI conversation engine

    Amazon Bedrock with Claude 3 Haiku powers natural language understanding. It is configured specifically for WhatsApp with instructions for short paragraphs, relevant emoji, and mobile-optimized responses.

    AI agents

    The agent maintains conversation context and handles structured actions such as course search, detail retrieval, and booking through defined functions. The following workflow is the agent action flow and sample code:

    Agent flow

    1. Greets user → Understands intent → Searches courses → Provides details → Facilitates booking
    2. Maintains context throughout the conversation
    3. Can switch between actions based on user responses
    4. Handles complex queries by combining multiple actions

    Sample code

    The following is sample code to create a Bedrock agent using AWS CDK:

        agent = bedrock.CfnAgent(foundation_model="anthropic.claude-3-haiku-20240307-v1:0",
         instruction="""
         Format for WhatsApp: short paragraphs,
            focus on technical courses only
           """,
          action_groups=[# Functions for search, details, booking]
    )

    Semantic search

    Traditional keyword search can miss the user’s intent. The application uses Amazon Titan Embeddings in Amazon Bedrock to convert courses and queries into vectors, enabling semantic understanding. When users ask for “cloud computing courses,” the system can understand related terms such as “AWS” and “serverless” without exact matches. Amazon OpenSearch Serverless handles vector similarity matching combined with traditional filters for course price, level, and duration.

    Analytics pipeline

    Every WhatsApp message interaction generates business intelligence. Messages are stored in Amazon Simple Storage Service (Amazon S3) with date partitioning, catalogued through AWS Glue, and made queryable using Amazon Athena. Teams can analyze user behavior, popular topics, and conversion rates through Quick Sight dashboards. The following dashboard shows example widgets displaying pie-chart breakdown of message delivery status and count of messages per day.

    Figure 2: Amazon Quick Sight dashboard

    As shown in the following dashboard, Amazon Q in QuickSight enables you to explore and analyze your data using conversational AI capabilities.

    Figure 3: Amazon Quick Sight dashboard showing chat window

    Error handling and resilience

    Such highly scalable and distributed solutions require robust error handling. The application has exponential backoff and retries for API calls, meaning the system can gracefully handle rate limits and temporary service unavailability.

    The following is sample code for error handling and resilience:

    python
    def retry_with_backoff(func, max_retries=5):
    retries = 0
    backoff = 1
    while retries < max_retries:
    try:
    return func()
    except ThrottlingException:
    sleep_time = backoff + random.uniform(0, 1)
    time.sleep(sleep_time)
    backoff = min(backoff * 2, 32)
    retries += 1
    raise Exception("Max retries exceeded")

    Business impact

    With the global EdTech market expected to reach $165 billion by 2026, educators and institutions are seeking solutions to prevent student dropouts, improve learning outcomes, and maintain their competitive advantage. Poor personalization can lead to decreased student engagement, lower course completion rates, and ultimately revenue loss.

    Implementing AI-driven personalization and communication systems means institutions can significantly improve student retention rates, boost learning outcomes, and create a more engaging educational experience, which directly impacts their bottom line and reputation in an increasingly competitive educational landscape. This solution could transform educational delivery through intelligent personalization and operational excellence. A serverless architecture can help educational institutions focus on content quality rather than infrastructure management while potentially maintaining rapid response times for course searches. The system’s analytics capabilities could offer insights into student behavior and course preferences, helping shape future curriculum development.

    With mobile optimization, institutions can better serve the growing population of digital-first learners. The combination of automated scaling and pay-per-use pricing could create opportunities for cost optimization, and real-time dashboards can be used to facilitate data-informed decision-making. Such improvements in user experience and operational efficiency could lead to enhanced student engagement and institutional growth in the evolving education environment.

    Sample conversation

    The following video shows how a user can interact with the generative AI-powered course recommendation system and receive course recommendations.

    Future enhancements

    We’re expanding to more messaging platforms, adding voice integration through Amazon Connect, and implementing predictive analytics for personalized recommendations. The serverless architecture makes these additions straightforward without infrastructure changes. Future scenarios could involve:

    • Educator and student support – This solution can be enhanced for student and educator experiences. For educators, it can automate administrative tasks. For students, it can create personalized engagement campaigns, a communication approach that could be significantly more effective than traditional methods.
    • Digital admission process flow – The solution integrates AWS Bedrock AI with WhatsApp Business API to streamline digital admissions. It can enable instant document verification, guide secure payments, and provide automated updates, all within the AWS End User Messaging WhatsApp channel. This AI-powered system could transform the complex admission process into an efficient, chat-based experience, benefiting both institutions and applicants.
    • Parental support and study material management – The system could intelligently distribute learning resources based on student needs, send automated schedule updates, and provide personalized progress reports to parents through WhatsApp. Parents could receive AI-curated study materials and real-time updates about their child’s academic performance, homework assignments, and upcoming assessments through familiar chat interactions. This integration could transform traditional parent-teacher communication into an efficient, automated system while providing timely access to relevant educational resources.

    Conclusion

    The WhatsApp course recommender agent demonstrates how modern AWS services can create sophisticated, AI-powered conversational experiences that scale automatically and provide rich business insights. The serverless architecture provides cost-effectiveness while maintaining enterprise-grade reliability. Key architectural principles that make this solution successful include event-driven design for scalability, AI integration for natural interactions, semantic search for superior user experience, customizable analytics for business intelligence, and infrastructure as code (IaC) for reliable deployments.

    For organizations considering similar implementations, we recommend focusing on user experience optimization, robust error handling, comprehensive monitoring, and gradual feature rollout. The conversational AI environment is rapidly evolving, and solutions that prioritize user experience while maintaining technical excellence can drive the most business value. This implementation can serve as a reference architecture for building production-ready conversational AI systems on AWS, demonstrating patterns that can apply across industries and use cases.


    About the authors

    Use Amazon MSK Connect and Iceberg Kafka Connect to build a real-time data lake

    Post Syndicated from Xiao Huang original https://aws.amazon.com/blogs/big-data/use-amazon-msk-connect-and-iceberg-kafka-connect-to-build-a-real-time-data-lake/

    As analytical workloads increasingly demand real-time insights, organizations need business data to enter the data lake immediately after generation. While various methods exist for real-time CDC data ingestion (such as AWS Glue and Amazon EMR Serverless), Amazon MSK Connect with Iceberg Kafka Connect provides a fully managed, streamlined approach that reduces operational complexity and enables continuous data synchronization.

    In this post, we demonstrate how to use Iceberg Kafka Connect with Amazon Managed Streaming for Apache Kafka (Amazon MSK) Connect to accelerate real-time data ingestion into data lakes, simplifying the synchronization process from transactional databases to Apache Iceberg tables.

    Solution overview

    In this post, we show you how to implement capturing transaction log data from Amazon Relational Database Service (Amazon RDS) for MySQL and writing it to Amazon Simple Storage Service (Amazon S3) in Iceberg table format using append mode, covering both single-table and multi-table synchronization, as shown in the following figure.

    Downstream consumers then process these change records to reconstruct the data state before writing to Iceberg tables.

    In this solution, you use the Iceberg Kafka Sink Connector to implement the business on the sink side. The Iceberg Kafka Sink Connector has the following features:

    • Supports exactly-once delivery
    • Support multi-table synchronization
    • Support schema changes
    • Field name mapping through Iceberg’s column mapping feature

    Prerequisites

    Before beginning the deployment, ensure you have the following components in place:

    Amazon RDS for MySQL: This solution assumes you already have an Amazon RDS for MySQL database instance running with the data you want to synchronize to your Iceberg data lake. Ensure that binary logging is enabled on your RDS instance to support Change Data Capture (CDC) operations.

    Amazon MSK Cluster: You need an Amazon MSK cluster provisioned in your target AWS Region. This cluster will serve as the streaming platform between your MySQL database and the Iceberg data lake. Ensure the cluster is properly configured with appropriate security groups and network access.

    Amazon S3 Bucket: Ensure you have an Amazon S3 bucket ready to host the custom Kafka Connect plugins. This bucket serves as the storage location from which AWS MSK Connect retrieves and installs your plugins. The bucket must exist in your target AWS Region, and you must have appropriate permissions to upload objects to it.

    Custom Kafka Connect Plugins: To enable real-time data synchronization with MSK Connect, you need to create two custom plugins. The first plugin uses the Debezium MySQL Connector to read transactional logs and produce Change Data Capture (CDC) events. The second plugin uses Iceberg Kafka Connect to synchronize data from Amazon MSK to Apache Iceberg tables.

    Build Environment: To build the Iceberg Kafka Connect plugin, you need a build environment with Java and Gradle installed. You can either launch an Amazon EC2 instance (recommended: Amazon Linux 2023 or Ubuntu) or use your local machine if it meets the requirements. Ensure you have sufficient disk space (at least 20GB) and network connectivity to clone the repository and download dependencies.

    Build Iceberg Kafka Connect from open source

    The connector ZIP archive is created as part of the Iceberg build. You can run the build using the following code:

    git clone https://github.com/apache/iceberg.git
    cd iceberg/
    ./gradlew -x test -x integrationTest clean build
    The ZIP archive will be saved in ./kafka-connect/kafka-connect-runtime/build/distributions.

    Create custom plugins

    The next step is to create custom plugins to read and synchronize the data.

    1. Upload the custom plugin ZIP file you compiled in the previous step to your designated Amazon S3 bucket.
    2. Go to the AWS Management Console and navigate to Amazon MSK and choose Connect in the navigation pane.
    3. Choose Custom plugins, then select the plugin file you uploaded to S3 by browsing or entering its S3 URI.
    4. Specify a unique, descriptive name for your custom plugin (such as my-connector-v1).
    5. Choose Create custom plugin.

    Configure MSK Connect

    With the plugins installed, you’re ready to configure MSK Connect.

    Configure data source access

    Start by configuring data source access.

    1. To create a worker configuration, choose Worker configurations in the MSK Connect console.
    2. Choose Create worker configuration and copy and paste the following configuration.
      key.converter.schemas.enable=false
      value.converter.schemas.enable=false
      key.converter=org.apache.kafka.connect.json.JsonConverter
      value.converter=org.apache.kafka.connect.json.JsonConverter
      # Enable topic creation by the worker
      topic.creation.enable=true
      # Default topic creation settings for debezium connector
      topic.creation.default.replication.factor=3
      topic.creation.default.partitions=1
      topic.creation.default.cleanup.policy=delete

    3. In the Amazon MSK console, choose Connectors under Amazon MSK Connect and choose Create connector.
    4. In the setup wizard, select the Debezium MySQL Connector plugin created in the previous step, enter the connector name and select the MSK cluster of the synchronization target. Copy and paste the following content in the configuration:
      
      connector.class=io.debezium.connector.mysql.MySqlConnector
      tasks.max=1
      include.schema.changes=false
      database.server.id=100000
      database.server.name=
      database.port=3306
      database.hostname=
      database.password=
      database.user=
      
      topic.creation.default.partitions=1
      topic.creation.default.replication.factor=3
      
      topic.prefix=mysqlserver
      database.include.list=
      
      ## route
      transforms=Reroute
      transforms.Reroute.type=io.debezium.transforms.ByLogicalTableRouter
      transforms.Reroute.topic.regex=(.*)(.*)
      transforms.Reroute.topic.replacement=$1all_records
      
      # schema.history
      schema.history.internal.kafka.topic
      schema.history.internal.kafka.bootstrap.servers=
      # IAM/SASL
      schema.history.internal.consumer.sasl.mechanism=AWS_MSK_IAM
      schema.history.internal.consumer.sasl.client.callback.handler.class=software.amazon.msk.auth.iam.IAMClientCallbackHandler
      schema.history.internal.consumer.security.protocol=SASL_SSL
      schema.history.internal.consumer.sasl.jaas.config=software.amazon.msk.auth.iam.IAMLoginModule required;
      schema.history.internal.producer.security.protocol=SASL_SSL
      schema.history.internal.producer.sasl.mechanism=AWS_MSK_IAM
      schema.history.internal.producer.sasl.client.callback.handler.class=software.amazon.msk.auth.iam.IAMClientCallbackHandler
      schema.history.internal.producer.sasl.jaas.config=software.amazon.msk.auth.iam.IAMLoginModule required;

      Note that in the configuration, Route is used to write multiple records to the same topic. In the parameter transforms.Reroute.topic.regex, the regular expression is configured to filter the table names that need to be written to the same topic. In the following example, the data containing <tablename-prefix> in the table name is written to the same topic.

      ## route
      transforms=Reroute
      transforms.Reroute.type=io.debezium.transforms.ByLogicalTableRouter
      transforms.Reroute.topic.regex=(.*)(.*)
      transforms.Reroute.topic.replacement=$1all_records

      For example, after transforms.Reroute.topic.replacement is specified as $1all_records, the topic name created in the MSK is < database.server.name>.all_records.

    5. After you choose Create, MSK Connect creates a synchronization task for you.

    Data synchronization (single table mode)

    Now, you can create a real-time synchronization task for the Iceberg table. Start by creating a real-time synchronization job for a single table.

    1. In the Amazon MSK console, choose Connectors under MSK Connect
    2. Choose Create connector.
    3. On the next page, select the previously created Iceberg Kafka Connect plugin
    4. Enter the connector name and select the MSK cluster of the synchronization target.
    5. Paste the following code in the configuration.
      
      connector.class=org.apache.iceberg.connect.IcebergSinkConnector
      tasks.max=1
      topics=
      iceberg.tables=
      iceberg.catalog.catalog-impl=org.apache.iceberg.aws.glue.GlueCatalog
      iceberg.catalog.warehouse=
      iceberg.catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO
      iceberg.catalog.client.region=
      iceberg.tables.auto-create-enabled=true
      iceberg.tables.evolve-schema-enabled=true
      iceberg.control.commit.interval-ms=120000
      transforms=debezium
      transforms.debezium.type=org.apache.iceberg.connect.transforms.DebeziumTransform
      key.converter=org.apache.kafka.connect.json.JsonConverter
      value.converter=org.apache.kafka.connect.json.JsonConverter
      value.converter.schemas.enable=false
      key.converter.schemas.enable=false
      iceberg.control.topic=control-iceberg

      For Iceberg Connector, it will create a topic named control-iceberg by default to record offset. Select the previously created worker configuration that includes topic.creation.enable = true. If you use the default worker configuration and auto-topic creation isn’t enabled at the MSK broker level, the connector will not be able to automatically create topics.

      You can also specify this topic name by setting the parameter iceberg.control.topic = <offset-topic>. If you want to use a custom topic, you can use the following code.

      $KAFKA_HOME/bin/kafka-topics.sh --bootstrap-server $MYBROKERS --create --topic <my-iceberg-offset-topic> --partitions 3 --replication-factor 2 --config cleanup.policy=compact

    6. Query the synchronized data results through Amazon Athena. From the table synchronized to Athena, you can see that, in addition to the source table field, an additional _cdc field has been added to store the metadata content of the CDC.

    Compaction

    Compaction is an essential maintenance operation for Iceberg tables. Although frequent ingestion of small files can negatively impact query performance, regular compaction mitigates this issue by consolidating small files, minimizing metadata overhead, and substantially improving query efficiency. To maintain optimal table performance, you should implement dedicated compaction workflows. AWS Glue offers an excellent solution for this purpose, providing automated compaction capabilities that intelligently merge small files and restructure table layouts for enhanced query performance.

    Schema Evolution Demonstration

    To demonstrate the schema evolution capabilities of this solution, we conducted a test to show how field changes at the source database are automatically synchronized to the Iceberg tables through MSK Connect and Iceberg Kafka Connect.

    Initial Setup:

    First, we created an RDS MySQL database with a customer information table (tb_customer_info) containing the following schema:

    +----------------+--------------+------+-----+-------------------+-----------------------------------------------+
    | Field          | Type         | Null | Key | Default           | Extra                                         |
    +----------------+--------------+------+-----+-------------------+-----------------------------------------------+
    | id             | int unsigned | NO   | PRI | NULL              | auto_increment                                |
    | user_name      | varchar(64)  | YES  |     | NULL              |                                               |
    | country        | varchar(64)  | YES  |     | NULL              |                                               |
    | province       | mediumtext   | NO   |     | NULL              |                                               |
    | city           | int          | NO   |     | NULL              |                                               |
    | street_address | varchar(20)  | NO   |     | NULL              |                                               |
    | street_name    | varchar(20)  | NO   |     | NULL              |                                               |
    | created_at     | timestamp    | NO   |     | CURRENT_TIMESTAMP | DEFAULT_GENERATED                             |
    | updated_at     | timestamp    | YES  |     | CURRENT_TIMESTAMP | DEFAULT_GENERATED on update CURRENT_TIMESTAMP |
    +----------------+--------------+------+-----+-------------------+-----------------------------------------------+

    We then configured MSK Connect using the Debezium MySQL Connector to capture changes from this table and stream them to Amazon MSK in real time. Following that, we set up Iceberg Kafka Connect to consume the data from MSK and write it to Iceberg tables.

    Schema Modification Test:

    To test the schema evolution capability, we added a new field named phone to the source table:

    ALTER TABLE tb_customer_info ADD COLUMN phone VARCHAR(20) NULL;

    We then inserted a new record with the phone field populated:

    INSERT INTO tb_customer_info (user_name,country,province,city,street_address,street_name,phone) values ('user_demo','China','Guangdong',755,'Street1 No.369','Street1','13099990001');

    Results:

    When we queried the Iceberg table in Amazon Athena, we observed that the phone field had been automatically added as the last column, and the new record was successfully synchronized with all field values intact. This demonstrates that Iceberg Kafka Connect’s self-adaptive schema capability seamlessly handles DDL changes at the source, eliminating the need for manual schema updates in the data lake.

    Data synchronization (multi-table mode)

    It’s common that data admins want to use a single connector for moving data in multiple tables. For example, you can use the CDC collection tool to write data from multiple tables to a topic and then write data from one topic to multiple Iceberg tables through the consumer side. In Configure data source access, you configured a MySQL synchronization Connector to synchronize tables with specified rules to a topic using Route. Now let’s review how to distribute data from this topic to multiple Iceberg tables.

    1. When using Iceberg Kafka Connect to synchronize multiple tables to Iceberg tables using AWS Glue Data Catalog, you must pre-create a database in the Data Catalog before starting the synchronization process. The database name in AWS Glue must exactly match the source database name, because the Iceberg Kafka Connect connector automatically uses the source database name as the target database name during multi-table synchronization. This naming consistency is required because the connector doesn’t provide an option to map source database names to different target database names in multi-table scenarios.
    2. If you want to use your custom topic name, you can create a new topic to store the MSK Connect record offset, see Data synchronization (single table mode).
    3. In the Amazon MSK console, create another connector using the following configuration.
      connector.class= org.apache.iceberg.connect.IcebergSinkConnector
      tasks.max=2
      topics=
      iceberg.catalog.catalog-impl=org.apache.iceberg.aws.glue.GlueCatalog
      iceberg.catalog.warehouse=
      iceberg.catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO
      iceberg.catalog.client.region=
      iceberg.tables.auto-create-enabled=true
      iceberg.tables.evolve-schema-enabled=true
      iceberg.control.commit.interval-ms=120000
      transforms=debezium
      transforms.debezium.type=org.apache.iceberg.connect.transforms.DebeziumTransform
      iceberg.tables.route-field=_cdc.source
      iceberg.tables.dynamic-enabled=true
      key.converter=org.apache.kafka.connect.json.JsonConverter
      value.converter=org.apache.kafka.connect.json.JsonConverter
      value.converter.schemas.enable=false
      key.converter.schemas.enable=false
      iceberg.control.topic=control-iceberg

      In this configuration, two parameters have been added:

      • iceberg.tables.route-field: Specifies the routing field that distinguishes between different tables, specified as cdc.source for CDC data parsed by Debezium
      • iceberg.tables.dynamic-enabled: If the iceberg.tables parameter isn’t set, it must be specified as true here
    4. After completion, MSK Connect will creates a sink connector for you.
    5. After the process is complete, you can view the newly created table through Athena.

    Other tips

    In this section, we share some more things that you can use to customize your deployment to fit your use case.

    • Specified table synchronizationIn the Data synchronization (multi-table mode) section, you specify iceberg.tables.route-field = _cdc.Source and iceberg.tables.dynamic-enabled=true, these two parameter settings can write multiple tables stored in the Iceberg table. If you want to synchronize only the specified tables, you can specify the table name you want to synchronize by setting iceberg.tables.dynamic-enabled = false and then setting the iceberg.tables parameter. For example,
      iceberg.tables.dynamic-enabled = false
      iceberg.tables = default.tablename1,default.tablename2
       
      iceberg.table.default.tablename1.route-regex = tablename1
      iceberg.table.default.tablename2.route-regex = tablename2

    • Performance Testing Results
      We conducted a performance test using sysbench to evaluate the data synchronization capabilities of this solution. The test simulated a high-volume write scenario to demonstrate the system’s throughput and scalability.Test Configuration:

      1. Database setup: Created 25 tables in the MySQL database using sysbench
      2. Data loading: Wrote 20 million records to each table (500 million total records)
      3. Real-time streaming: Configured MSK Connect to stream data from MySQL to Amazon MSK in real time during the write process
      4. Kafka Connect configuration:
        • Started Kafka Iceberg Connect
        • Minimum workers: 1
        • Maximum workers: 8
        • Allocated two MCUs per worker

      Performance Results:

      In our test using the configuration above, each MCU achieved peak writing performance of approximately 10,000 records per second, as shown in the following figure. This demonstrates the solution’s ability to handle high-throughput data synchronization workloads effectively.

    Clean up

    To clean up your resources, complete the following steps:

    1. Delete MSK Connect connectors: Remove both the Debezium MySQL Connector and Iceberg Kafka Connect connector created for this solution.
    2. Delete the Amazon MSK cluster: If you created a new MSK cluster specifically for this demonstration, delete it to stop incurring charges.
    3. Delete the S3 buckets: Remove the S3 buckets used to store the custom Kafka Connect plugins and Iceberg table data. Ensure you have backed up any data you need before deletion.
    4. Delete the EC2 instance: If you launched an EC2 instance to build the Iceberg Kafka Connect plugin, terminate it.
    5. Delete the RDS MySQL instance (optional): If you created a new RDS instance specifically for this demonstration, delete it. If you’re using an existing production database, skip this step.
    6. Remove IAM roles and policies (if created): Delete any IAM roles and policies that were created specifically for this solution to maintain security best practices.

    Conclusion

    In this post, we presented a solution to achieve real-time, efficient data synchronization from transactional databases to data lakes using Amazon MSK Connect and Iceberg Kafka Connect. This solution provides a low-cost and efficient data synchronization paradigm for enterprise-level big data analysis. Whether you’re working with ecommerce transactions, financial transactions, or IoT device logs, this solution can help you achieve quick access to a data lake, enabling analytical businesses to quickly obtain the latest business data. We encourage you to try this solution in your own environment and share your experiences in the comments section. For more information, visit Amazon MSK Connect.


    About the author

    Huang Xiao

    Huang Xiao

    Huang is a Senior Specialist Solution Architect with Analytics at AWS. He focuses on big data solution architecture design, with years of experience in development and architectural design within the big data field.

    Optimizing Flink’s join operations on Amazon EMR with Alluxio

    Post Syndicated from Qingyuan Tang original https://aws.amazon.com/blogs/big-data/optimizing-flinks-join-operations-on-amazon-emr-with-alluxio/

    When you’re working with data analysis, you often face the challenge of effectively correlating real-time data with historical data to gain actionable insights. This becomes particularly critical when you’re dealing with scenarios like e-commerce order processing, where your real-time decisions can significantly impact business outcomes. The complexity arises when you need to combine streaming data with static reference information to create a comprehensive analytical framework that supports both your immediate operational needs and strategic planning

    To tackle this challenge, you can employ stream processing technologies that handle continuous data flows while seamlessly integrating live data streams with static dimension tables. These solutions enable you to perform detailed analysis and aggregation of data, giving you a comprehensive view that combines the immediacy of real-time data with the depth of historical context. Apache Flink has emerged as a leading stream computing platform that offers robust capabilities for joining real-time and offline data sources through its extensive connector ecosystem and SQL API.

    In this post, we show you how to implement real-time data correlation using Apache Flink to join streaming order data with historical customer and product information, enabling you to make informed decisions based on comprehensive, up-to-date analytics.

    We also introduce an optimized solution to automatically load Hive dimension table data into Alluxio Universal Flash Storage (UFS) through the Alluxio cache layer. This enables Flink to perform temporal joins on changing data, accurately reflecting the content of a table at specific points in time.

    Solution architecture

    When it comes to joining Flink SQL tables with stream tables, the lookup join is a go-to method. This approach is particularly effective when you need to correlate streaming data with static or slowly changing data. In Flink, you can use connectors like the Flink Hive SQL connector or the FileSystem connector to archive the scenario.

    The following architecture shows general approach which we describe ahead:

    Here’s how we do this:

    1. We use offline data to construct a Flink table. This data could be from an offline Hive database table or from files stored in a system like Amazon S3. Concurrently, we can create a stream table from the data flowing in through a Kafka message stream
    2. Use a batch cluster for offline data processing. In this example, we use an Amazon EMR cluster which creates a fact table in it. It also provides a Detail Wide Data (DWD) table which has been used as a Flink dynamic table to perform consequence processing after a lookup join
      • It is typically located in the middle layer of a data warehouse, between the raw data contained in the Operational Data Store (ODS) and the highly aggregated data found in the Data Warehouse (DW), or Data Mart (DM).
      • The primary purpose of the DWD layer is to support complex data analysis and reporting needs by providing a detailed and comprehensive data view.
      • Both the fact table and DWD table are hive tables on Hadoop
    3. Use a streaming cluster for the real-time processing. In this example, we use an Amazon EMR cluster to stream event ingestion and analyze it using Flink, using Flink Kafka connector and Hive connector to join the streaming event data and statics dimension data (fact table)

    One of the key challenges encountered with this approach is related to the management of the lookup dimension table data. Initially, when the Flink application is started, this data is stored in the task manager’s state. However, during subsequent operations like continuous queries or window aggregations, the dimension table data isn’t automatically refreshed. This means that the operator must either restart the Flink application periodically or manually refresh the dimension table data in the temporary table. This step is crucial to ensure that the join operations and aggregations are always performed with the most current dimension data.

    Another significant challenge with this approach is needing to pull the entire dimension table data and perform a cold start each time. This becomes particularly problematic when dealing with a large volume of dimension table data. For instance, when handling tables with tens of millions of registered users or tens of thousands of product SKU attributes, this process generates substantial input/output (IO) overhead. Consequently, it leads to performance bottlenecks, impacting the efficiency of the system.

    Flink’s checkpointing mechanism processes the data and stores checkpoint snapshots of all the states during continuous queries or window aggregations, resulting in state snapshots data bloat.

    Optimizing the solution

    This post includes an optimized solution to address the aforementioned challenges, by automatically loading Hive dimension table data into the Alluxio UFS via the Alluxio cache layer. We join this data with Flink’s temporal joins to create a view on a changing table. This view reflects the content of a table at a specific point in time

    Alluxio is a distributed cache engine for big data technology stacks. It provides a unified UFS that can connect to the underlying Amazon S3 and HDFS data. Alluxio UFS read and write operations warm up the distributed storage layers on S3 and HDFS and thus significantly increase throughput and reducing network overhead. Deeply integrated with upper level computing engines such as Hive, Spark, and Trino, Alluxio is an excellent cache accelerator for offline dimension data.

    Additionally, we utilize Flink’s temporal table function to pass a time parameter. This function returns a view of the temporal table at the specified time. By doing so, when the main table of the real-time dynamic table is correlated with the temporal table, it can be associated with a specific historical version of the dimension data

    Solution implementation details

    For this post, we use “user behavior” log data in Kafka as real-time stream fact table data, and user information data on Hive as offline dimension table data. A demo with Alluxio + Flink temporal join is used to verify the Flink join optimized solution.

    Real-time fact tables

    For this demonstration, we utilize user behavior JSON data simulated by the open-source component json-data-generator. We write the data to Amazon Managed Kafka (Amazon MSK) in real-time. Using the Flink Kafka Connector, we convert this stream into a Flink stream table for continuous queries. This served as our fact table data for real-time joins.

    A sample of the user behavior simulation data in JSON format is as follows:

    [{          
    	"timestamp": "nowTimestamp()",
    	"system": "BADGE",
    	"actor": "Agnew",
    	"action": "EXIT",
    	"objects": ["Building 1"],
    	"location": "45.5,44.3",
    	"message": "Exited Building 1"
    }]
    

    It includes user behavior information such as operation time, login system, user signature, behavioral activities, and service objects, locations, and related text fields. We create a fact table in Flink SQL with the main fields as follows:

    CREATE TABLE logevent_source (`timestamp`  string, 
    `system` string,
     actor STRING,
     action STRING
    ) WITH (
    'connector' = 'kafka',
    'topic' = 'logevent',
    'properties.bootstrap.servers' = 'b-6.msk06.dr04w4.c3.kafka.ap-southeast-1.amazonaws.com:9092,b-5.msk06.dr04w4.c3.kafka.ap-southeast-1.amazonaws.com:9092,b-1.msk06.dr04w4.c3.kafka.ap-southeast-1.amazonaws.com:9092 (http://b-6.msk06.dr04w4.c3.kafka.ap-southeast-1.amazonaws.com:9092%2Cb-5.msk06.dr04w4.c3.kafka.ap-southeast-1.amazonaws.com:9092%2Cb-1.msk06.dr04w4.c3.kafka.ap-southeast-1.amazonaws.com:9092/)',
    'properties.group.id' = 'testGroup6',
    'scan.startup.mode'='latest-offset',
    'format' = 'json'
    );

    Caching dimension tables with Alluxio

    Amazon EMR provides solid integration with Alluxio. You can use the Amazon EMR bootstrap startup script to automatically deploy Alluxio components and start the Alluxio master and worker processes when an Amazon EMR cluster is created. For detailed installation and deployment steps, refer to the article Integrating Alluxio on Amazon EMR.

    In an Amazon EMR cluster that integrates Alluxio, you may use Alluxio to create a cache table for the Hive offline dimension table as follows:

    ##Set up the client jar package in hive-env.sh:
    $ export HIVE_AUX_JARS_PATH=/<PATH_TO_ALLUXIO>/client/alluxio-2.2.0-client.jar:${HIVE_AU
    
    ##Make sure the UFS is configured on the EMR cluster where Alluxio is installed and that the table/db path has been created:
    alluxio fs mkdir alluxio://ip-xxx-xxx-xxx-xxx.ap-southeast-1.compute.internal:19998/s3/customer
    alluxio fs chown hadoop:hadoop alluxio://ip-xxx-xxx-xxx-xxx.ap-southeast-1.compute.internal:19998/s3/customer
    
    ##On the AWS EMR cluster, create a Hive table path pointing to Alluxio namespace URI:
    !connect jdbc:hive2://xxx.xxx.xxx.xxx:10000/default;
    hive> CREATE TABLE customer(
        c_customer_sk             bigint,
        c_customer_id             string,
        c_current_cdemo_sk        bigint,
        c_current_hdemo_sk        bigint,
        c_current_addr_sk         bigint,
        c_first_shipto_date_sk    bigint,
        c_first_sales_date_sk     bigint,
        c_salutation              string,
        c_first_name              string,
        c_last_name               string,
        c_preferred_cust_flag     string,
        c_birth_day               int,
        c_birth_month             int,
        c_birth_year              int,
        c_birth_country           string,
        c_login                   string,
        c_email_address           string
    )
        ROW FORMAT DELIMITED
        FIELDS TERMINATED BY '|'
        STORED AS TEXTFILE
        LOCATION 'alluxio://ip-xxx-xxx-xxx-xxx.ap-southeast-1.compute.internal:19998/s3/customer';
    OK
    Time taken: 3.485 seconds

    As shown in the previous section, the Alluxio table location alluxio://ip-xxx-xx:19998/s3/customer points to the S3 path where the Hive dimension table is located; writing to the customer dimension table is automatically synchronized to the Alluxio cache.

    After creating the Alluxio Hive offline dimension table, you can view the details of the Alluxio cache table by connecting to the Hive metadata through the Hive catalog in Flink SQL:

    CREATE CATALOG hiveCatalog WITH (  'type' = 'hive',
        'default-database' = 'default',
        'hive-conf-dir' = '/etc/hive/conf/',
        'hive-version' = '3.1.2',
        'hadoop-conf-dir'='/etc/hadoop/conf/'
    );
    -- set the HiveCatalog as the current catalog of the session
    USE CATALOG hiveCatalog;
    show create table customer;
    create external table customer(
        c_customer_sk             bigint,
        c_customer_id             string,
        c_current_cdemo_sk        bigint,
        c_current_hdemo_sk        bigint,
        c_current_addr_sk         bigint,
        c_first_shipto_date_sk    bigint,
        c_first_sales_date_sk     bigint,
        c_salutation              string,
        c_first_name              string,
        c_last_name               string,
        c_preferred_cust_flag     string,
        c_birth_day               int,
        c_birth_month             int,
        c_birth_year              int,
        c_birth_country           string,
        c_login                   string,
        c_email_address           string
    ) 
    row format delimited fields terminated by '|'
    location 'alluxio://ip-xxx-xxx-xxx-xxx.ap-southeast-1.compute.internal:19998/s3/30/customer' 
    TBLPROPERTIES (
      'streaming-source.enable' = 'false',  
      'lookup.join.cache.ttl' = '12 h'
    )

    As shown in the preceding code, the location path of the dimension table is the UFS cache path Uniform Resource Identifier (URI). When the business program reads and writes the dimension table, Alluxio automatically updates the customer dimension table data in the cache and asynchronously writes it to the Alluxio backend storage path of the S3 table to achieve table data synchronization in the data lake.

    Flink temporal table join

    Flink temporal table is also a type of dynamic table. Each record in the temporal table is correlated with one or more time fields. When we join the fact table and the dimension table, we usually need to obtain real-time dimension table data for the lookup join. Thus, when creating or joining a table, we usually need to use the proctime() function to specify the time field of the fact table. When we join the tables, we use the syntax of FOR SYSTEM_TIME AS OF to specify the time version of the fact table that corresponds to the time of the lookup dimension table.

    For this post, the customer information is a changing dimension table in the Hive offline table, whereas the customer behavior is the fact table in Kafka. We specified the time field with proctime() in the Flink Kafka source table. Then when joining the Flink Hive table, we used FOR SYSTEM_TIME AS OF to specify the time field of the lookup Kafka source table to allow us to realize the Flink temporal table join operation

    As shown in the following code, a fact table of user behavior is created through the Kafka Connector in Flink SQL. The ts field refers to the timestamp when the temporal table is joined:

    CREATE TABLE logevent_source (`timestamp`  string, 
    `system` string,
     actor STRING,
     action STRING,
     ts as PROCTIME()
    ) WITH (
    'connector' = 'kafka',
    'topic' = 'logevent',
    'properties.bootstrap.servers' = 'b-6.msk06.dr04w4.c3.kafka.ap-southeast-1.amazonaws.com:9092,b-5.msk06.dr04w4.c3.kafka.ap-southeast-1.amazonaws.com:9092,b-1.msk06.dr04w4.c3.kafka.ap-southeast-1.amazonaws.com:9092 (http://b-6.msk06.dr04w4.c3.kafka.ap-southeast-1.amazonaws.com:9092%2Cb-5.msk06.dr04w4.c3.kafka.ap-southeast-1.amazonaws.com:9092%2Cb-1.msk06.dr04w4.c3.kafka.ap-southeast-1.amazonaws.com:9092/)',
    'properties.group.id' = 'testGroup-01',
    'scan.startup.mode'='latest-offset',
    'format' = 'json'
    );

    The Flink offline dimension table and the streaming real-time table are joined as follows:

    select a.`timestamp`,a.`system`,a.actor,a.action,b.c_login from 
           (select *, proctime() as proctime from user_logevent_source) as a 
     left join customer  FOR SYSTEM_TIME AS OF a.proctime as b on a.actor=b.c_last_name;

    When the fact table logevent_source joins the lookup dimension table, the proctime function ensures real-time joins by obtaining the latest dimension table version. This dimension data, cached in Alluxio, delivers significantly better read performance than direct S3 access.

    At the same time, the dimension table data is already cached in Alluxio; the read performance is much better than offline data read on S3.

    The comparison test shows that Alluxio cache brings a clear performance advantage by switching the S3 and Alluxio paths of the customer dimension table through Hive

    You can easily switch the local and cache location paths with alter table in hive cli:

    alter table customer set location "s3://xxxxxx/data/s3/30/customer";
    alter table customer  set location "alluxio://ip-xxx-xxx-xxx-xxx.ap-southeast-1.compute.internal:19998/s3/30/customer";

    You can also select the Task Manager log from the Flink dashboard for a split test.

    The performance of the fact table load was doubled through the implementation of optimized data processing techniques.

    1. Before caching (S3 path read): 5s load time
      2022-06-29 02:54:34,791 INFO  com.amazon.ws.emr.hadoop.fs.s3n.S3NativeFileSystem           [] - Opening 's3://salunchbucket/data/s3/30/customer/data-m-00029' for reading
      2022-06-29 02:54:39,971 INFO  org.apache.flink.table.filesystem.FileSystemLookupFunction   [] - Loaded 433000 row(s) into lookup join cache

    2. After caching (Alluxio read): 2s load time
      2022-06-29 03:25:14,476 INFO  com.amazon.ws.emr.hadoop.fs.s3n.S3NativeFileSystem           [] - Opening 's3://salunchbucket/data/s3/30/customer/data-m-00029' for reading
      2022-06-29 03:25:16,397 INFO  org.apache.flink.table.filesystem.FileSystemLookupFunction   [] - Loaded 433000 row(s) into lookup join cache

    The timeline on JobManager clearly shows the difference in execution duration under Alluxio and S3 paths.

    For single task query ,we accelerate by more than 1 times using this solution. The overall job performance improvement is even more visible.

    Other optimalizations to consider

    Implementing a continuous join requires pulling dimension data every time. Does it lead to Flink’s checkpoint state bloat that can cause Flink TaskManager RocksDB to explode or memory overflow.
    In Flink, the state comes with a TTL mechanism. You can set a TTL expiration policy to trigger Flink to clean up expired state data. Flink SQL can be set using the hint method.

    insert into logevent_sink
    select a.`timestamp`,a.`system`,a.actor,a.action,b.c_login from 
    (select *, proctime() as proctime from logevent_source) as a 
      left join 
    customer/*+ OPTIONS('lookup.join.cache.ttl' = '5 min')*/  FOR SYSTEM_TIME AS OF a.proctime as b 
    on a.actor=b.c_last_name;

    Flink Table/Streaming API is similar:

    StateTtlConfig ttlConfig = StateTtlConfig
        .newBuilder(Time.days(7))
        .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
        .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
        .cleanupInRocksdbCompactFilter() 
        .build();
    ValueStateDescriptor<Long> lastUserLogin = 
        new ValueStateDescriptor<>("lastUserLogin", Long.class);
    lastUserLogin.enableTimeToLive(ttlConfig);
    StreamTableEnvironment.getConfig().setIdleStateRetentionTime(min, max);

    Restart the lookup join after the configuration. As you can see from the Flink TM log, after TTL expires, it triggers clean-up and re-pull the Hive dimension table data:

    2022-06-29 04:17:09,161 INFO  org.apache.flink.table.filesystem.FileSystemLookupFunction   
    [] - Lookup join cache has expired after 5 minute(s), reloading

    In addition, you can reduce the number of checkpoint snapshots by configuring Flink state retention and thereby reduce the amount of space taken up by state at the time of snapshot.

    Flink job configuration as follow:
    -D state.checkpoints.num-retained=5 

    After the configuration, you can see that in the S3 checkpoint path, the Flink job automatically cleans up historical snapshots and keeps the most recent 5 snapshots, thus ensuring that checkpoint snapshots do not accumulate.

    [hadoop@ip-172-31-41-131 ~]$ aws s3 ls s3://salunchbucket/data/checkpoints/7b9f2f9becbf3c879cd1e5f38c6239f8/
                               PRE chk-3/
                               PRE chk-4/
                               PRE chk-5/
                               PRE chk-6/
                               PRE chk-7/

    Summary

    Customers implementing Flink streaming framework to join dimension and real-time fact tables frequently encounter performance challenges. In this post, we presented an optimized solution that uses Alluxio’s caching capabilities to automatically load Hive dimension table data into the UFS cache. By integrating with Flink temporal table joins, dimension tables are transformed into time-versioned views, effectively addressing performance bottlenecks in traditional implementations.


    About the author

    Jeff Tang

    Jeff Tang

    Jeff is a Data Analytics Solutions Architect at AWS. He’s responsible for designing and optimizing Amazon Data Analytic services, with over 10 years of experience in data architecture and development. Former roles include Senior Consulting Advisor at Oracle, Senior Architect at Migu Culture Data Market, and Data Analytics Architect at ANZ Bank. Extensive experience in big data, data lakes, intelligent lakehouses, and MLOps platforms

    Federate access to Amazon SageMaker Unified Studio with AWS IAM Identity Center and Ping Identity

    Post Syndicated from Raghavarao Sodabathina original https://aws.amazon.com/blogs/big-data/federate-access-to-amazon-sagemaker-unified-studio-with-aws-iam-identity-center-and-ping-identity/

    With an identity provider (IdP), you can manage your user identities outside of AWS and give these external user identities permissions to use AWS resources in your AWS accounts. External IdPs, such as Ping Identity, can integrate with AWS IAM Identity Center to be the source of truth for Amazon SageMaker Unified Studio. SageMaker Unified Studio also supports trusted identity propagation for SQL analytics, including Amazon Athena and Amazon Redshift.

    SageMaker Unified Studio provides an integrated experience to use your data and tools for analytics and AI. You can use SageMaker Unified Studio to discover your data and put it to work using familiar AWS analytics and machine learning (ML) services for model development, generative AI, big data processing, and SQL analytics, assisted by Amazon Q Developer. By default, SageMaker domains support AWS Identity and Access Management (IAM) user credentials. You can also enable access to SageMaker domains in SageMaker Unified Studio for users with single sign-on (SSO) with IAM Identity Center and direct SAML integration with SageMaker Unified Studio.

    Users can access SageMaker Unified Studio with their existing corporate credentials. With IAM Identity Center, administrators can connect their existing external IdPs and continue to manage users and groups in those existing identity systems, which can then be synchronized with IAM Identity Center using System for Cross-domain Identity Management (SCIM).In this post, we show how to set up workforce access with SageMaker Unified Studio using Ping Identity as an external IdP with IAM Identity Center.

    In this post, we show how to set up workforce access with SageMaker Unified Studio using Ping Identity as an external IdP with IAM Identity Center.

    Solution overview

    We walk through the following high-level steps to implement this solution:

    1. Enable IAM Identity Center.
    2. Create a SageMaker Unified Studio domain.
    3. Set up your IdP (for this example, Ping Identity).
    4. Connect Ping Identity and IAM Identity Center.
    5. Set up automatic provisioning of users and groups in IAM Identity Center.
    6. Configure SageMaker Unified Studio SSO user access.

    Prerequisites

    For this walkthrough, you should have the following prerequisites:

    • An AWS account with IAM Identity Center enabled. It is recommended to use an organization-level IAM Identity Center instance for best practices and centralized identity management across your AWS organization.
    • A Ping Identity account.
    • A browser with network connectivity to Ping Identity and SageMaker Unified Studio.

    Enable IAM Identity Center

    To enable IAM Identity Center, follow the instructions in Enable IAM Identity Center.

    Create a SageMaker Unified Studio domain

    To create a SageMaker Unified Studio domain, refer to the instructions in Create a Amazon SageMaker Unified Studio domain – manual setup.

    On the SageMaker console, go to the domain details and copy the Amazon Resource Name (ARN) under Domain ARN. You will use this value when you add your trust policy and when you connect your IAM IdP to your Ping Identity instance.

    Create a SageMaker Unified Studio domain

    Set up your IdP (Ping Identity)

    In this section, we walk through the procedure to set up your IdP (for this example, Ping Identity).

    Create an environment in Ping Identity

    Complete the following steps to create an environment for Ping Identity:

    1. Log in to your Ping Identity account.
    2. Choose Create Environment.
    3. Choose Create a Customer Solution.
    4. In the Tailor your experiences pop-up, choose Skip.
      Create an environment in Ping Identity

    Create a group in Ping Identity

    Complete the following steps to create a group in Ping Identity:

    1. On the Environments page, choose Manage Environments.
    2. In the navigation pane, choose Directory, then choose Groups.
    3. Choose the plus sign to add a group.
    4. For Group Name, enter sagemaker
    5. For Description, enter an optional description (for example, Amazon SageMaker Unified Studio).
    6. For Population, choose Default.
    7. Choose Save.
      Create a group in Ping Identity
    8. On the Roles tab for the sagemaker group, assign the Environment Admin role to the group.
      Assigning roles for the sagemaker group

    Create a user in Ping Identity

    Complete the following steps to create a user:

    1. In the navigation pane, choose Directory, then choose Users.
    2. Choose the plus sign to create a user.
    3. Provide values for Given name, Family name, Username, and Email.
    4. For Password, choose First time password.
    5. Choose Save.

    You can add more users as needed.

    Assign group to user

    Complete the following steps to assign your group to your user:

    1. In the navigation pane, choose Directory, then choose Groups.
    2. Choose the sagemaker group you created.
    3. On the Users tab, choose the plus sign to add a user.
    4. Add the user you created.

    Connect Ping Identity and IAM Identity Center

    To configure the integration between Ping Identity and IAM Identity Center, you need access to both management consoles. Although Ping Identity’s application catalog includes IAM Identity Center, we recommend configuring a standard SAML application for greater control over settings and attribute mappings.

    Complete the following steps:

    1. Go to the Ping Identity environment you created and choose Applications in the navigation pane.
    2. Choose the plus sign to add an application:
      1. For Application name, enter a name (for this example, we use unifiedstudio).
      2. For Description, enter an optional description.
      3. For Application Type, choose SAML Application.
      4. Choose Configure.

      Creating a SAML app integration in Ping Identity

    3. Sign in to the IAM Identity Center console as a user with administrative privileges.
    4. In the navigation pane, choose Settings to update your settings:
      1. On the Identity source tab, choose Change identity source on the Actions dropdown menu.
        Selecting identity source in AWS IAM Identity Center
      2. For Choose identity source, select External identity provider, then choose Next.

        Choosing External Identity provider in AWS IAM Identity Center

      3. In the Service provider metadata section, choose Download metadata file to download the IAM Identity Center metadata file.

        You will use this service provider metadata file in the next step when you connect Ping Identity with IAM Identity Center.

      Downloading service provider metadata from AWS IAM Identity Center

    5. Return to the Ping Identity console and the SAML application page.
    6. In the SAML Configuration section, select Import Metadata, upload the metadata file you downloaded, then choose Save.

      Importing service provider metadata into Ping Identity

    7. On the Overview tab of the application page, choose Download Metadata under Connection details to download the Ping Identity IdP metadata.
      You will use this for the SAML configuration in IAM Identity Center to set up Ping Identity as an IdP in the next step.

      Downloading Identity provider metadata from Ping Identity

    8. Return to the IAM Identity Center console and continue configuring your identity source:
      1. In the Identity provider metadata section, choose Choose file under IdP SAML metadata, upload the metadata file you downloaded from Ping Identity, then choose Next.

        Configuring Ping Identity as Identity Provider in AWS IAM Identity Center

      2. Choose Accept to accept the disclaimer.
      3. Choose Change identity source.
    9. Return to the Ping Identity console to complete the SAML configuration.
    10. On the Configuration tab, choose the edit icon to update the configuration:
      1. For Sign, choose Sign Assertion & Response.
      2. For Subject Name ID, enter urn:oasis:names:tc:SAML:1.1:nameid-format:emailAddress.
      3. For Assertion Validity Duration, enter 300.
      4. Leave the remaining values as default.

      Ping Identity SAML Configurations

    11. On the Attributes tab, choose the edit icon.
    12. Choose +Add to add two attribute mappings:
      1. Map the attribute saml-subject to Username, and leave Name format as default.
      2. Map the attribute https://aws.amazon.com/SAML/Attributes/PrincipalTag:Email to Email Address, and set Name format to Unspecified.
      3. Choose Save.

      Ping Identity SAML attributes mapping

    13. On the PingOne Policies tab, select Single Factor, then choose Save.
      This post uses single-factor authentication for demonstration purposes only. In your environments, follow your organization’s security standards and governance framework.

      Ping Identity policy configuration

    14. On the Access tab, search for the sagemaker group under Group Membership Policy, and assign the unifiedstudio SAML application to the group.
    15. Enable the application.
      Enabling Ping Identity SMAL application

    Set up automatic provisioning of users and groups from Ping Identity into IAM Identity Center

    To configure the automatic provisioning of users and groups between Ping Identity and IAM Identity Center through SCIM, you must have access to both management consoles. Complete the following steps:

    1. On the IAM Identity Center console, choose Settings in the navigation pane.
    2. In the Automatic provisioning section, choose Enable.
      Enabling automatic provisioning in AWS IAM Identity Center

      This enables automatic provisioning in IAM Identity Center and displays the necessary SCIM endpoint and access token information.

    3. In the Inbound automatic provisioning dialog box, copy the values for SCIM endpoint and Access token, then choose Close.
      You will use these values to configure provisioning in Ping Identity in the next step.

      Automatic provisioning configuration parameters in IAM Identity Center

      This completes the setup process in IAM Identity Center.

    4. Log in to the Ping Identity console.
    5. In the navigation pane, choose Integrations, then choose Provisioning.
    6. Choose the plus sign to add a new connection.
      Creating a new SCIM connection
    7. For Choose a connection type, choose Select next to Identity Store.
      Choosing connection type
    8. Provide a name (for this example, we use Identitycenter) and an optional description, then choose Next.
      Creating new connection
    9. Under Configuration Authentication, provide the following configuration:
      1. For SCIM BASE URL, enter the SCIM endpoint from IAM Identity Center.
      2. For Authentication Method, choose OAuth 2 Bearer Token.
      3. For Oauth Access Token, enter the access token from IAM Identity Center.
      4. For Auth Type Header, choose Bearer (default option).
      5. Choose Test Connection to validate the connection between Ping Identity and IAM Identity Center, then choose Next.

      Configuring authentication between Ping Identity and IAM Identity Center

    10. Under Configuration Preference, provide the following configuration:
      1. For User Filter Expression, enter userName Eq “%s”.
      2. For Group Membership Handling, select Merge.
      3. Leave the remaining settings as default and choose Save.

      SCIM connection preferences

    11. On the Provisioning tab, choose the plus sign, then choose New Rule to create a rule for the SCIM connection.
      Creating a new SCIM rule
    12. Enter a name (for this example, unifiedstudio) and an optional description, then choose Create Rule.
    13. Under the newly created rule, choose the plus sign next to Available Connections to add the connection identitycenter, then choose Save.
    14. Edit the user filter:
      1. For Attribute, choose Enabled.
      2. For Operator, choose Equals.
      3. For Value, choose true.
      4. Choose Save.

      User Filter attributes mapping

    15. Choose the edit icon next to Attribute Mapping and set the attribute mappings as shown in the following screenshot:
      1. Delete the Primary Phone attribute mapping because it’s optional in AWS. Leaving this field blank can cause Ping Identity’s SCIM connector to generate errors during user provisioning.
      2. Add a new attribute called Username under PingOne Directory and then map to displayName under Identitycenter.

      Attributes mapping between Ping Identity SCIM and AWS IAM Identity Center

    16. Under Group Provisioning, choose the sagemaker group if you want to sync all sagemaker group users with auto provisioning.
      1. In the pop-up, select I understand and want to continue, then choose Save.

      Assigning groups to SCIM rule

      Assigning groups to SCIM rule

    17. On the Provisioning page, choose the Connections tab.
    18. Enable the SCIM connection Identitycenter and rule unifiedstudio.

      Enabling the SCIM connection

      Enabling the SCIM rule

    This completes the SCIM setup process between Ping Identity and IAM Identity Center.

    Configure SageMaker Unified Studio SSO user access

    Complete the following steps to configure SSO user access to SageMaker Unified Studio for your SageMaker domain:

    1. On the SageMaker console, choose Domains in the navigation pane.
    2. Choose the domain for which you want to configure SAML user access.
    3. On the domain details page, you can find the SSO configuration in two locations:
      1. From the main domain view, choose Configure next to Configure SSO user access.
      2. Alternatively, scroll down to the User management tab and choose Configure SSO user access.

      SageMaker Unified Studio SSO configuration

    4. On the Choose user authentication method page, select IAM Identity Center, then choose Next.
      Choosing authentication
    5. For Choose user and group assignment method, choose from the following options, then choose Next:
      1. Require assignments: Users and groups must be explicitly added to the domain to gain access. This provides more granular control over who can access the domain.
      2. Do not require assignments: All authorized Ping Identity users and groups can access this domain if they have been assigned to the SAML application in Ping Identity.

      For either option, users or groups must have access to the Ping Identity SAML application (unifiedstudio in this example) to authenticate successfully.

      SageMaker Unified Studio SAML configuration

    6. On the Review and save page, review your choices and choose Save. These settings can’t be changed after you save them.
      Review and confirm SAML configuration
    7. If you’ve chosen to require assignments, use the Add users and groups section to add SAML users and groups to your domain.
      Add users and groups to SageMaker Unified Studio domain

    Now, users will be able to access SageMaker Unified Studio using the domain URL with their SSO credentials.

    You can explore different projects for your users and assign those projects based on your IdP user groups for fine-grained access controls. For example, you can create different SAML user groups based on their job function in Ping Identity, then assign those Ping Identity groups to the unifiedstudio SAML application in Ping Identity, and then assign those Ping Identity SAML groups to their respective project profiles in SageMaker Unified Studio. To assign project profiles for their respective groups, choose the Project profiles tab and choose your project profile. On the Authorized users and groups page, choose Add, then choose SSO groups. Choose Add users and groups button to complete the project profile assignment.

    Assigning a project profile to Ping Identity group

    Validate access with Ping Identity users

    Complete the following steps to validate access:

    1. On the SageMaker domain details page, choose the link for the SageMaker Unified Studio URL.
      Validating Ping Identity user access with Amazon SageMaker Unified Studio
    2. Log in with your user credentials.
      After successful login, you will be redirected to the SageMaker Unified Studio home page. Here, you can explore different projects to your users and assign those projects based on your SAML user groups for fine-grained access control.

      SAML authenticated Amazon SageMaker Unified Studio

    3. To assign an authorization policy, those Govern and then Domain units.
    4. Choose your SageMaker domain, then choose a suitable authorization policy. For this example, we choose Project creation policy.
      Amazon SageMaker unified studio authorization policies
    5. Choose Add policy grant to assign user groups or users to their respective project profiles.
      Amazon SageMaker unified studio authorization policies assignment

    You have successfully federated SageMaker Unified Studio with Ping Identity as an IdP with IAM Identity Center. You can connect to SageMaker Unified Studio by using your Ping Identity credentials.

    Clean up

    After you test out this solution, remember to delete the resources you created to avoid incurring future charges. For instructions to delete your SageMaker Unified Studio domain, refer to Delete domains. If you want to delete your Ping Identity account, reach out to Ping Identity for assistance.

    Conclusion

    In this post, we demonstrated how to set up Ping Identity as an IdP over SAML authentication for SageMaker Unified Studio access through IAM Identity Center federation. To learn more, refer to the Amazon SageMaker Unified Studio User Guide, which provides guidance on how to build data and AI applications using SageMaker.


    About the authors

    Raghavarao Sodabathina

    Raghavarao Sodabathina

    Raghavarao is a Principal Solutions Architect at AWS, focusing on data analytics, AI/ML, and cloud security. He engages with customers to create innovative solutions that address customer business problems and accelerate the adoption of AWS services. In his spare time, Raghavarao enjoys spending time with his family, reading books, and watching movies.

    Matt Nispel

    Matt Nispel

    Matt is an Enterprise Solutions Architect at AWS. He has more than 10 years of experience building cloud architectures for large enterprise companies. At AWS, Matt helps customers rearchitect their applications to take full advantage of the cloud. Matt lives in Minneapolis, Minnesota, and in his free time enjoys spending time with friends and family.

    Himanshu Sarda

    Himanshu Sarda

    Himanshu is a Solutions Architect at AWS who specializes in generative AI and autonomous agent architectures, helping enterprise customers revolutionize their businesses through cutting-edge AI solutions. When not pioneering AI innovations, Himanshu recharges by exploring the outdoors and creating memories with family and friends.

    Nicholaus Lawson

    Nicholaus Lawson

    Nicholaus is a Solutions Architect at AWS and part of the AI/ML specialty group. He has a background in software engineering and AI research. Outside of work, Nicholaus is often coding, learning something new, or woodworking.

    Krupanidhi Jay

    Krupanidhi Jay

    Krupanidhi is a Boston-based Enterprise Solutions Architect at AWS. He is a seasoned architect with over 20 years of experience in helping customers with digital transformation and delivering seamless digital user experiences. He enjoys working with customers to help them build scalable, cost-effective solutions in AWS. Outside of work, Jay enjoys spending time with family and traveling.

    Build a trusted foundation for data and AI using Alation and Amazon SageMaker Unified Studio

    Post Syndicated from Anthony Lempelius, James Mesney original https://aws.amazon.com/blogs/big-data/build-a-trusted-foundation-for-data-and-ai-using-alation-and-amazon-sagemaker-unified-studio/

    This post was co-written with Anthony Lempelius and James Mesney from Alation.

    When a team wants to reuse a dataset, whether it is to build a new pipeline, launch a dashboard, run an analysis, or power an AI application, the first challenge is rarely the code. Data engineers need to understand lineage, transformations, and operational expectations. Data analysts and BI engineers need consistent definitions, metrics, and trusted sources. Data scientists and AI engineers need to know provenance, quality, access constraints, and how data or features were derived. In many organizations, that context is captured in different places by different teams, often across solutions like Alation and SageMaker Unified Studio, both of which can serve as a system of record for business context depending on who is doing the work and where they operate day to day. When those perspectives are not connected, people revalidate the same information, debate definitions, and duplicate documentation across tools. A unified metadata foundation brings these role specific views together so business context, technical metadata, and governance stay aligned across platforms, making data easier to trust, easier to find, and easier to use across analytics and AI.

    The new Alation integration with Amazon SageMaker Unified Studio addresses these challenges by synchronizing catalog metadata between both systems. This synchronization creates a unified metadata experience where technical teams working in SageMaker Unified Studio and business teams working in Alation collaborate on top of the same metadata. You can verify how ML and analytics assets are created, understand dependencies, and maintain traceability across your data lifecycle regardless of which system your teams prefer to use.

    In this post, we demonstrate who benefits from this integration, how it works, the specific metadata it synchronizes, and provide a complete deployment guide for your environment.

    The value of unified metadata governance

    Organizations managing large-scale analytics and ML workloads face critical challenges when metadata is fragmented across multiple systems. When metadata exists in silos, data scientists spend valuable time searching for the right datasets. Teams duplicate metadata management efforts, creating inconsistent definitions and conflicting metrics across the organization.

    Regulatory requirements demand clear provenance. Without unified metadata governance, organizations struggle to demonstrate compliance, trace data origins, and maintain audit trails across their ML and analytics pipelines. Data discovery becomes a bottleneck when teams can’t quickly find, understand, and trust the data they need, delaying model development and reducing the overall business value of data investments.

    Applying consistent governance policies across disparate systems is nearly impossible without a unified metadata layer. This creates security vulnerabilities, data quality issues, and compliance blind spots. A unified metadata governance approach alleviates these challenges by providing a single source of truth for metadata across ML and analytics systems, enabling faster data discovery, consistent governance, and confident compliance while reducing the operational burden on data and ML teams.

    Solution overview

    The Alation and SageMaker Unified Studio integration unifies the user experience, synchronizing metadata from cataloged assets between both systems.

    This Phase 1 integration extracts metadata from Amazon SageMaker Catalog into Alation, giving you one place to discover assets.

    The integration connects through AWS Identity and Access Management (IAM) authentication and synchronizes key metadata elements, including domains, projects, asset names, descriptions, owners, glossary terms, and custom metadata fields. Every metadata update includes provenance information: the originating service, the person who made the change, and the timestamp, creating comprehensive audit trails for compliance.

    You can run metadata extractions on demand or schedule them to run automatically. The system performs an initial bulk extraction of your selected domains and projects, then keeps it up-to-date through incremental updates using either event-driven triggers or scheduled polling. Communication uses encrypted APIs with scoped IAM permissions following least-privilege principles.

    This integration helps organizations in financial services, telecommunications, retail, manufacturing, and transportation that manage large numbers of analytics and ML workloads across many systems and teams. You can reduce metadata duplication, accelerate data discovery, and enable your data scientists, analysts, and engineers to find trusted data faster so they can focus on building insights rather than validating data quality.

    The following diagram illustrates the solution architecture.

    The following screenshot showcases the Alation catalog displaying the SageMaker Unified Studio project and its synchronized assets.

    Metadata synchronization

    This integration automatically synchronizes essential metadata between SageMaker Unified Studio and Alation, facilitating consistent information across both systems. The synchronization brings together the types of metadata you need for discovery, governance, and audit workflows, giving you clearer insight into how datasets, features, and models relate across your services.

    The integration synchronizes catalog metadata, including domains, projects, asset names, descriptions, owners, glossary terms, and metadata forms. Additionally, the integration synchronizes provenance metadata, which includes information about the originating service, the actor who made the change, and the timestamp, to support traceability and audit workflows.

    Integration mechanics

    The integration connects SageMaker Unified Studio and Alation through a scoped IAM role that provides secure, encrypted communication. After you configure this connection within Alation, the system performs an initial extraction of your selected domains and projects, then keeps information current through incremental updates using either event-driven triggers or scheduled polling.

    The integration synchronizes metadata forms from SageMaker Unified Studio into Alation through automated field mapping between both systems’ schemas. Metadata forms can capture various asset specific details like feature store references, training run identifiers, model versions, and evaluation metrics.

    Every metadata update includes provenance information: the originating service, the person who made the change, and when it occurred. This supports audit and stewardship workflows. Access controls follow least-privilege principles through IAM while applying Alation’s role-based permissions, letting you limit synchronization by project, namespace, or tag as needed.

    Security and compliance

    Security and compliance are critical when synchronizing metadata across systems. This integration follows enterprise security practices to facilitate safe, controlled metadata synchronization. The connector uses least-privilege access, encrypted transport, and clear separation between metadata and data, so you can maintain governance without disrupting existing workflows.

    You configure a scoped IAM role to define which accounts, projects, and namespaces the connector can access, making sure access follows your organization’s security policies. Metadata moves over TLS-protected APIs, and you control which domains and projects to include in Alation. By default, the integration synchronizes only metadata; your data files and artifacts remain in their original AWS locations unless you explicitly choose to export them.

    Alation maintains a complete audit trail by recording extraction events, mapping changes, and stewardship activities. These security controls support compliant metadata governance while preserving your existing operational practices.

    Prerequisites

    Before setting up this integration, ensure you have the following:

    • An Alation Cloud Service (ACS) instance
    • Alation server admin access
    • An AWS account
    • A SageMaker Unified Studio domain and project with existing metadata

    Configure authentication

    Before configuring the Alation connector, you must set up the required AWS resources and permissions. The first step is to configure authentication. The Alation connector supports two authentication methods to access SageMaker Unified Studio. Choose the method that best fits your security requirements.

    Option 1: IAM role (Recommended)

    Create an IAM role that the Alation connector will assume to access SageMaker Unified Studio. For detailed instructions on creating IAM roles, see IAM role creation.

    The following is an example IAM permission policy for SageMaker Catalog access:

    {
       "Version": "2012-10-17",
        "Statement": [
            {
                "Sid": "AlationSageMakerAccess",
                "Effect": "Allow",
                "Action": [
                    "datazone:ListDomains",
                    "datazone:GetFormType",
                    "datazone:Search",
                    "datazone:ListProjects",
                    "datazone:GetAsset"
                ],
                "Resource": "arn:aws:datazone:<region>:<account-id>:domain/*”
            }
        ]
    }

    The following is an example trust policy for the IAM role:

    {
        "Version": "2012-10-17",
        "Statement": [
            {
                "Sid": "AlationSageMakerAccessAssumeRole",
                "Effect": "Allow",
                "Principal": {
                    "AWS": "<alation_provided_role_arn>"
                },
                "Action": "sts:AssumeRole"
            }
        ]
    }     

    Option 2: IAM user with access keys

    Create an IAM user with programmatic access and attach the necessary permissions. For detailed instructions on creating IAM users, see Create an IAM user in your AWS account.

    Create an IAM user with programmatic access enabled, attach the following policy, and generate access keys for use in Alation configuration:

    {
       "Version": "2012-10-17",
        "Statement": [
            {
                "Sid": "AlationSageMakerAccess",
                "Effect": "Allow",
                "Action": [
                    "datazone:ListDomains",
                    "datazone:GetFormType",
                    "datazone:Search",
                    "datazone:ListProjects",
                    "datazone:GetAsset"
                ],
                "Resource": "arn:aws:datazone:<region>:<account-id>:domain/*"
            }
        ]
    }

    Add IAM role or user to SageMaker Unified Studio domain

    Add the IAM role or user you created to the SageMaker Unified Studio domain. For detailed instructions on adding users to a domain, see User management in Amazon SageMaker Unified Studio. The following screenshot shows an example of adding IAM users on the SageMaker dashboard.

    Add IAM role or user to SageMaker Unified Studio projects

    The IAM role or user must be added as a member to all SageMaker Unified Studio projects that contain metadata you want to synchronize with Alation. Projects without this member will not be included in the synchronization process.

    Add the IAM role or user as a project member with Contributor or Owner permissions for each project you want to include in the sync, as illustrated in the following screenshot. For detailed instructions on adding project members, see Add project members.

    Install SageMaker enhanced connector

    After completing the AWS setup, you can configure the Alation connector to establish the integration. The connector is distributed as a .zip package for upload and installation in the Alation application. To obtain the connector, contact the Forward Deployed Engineering team or your Alation Account Manager.

    When you have the .zip package, follow the installation procedures to add the connector.

    Create and configure Alation’s data source

    Navigate to the Data Sources section in Alation, create a new data source, and select SageMaker Catalog as the source type. Configure the connection settings with the authentication method chosen in the AWS setup.

    For IAM role authentication, use the following configuration:

    • Connection Type: IAM Role
    • Role ARN: ARN of the IAM role created in AWS setup
    • External ID: External ID configured in the trust policy
    • AWS Region: Region where your SageMaker Unified Studio domain is located

    For IAM user authentication, use the following configuration:

    • Connection Type: Access Keys
    • Access Key ID: Access key from AWS setup
    • Secret Access Key: Secret key from AWS setup
    • AWS Region: Region where your SageMaker Unified Studio domain is located

    Test the connection to verify authentication and network connectivity, as shown in the following screenshot.

    Configure metadata extraction settings

    Configure the extraction scope by selecting the SageMaker domains and projects to synchronize, as shown in the following screenshot. Only projects where the IAM role or user is a member will be available for synchronization.

    Run initial extraction

    Execute the first metadata synchronization to import existing metadata from SageMaker Unified Studio into Alation. Monitor the extraction progress through Alation’s status indicators and validate that SageMaker assets appear correctly in the catalog.

    The following screenshot shows the job history page with job status Running.

    The following screenshot shows the job history page with job status Succeeded.

    The following screenshot shows the Alation catalog displaying the SageMaker Unified Studio project and its synchronized assets.

    Operate and tune

    Configure ongoing operations by setting extraction cadence, configuring reconciliation alerts, and monitoring logs regularly. Add data stewards to synchronized assets, and consider enabling AI-generated descriptions or working with Alation Professional Services for advanced governance design.

    Enhanced capabilities

    The next phase of the integration introduces three key capabilities: bi-directional metadata synchronization, lineage replication, and data quality metadata replication. The bi-directional capability gives you the flexibility to control where metadata updates originate, either in Alation or in SageMaker Unified Studio, so you can manage metadata changes in the service that best aligns with your organizational workflows and governance processes.

    The feature set is rolling out in phases. Phase 1 is available at the time of writing this post and provides extraction from SageMaker Unified Studio into Alation, including initial and incremental updates and audit logging. Phase 2 is coming soon and will offer configurable principal catalogs, advanced scoped syncs, and reconciliation workflows for Alation Cloud Service customers.

    These enhancements will support governed, scalable ML operations with increasing depth and automation.

    Conclusion

    The Alation and SageMaker Unified Studio integration helps organizations bridge the gap between fast analytics and ML development and the governance requirements most enterprises face. By cataloging metadata from SageMaker Unified Studio in Alation, you gain a governed, discoverable view of how assets are created and used. This supports leaders, stewards, compliance teams, and ML practitioners who depend on accurate, well-documented data to scale analytics and AI responsibly.

    To learn more about this integration and explore additional resources, refer to the Amazon SageMaker Unified Studio User Guide and Alation Documentation.


    About the authors

    Anthony Lempelius

    Anthony Lempelius

    Anthony is the Director of Channel and Alliances at Alation, where he leads strategic partnerships with independent software vendor (ISV) and systems integrator (SI) partners. He focuses on bringing joint integrations and solutions to market that help customers unlock value from trusted, well-governed data. Anthony is passionate about building the AWS Partner Network that accelerates innovation across the data and AI landscape.

    James Mesney

    James Mesney

    James is a Principal Product Manager at Alation, where he leads product strategy for advancing Alation’s Agentic capabilities. He focuses on helping organizations make their data more discoverable, governed, and actionable by shaping features that improve metadata quality, user experience, and AI-driven insights. James is passionate about building products that empower enterprises to fully unlock the value of trusted data.

    Divij Bhatia

    Divij Bhatia

    Divij is a Software Development Engineer at AWS. He is passionate about building resilient and scalable cloud-based solutions that solve real-world problems for customers. His free time often takes him outdoors, traveling and shooting landscapes.

    Leonardo Gomez

    Leonardo Gomez

    Leonardo is a Principal Analytics Specialist Solutions Architect at AWS. He has over a decade of experience in data management, helping customers around the globe address their business and technical needs.

    Reduce EMR HBase upgrade downtime with the EMR read-replica prewarm feature

    Post Syndicated from Suthan Phillips original https://aws.amazon.com/blogs/big-data/reduce-emr-hbase-upgrade-downtime-with-the-emr-read-replica-prewarm-feature/

    HBase clusters on Amazon Simple Storage Service (Amazon S3) need regular upgrades for new features, security patches, and performance improvements. In this post, we introduce the EMR read-replica prewarm feature in Amazon EMR and show you how to use it to minimize HBase upgrade downtime from hours to minutes using blue-green deployments. This approach works well for single-cluster deployments where minimizing service interruption during infrastructure changes is important.

    Understanding HBase operational challenges

    HBase cluster upgrades have required complete cluster shutdowns, resulting in extended downtime while regions initialize and RegionServers come online. Version upgrades require a complete cluster switchover, with time-consuming steps that include loading and verifying region metadata, performing HFile checks, and confirming proper region assignment across RegionServers. During this critical period—which can extend to hours depending on cluster size and data volume—your applications are completely unavailable.

    The challenge doesn’t stop at version upgrades. You must regularly apply security patches and kernel updates to maintain compliance. For Amazon EMR 7.0 and later clusters running on Amazon Linux 2023, instances don’t automatically install security updates after launch; they remain at the patch level from cluster creation time. AWS recommends periodically recreating clusters with newer AMIs, requiring the same hard cutover and downtime risks as a full version upgrade. Similarly, when you need to use different instance types, traditional approaches mean taking your cluster offline.

    Solution overview

    Amazon EMR 7.12 introduces read-replica prewarm, a new feature that tackles these challenges. This feature lets you make infrastructure changes to Apache HBase on Amazon S3 at scale while reducing downtime risk and maintaining data consistency.

    With read-replica prewarm, you can prepare and validate your changes in a read-replica cluster before promoting it to active status, cutting service interruption from hours to minutes. You will learn how to prepare your read-replica cluster with the target version, execute cutover procedures that minimize downtime, and verify successful migration before completing the switchover.

    Read-replica prewarm architecture

    The following diagram shows the architecture and workflow. Both primary and read-replica clusters interact with the same Amazon S3 storage, accessing the same S3 bucket and root directory.

    Amazon EMR HBase architecture diagram showing primary cluster in Availability Zone 1 with read/write access to Amazon S3, and read-replica cluster in Availability Zone 2 with read access to S3.

    Distributed locking confirms only one HBase cluster can write at a time (for clusters version 7.12.0 and later). The read-replica cluster performs full HBase region initialization without time pressure, and after promotion, the read replica becomes the active writer as shown in the following diagram.

    Amazon EMR HBase failover scenario showing primary cluster unavailable in Availability Zone 1, with read-replica cluster in Availability Zone 2 promoted to handle read and write operations after failover.

    Implementation steps HBase cluster upgrade

    Now that you understand how read-replica prewarm works and the architecture behind it, let’s put this knowledge into practice. You will follow a process that consists of three main phases: preparation, cutover, and verification. Each phase includes specific steps, shown in the following figure, that you will execute in sequence to complete the migration.

    Process flow diagram showing three-phase HBase cluster migration: Phase 1 preparation and validation, Phase 2 cutover and DNS update, Phase 3 post-migration verification.

    Phase 1: Preparation

    Before starting the migration, prepare both your primary cluster and launch a new read-replica cluster. Each step in this phase builds toward confirming that your new cluster can properly access and serve your existing data.

    1. Run major compactions on tables to verify regions are not in SPLIT state
      Run major compactions to consolidate data files and verify regions are not in SPLIT state. Split regions can cause assignment conflicts during migration, so resolving them at the start helps maintain cluster stability throughout the transition.

      echo “major_compact 'tablename'” | hbase shell

    2. Run catalog_janitor to clean up stale regions
      Execute the catalog_janitor process (HBase’s built-in maintenance tool) to remove stale region references from the metadata. Cleaning up these references prevents confusion during region assignment in the read-replica cluster.

      echo “catalogjanitor_run” | hbase shell

    3. Confirm no inconsistencies in the primary HBase cluster
      Verify cluster integrity before migration:

      sudo -u hbase hbase hbck > hbck_report.txt

      Running the HBase Consistency Check tool version 2 (HBCK2) performs a diagnostic scan that identifies and reports problems in metadata, regions, and table states, confirming your cluster is ready for migration.

    4. Launch HBase read-replica cluster with the target version connecting to the same HBase root directory in Amazon S3 as the primary cluster
      Launch a new HBase cluster with the target version and configure it to connect to the same S3 root directory as the primary cluster. Confirm that read-only mode is enabled by default as shown in the following screenshot.

      AWS console screenshot showing Amazon EMR data durability and availability configuration options, with "Create a read-replica cluster" option selected and S3 location settings.

      If you are using AWS Command Line Interface (AWS CLI), you can enable the read replica while launching the Amazon EMR HBase on the Amazon S3 cluster by setting the hbase.emr.readreplica.enabled.v2 parameter to true in the HBase classification as shown in the following example:

      {
          "Classification": "hbase",
          "Properties": {
            "hbase.emr.readreplica.enabled.v2": "true",
            "hbase.emr.storageMode": "s3"
          }
      }

    5. Run meta refresh in this read-replica HBase cluster
      echo "refresh_meta" | hbase shell

      You’re creating a parallel environment with the new version that can access existing data without modification risk, allowing validation before committing to the upgrade.

    6. Validate the read-replica and verify that regions show OPEN status and are properly assigned:
      Execute sample read operations against your key tables to confirm the read replica can access your data correctly. In the HBase Master UI, verify that regions show OPEN status and are properly assigned to RegionServers. You should also confirm that the total data size matches your previous cluster to verify complete data visibility.
    7. Prepare for cutover on primary cluster
      Disable balancing and compactions on the primary cluster:

      echo "balance_switch false" | hbase shell
      echo "compaction_switch false" | hbase shell

      Preventing background operations from changing data layout or triggering region movements maintains a consistent state during the migration window.

      Take snapshots of your tables for rollback capability:

      # For each table
      echo "snapshot 'table_name', 'table_name_pre_migration_$(date +%Y%m%d)'" | hbase shell
      # For system tables
      echo "snapshot 'hbase:meta', 'meta_pre_migration_$(date +%Y%m%d)'" | hbase shell
      echo "snapshot 'hbase:namespace', 'namespace_pre_migration_$(date +%Y%m%d)'" | hbase shell

      These snapshots enable point-in-time recovery if you discover issues after migration.

    8. Run meta refresh and refresh hfiles on the read replica:
      echo "refresh_meta" | hbase shell
      hbase org.apache.hadoop.hbase.client.example.RefreshHFilesClient "table_name'"

      Refreshing confirms the read replica has the most current region assignments, table structure, and HFile references before taking over production traffic.

    9. Check for inconsistencies in the read-replica cluster
      Run the HBCK2 tool on the read-replica cluster to identify potential issues:

      sudo -u hbase hbase hbck > hbck_report.txt

      When a read replica is created, both the primary and replica clusters show metadata inconsistencies referencing each other’s meta folders: “There is a hole in the region chain”. The primary cluster complains about meta_<read-replica-cluster-id>, while the read replica complains about the primary’s meta folder. This inconsistency doesn’t impact cluster operations but shows up in hbck reports. For a clean hbck report after switching to the read replica and terminating the primary cluster, manually delete the old primary’s meta folder from Amazon S3 after taking a backup of it.

      Additionally, check the HBase Master UI to visually confirm cluster health. Verifying the read-replica cluster has a clean, consistent state before promotion prevents potential data access issues after cutover.

    Phase 2: Cutover

    Perform the actual migration by shutting down the primary cluster and promoting the read replica. The steps in this phase minimize the window when your cluster is unavailable to applications.

    1. Remove the primary cluster from DNS routing
      Update DNS entries to direct traffic away from the primary cluster, preventing new requests from reaching it during shutdown.
    2. Flush in-memory data to Amazon S3
      Flush in-memory data to confirm durability in Amazon S3:

      # Flush application data  
      echo "flush 'usertable'" | hbase shell
      # Flush system tables
      echo "flush 'hbase:meta'" | hbase shell
      echo "flush 'hbase:namespace'" | hbase shell

      Flushing forces data still in memory (in MemStores, HBase’s write cache) to be written to persistent storage (Amazon S3), preventing data loss during the transition between clusters.

    3. Terminate the primary cluster
      Terminate the primary cluster after confirming the data is persisted to Amazon S3. This step releases resources and eliminates the possibility of split-brain scenarios where both clusters might accept writes to the same dataset.
    4. Promote the read replica to active status
      Convert the read replica to read-write mode:

      echo "readonly_switch false" | hbase shell  
      echo "readonly_state" | hbase shell  # Verify the switch was successful

      The promotion process automatically refreshes meta and HFiles, capturing final changes from the flush operations and confirming complete data visibility.

      When you promote the cluster, it transitions from read-only to read-write mode, allowing it to accept application write operations and fully replace the old cluster’s functionality.

    5. Update DNS to point to the new active cluster
      Update DNS entries to direct traffic to the new active cluster. Routing client traffic to the new cluster restores service availability and completes the migration from the application perspective.

    Phase 3: Validation

    With your new cluster now active, you’re ready to verify that everything is working correctly before declaring the migration complete.

    Execute test write operations to confirm the cluster accepts writes properly. Check the HBase Master UI to verify regions are serving both read and write requests without errors. At this point, your migration to the new Amazon EMR release is complete, and your applications can connect to the new cluster and resume normal read-write operations.

    Key benefits

    The read-replica prewarm approach delivers several important advantages over traditional HBase upgrade methods. Most notably, you can reduce service interruption from hours to minutes by preparing your new cluster in parallel with your running production environment.

    Before committing to the upgrade, you can thoroughly test that data is readable and accessible in the new version. The system loads and assigns regions before activation, eliminating the lengthy startup time that traditionally causes extended downtime. This pre-warming process means your new cluster is ready to serve traffic immediately upon promotion.

    You also gain the ability to validate multiple aspects of your deployment before cutover, including data integrity, read performance, cluster stability, and configuration correctness. This validation happens while your production cluster continues serving traffic, reducing the risk of discovering issues during your maintenance window.

    For testing and validation workflows, you can run parallel testing environment by creating multiple HBase read replicas. However, you should verify that only one HBase cluster remains in read-write mode to the Amazon S3 data store to prevent data corruption and consistency issues.

    Rollback procedures

    Always thoroughly test your HBase rollback procedures before implementing upgrades in production environments.

    When rolling back HBase clusters in Amazon EMR, you have two primary options.

    • Option 1 involves launching a new cluster with the previous HBase version that points to the same Amazon S3 data location as the upgraded cluster. This approach is straightforward to implement, preserves data written before and after the upgrade attempt, and offers faster recovery with no additional storage requirements. However, it risks encountering data compatibility issues if the upgrade modified data formats or metadata structures, potentially leading to unexpected behavior.
    • Option 2 takes a more cautious approach by launching a new cluster with the previous HBase version and restoring from snapshots taken before the upgrade. This method guarantees a return to a known, consistent state, eliminates version compatibility risks, and provides complete isolation from corruption introduced during the upgrade process. The tradeoff is that data written after the snapshot was taken will be lost, and the restoration process requires more time and planning.

    For production environments where data integrity is paramount, the snapshot-based approach (option 2) is generally preferred despite the potential for some data loss.

    Considerations

    • Store file tracking migration: Migrating from Amazon EMR 7.3 (or earlier) requires disabling and dropping the hbase:storefile table on the primary cluster, then flushing metadata. When launching the new read-replica cluster, configure the DefaultStoreFileTracker implementation using the hbase.store.file-tracker.impl property. When operational, run change_sft commands to switch tables to FILE tracking method, providing seamless data file access during migration.
    • Multi-AZ deployments: Consider network latency and Amazon S3 access patterns when deploying read replicas across Availability Zones. Cross-AZ data transfer might impact read latency for the read-replica cluster.
    • Cost impact: Running parallel clusters during migration incurs additional infrastructure costs until the primary cluster is terminated.
    • Disabled tables: The disabled state of tables in the primary cluster is a cluster-specific administrative property that isn’t propagated to the read-replica cluster. If you want them disabled in the read replica, you must explicitly disable them.
    • Amazon EMR 5.x cluster upgrade: Direct upgrade from Amazon EMR 5.x to Amazon EMR 7.x using this feature isn’t supported because of the major HBase version change from 1.x to 2.x. For upgrading from Amazon EMR 5.x to Amazon EMR 7.x, follow the steps in our best practices: AWS EMR Best Practices – HBase Migration

    Conclusion

    In this post, we showed you how the read-replica prewarm feature of Amazon EMR 7.12 improves HBase cluster operations by minimizing the hard cutover constraints that make infrastructure changes challenging. This feature gives you a consistent blue-green deployment pattern that reduces risk and downtime for version upgrades and security patches.

    When you can thoroughly validate changes before committing to them and reduce service interruption from hours to minutes, you can maintain HBase infrastructure more confidently and efficiently. You can now take a more proactive approach to cluster maintenance, security compliance, and performance optimization with greater confidence in your operational processes.

    To learn more about Amazon EMR and HBase on Amazon S3, visit the Amazon EMR documentation. To get started with read replicas, see the HBase on Amazon S3 guide .


    About the authors

    Suthan Phillips

    Suthan Phillips

    Suthan is a Senior Analytics Architect at AWS, where he helps customers design and optimize scalable, high-performance data solutions that drive business insights. He combines architectural guidance on system design and scalability with best practices to provide efficient, secure implementation across data processing and experience layers. Outside of work, Suthan enjoys swimming, hiking, and exploring the Pacific Northwest.

    Ramesh Kandasamy

    Ramesh Kandasamy

    Ramesh is an Engineering Manager at Amazon EMR. He is a long tenured Amazonian dedicated to solve distributed systems problems.

    Mehul Gulati

    Mehul Gulati

    Mehul is a Software Development Engineer for Amazon EMR at Amazon Web Services. His expertise spans big data systems including HBase, Hive, Tez, and distributed storage solutions. His customer obsession and focus on reliability helps Amazon EMR deliver reliable and efficient big data processing capabilities to customers.

    Modernize game intelligence with generative AI on Amazon Redshift

    Post Syndicated from Narendra Gupta original https://aws.amazon.com/blogs/big-data/modernize-game-intelligence-with-generative-ai-on-amazon-redshift/

    Game studios generate massive amounts of player and gameplay telemetry, but transforming that data into meaningful insights is often slow, technical, and dependent on SQL expertise. With the new Amazon Redshift integration for Amazon Bedrock Knowledge Bases, teams can unlock instant, AI-powered analytics by asking questions in natural language. Analysts, product managers, and designers can now explore Amazon Redshift data conversationally—no query writing required—and Amazon Bedrock automatically generates optimized SQL, executes it on Amazon Redshift, and returns clear, actionable answers. This brings together the scale and performance of Amazon Redshift with the intelligence of Amazon Bedrock, enabling faster decisions, deeper player understanding, and more engaging game experiences.

    Amazon Redshift can be used as a structured data source for Amazon Bedrock Knowledge Bases, allowing for natural language querying and retrieval of information from Amazon Redshift. Amazon Bedrock Knowledge Bases can transform natural language queries into SQL queries, so users can retrieve data directly from the source without needing to move or preprocess the data. A game analyst can now ask, “How many players completed all the levels in a game?” or “List the top 5 players by the number of times the game was played,” and Amazon Bedrock Knowledge Bases automatically translates that query into SQL, runs the query against Amazon Redshift, and returns the results—or even provides a summarized narrative response.

    To generate accurate SQL queries, Amazon Bedrock Knowledge Bases uses database schema, previous query history, and other domain or business knowledge such as table and column annotations that are provided about the data sources. In this post, we discuss some of the best practices to improve accuracy while interacting with Amazon Bedrock using Amazon Redshift as the knowledge base.

    Solution overview

    In this post, we illustrate the best practices using gaming industry use cases. You will converse with players and their game attempts data in natural language and get the response back in natural language. In the process, you will learn the best practices. To follow along with the use case, follow these high-level steps:

    1. Load game attempts data into the Redshift cluster.
    2. Create a knowledge base in Amazon Bedrock and sync it with the Amazon Redshift data store.
    3. Review the approaches and best practices to improve the accuracy of response from the knowledge base.
    4. Complete the detailed walkthrough for defining and using curated queries to improve the accuracy of responses from the knowledge base.

    Prerequisites

    To implement the solution, you need to complete the following prerequisites:

    Load game attempts and players data

    To load the datasets to Amazon Redshift, complete the following steps:

    1. Open Amazon Redshift Query Editor V2 or another SQL editor of your choice and connect to the Redshift database.
    2. Run the following SQL to create the data tables to store games attempts and player details:
      CREATE TABLE game_attempts (
          player_id numeric(10, 0), -- Player ID.
          level_id numeric(5, 0), -- Game level ID
          f_success integer, -- Indicates whether user completed the level (1: completed, 0: fails).
          f_duration real, -- duration of the attempt.  Units in seconds
          f_reststep real, -- The ratio of the remaining steps to the limited steps.  Failure is 0.
          f_help integer, -- Whether extra help, such as props and hints, was used.  1- used, 0- not used
          game_time timestamp, -- Attempt timestamp
          bp_used boolean -- Whether bonus packages used or not.  true: used, false: not used.
      );
      CREATE TABLE players (
      	player_id numeric(10, 0), -- Player ID
      	lost_label boolean, -- Indicated if user retained or lost.  true: lost ,  false: retained
      	bp_category integer -- bonus package category codes
      );

    3. Download the game attempts and players datasets to your local storage.
    4. Create an Amazon Simple Storage Service (Amazon S3) bucket with a unique name. For instructions, refer to Creating a general purpose bucket.
    5. Upload the downloaded files into your newly created S3 bucket.
    6. Using the following COPY command statements, load the datasets from Amazon S3 into the new tables you created in Amazon Redshift. Replace <<your_s3_bucket>> with the name of your S3 bucket and <<your_region>> with your AWS Region:
      COPY game_attempts 
      FROM 's3://<<your_s3_bucket>>/game_attempts.csv' 
      IAM_ROLE DEFAULT 
      FORMAT AS CSV 
      IGNOREHEADER 1;
      COPY players
      FROM 's3://<<your_s3_bucket>>/players.csv' 
      IAM_ROLE DEFAULT 
      FORMAT AS CSV 
      IGNOREHEADER 1;

    Create knowledge base and sync

    To create a knowledge base and sync your data store with your knowledge base, complete these steps:

    1. Follow the steps at Create a knowledge base by connecting to a structured data store.
    2. Follow the steps at Sync your structured data store with your Amazon Bedrock knowledge base.

    Alternatively, you can refer Step 4: Set up Bedrock Knowledge Bases in Accelerating Genomic Data Discovery with AI-Powered Natural Language Queries in the AWS for Industries blog.

    Approaches to improve the accuracy

    If you’re not getting the expected response from the knowledge base, you can consider these key strategies:

    1. Provide additional information in the Query Generation Configuration. The knowledge base’s response accuracy can be improved by providing supplementary information and context to help it better understand your specific use case.
    2. Use representative sample queries. Running example queries that reflect common use cases helps train the knowledge base on your database’s specific patterns and conventions.

    Consider a database that stores player information using country codes rather than full country names. By running sample queries that demonstrate the relationship between country names and their corresponding codes (for example, “USA” for “United States”), you help the knowledge base understand how to properly translate user requests that reference full country names into queries using the correct country codes. This approach helps connect natural language requests and your database’s specific implementation details, resulting in more accurate query generation.

    Before we dive into more optimizations options, let’s explore how you can personalize the query engine to generate queries for a specific query engine. In this walkthrough, we use Amazon Redshift. Amazon Bedrock Knowledge Bases analyzes three key components to generate accurate SQL queries:

    • Database metadata
    • Query configurations
    • Historical query and conversation data

    The following graphic illustrates this flow.

    Amazon Bedrock Knowledge Bases architecture diagram showing structured data retrieval workflow with generative AI

    You can configure these settings to enhance query accuracy in two ways:

    • When creating a new Amazon Redshift knowledge base
    • By editing the query engine settings of an existing knowledge base

    To configure setting when creating new knowledge base, follow steps on Create a knowledge base by connecting to a structured data store and configure below parameters in (Optional) Query configurations section as shown in following screenshot:

    1. Table and column descriptions
    2. Table and column inclusions/exclusions
    3. Curated queries

    Amazon Bedrock Knowledge Base creation interface showing Redshift database configuration options

    To configure setting when editing the query engine of an existing knowledge base, follow these steps:

    1. On the Amazon Bedrock console in the left navigation pane, choose Knowledge Bases and select your Redshift Knowledge Base.
    2. Choose your query engine and choose Edit,
    3. Configure below parameters in (Optional) Query configurations section as shown in following screenshot:
      1. Table and column descriptions
      2. Table and column inclusions/exclusions
      3. Curated queries

    Edit query engine configuration page for Amazon Bedrock Knowledge Base with Redshift settings

    Let’s explore the available query configuration options in more detail to understand how these help the knowledge base generate a more accurate response.

    Table and column descriptions provide essential metadata that helps Amazon Bedrock Knowledge Bases understand your data structure and generate more accurate SQL queries. These descriptions can include table and column purposes, usage guidelines, business context, and data relationships.

    Follow these best practices for descriptions:

    • Use clear, specific names instead of abstract identifiers
    • Include business context for technical fields
    • Define relationships between related columns

    For example, consider a gaming table with timestamp columns named t1, t2, and t3. Adding these descriptions helps the knowledge base generate appropriate queries. For example, if t1 is play start time, t2 is play end time, and t3 is record creation time, adding these descriptions will indicate to the knowledge base to use t2–t1 for finding the game duration.

    Curated queries are a set of predefined question and answer examples. Questions are written as natural language queries (NLQs) and answers are the corresponding SQL query. These examples help the SQL generation process by providing examples of the kinds of queries that should be generated. They serve as reference points to improve the accuracy and relevance of generative SQL outputs. Using this option, you can provide some example queries to the knowledge base for it understand custom vocabulary also. For example, if the country field in the table is populated with a country code, adding an example query will help the knowledge base to convert the country name to a country code before running the query to answer questions on the data of players in a specific country. You can also provide some example complex queries to help the knowledge base to respond to more complex questions. The following is an example query that can be added to the knowledge base:

    Select count(*) from players_address where country = ‘USA’;
    

    With table and column inclusion and exclusion, you can specify a set of tables or columns to be included or excluded for SQL generation. This field is crucial if you want to limit the scope of SQL queries to a defined subset of available tables or columns. This option can help optimize the generation process by reducing unnecessary table or column references. You can also use this option to:

    • Exclude redundant tables, for example, those generated by copying the original table to run a complex analysis
    • Exclude tables and columns containing sensitive data

    If you specify inclusions, all other tables and columns are ignored. If you specify exclusions, the tables and columns you specify are ignored.

    Walkthrough for defining and using curated queries to improve accuracy

    To define and use curated queries to improve accuracy, complete the following steps.

    1. On the AWS Management Console, navigate to Amazon Bedrock and in the left navigation pane, choose Knowledge Bases. Select the knowledge base you created with Amazon Redshift.
    2. Choose Test Knowledge Base, as shown in the following screenshot, to validate the accuracy of the knowledge base response.
      Amazon Bedrock Knowledge Base overview page showing game-rs-kb configuration and status details
    3. On the Test Knowledge Base screen under Retrieval and response generation, choose Retrieval and response generation: data sources and model.
    4. Choose Select model to pick a large language model (LLM) to convert the SQL query response from the knowledge base to a natural language response.
    5. Choose Nova Pro in the popup and choose Apply, as shown in the following screenshot.
      Model selection dialog showing Amazon Nova Pro and other foundation models for Bedrock Knowledge Base

    Now you have Amazon Nova Pro connected to your knowledge base to respond to your queries based on the data available in Amazon Redshift. You can ask some questions and verify them with actual data in Amazon Redshift. Follow these steps:

    1. In the Test section on the right, enter the following prompt, then choose the send message icon, as shown in the following screenshot.
      What is the latest attempt status for player 12004?

      Amazon Bedrock Knowledge Base test interface with configuration panel and preview section

    2. Amazon Nova Pro generates a response using the data stored in the Redshift knowledge base.
    3. Choose Details to see the SQL query generated and used by Amazon Nova Pro, as shown in the following screenshot.
      Test results showing AI-generated response with source details for player attempt status query
    4. Copy the query and enter it in query editor v2 of the Redshift knowledge base, as shown in the following screenshot.
      AWS Redshift Query Editor showing SQL query execution with player game attempt results
    5. Verify that the response generated by Amazon Nova Pro in natural language matches the data in Amazon Redshift and that the generated SQL query is also accurate.

    You can try some more questions to verify the Amazon Nova Pro response, for example:

    What is the lost status for player ID 12004?
    How many levels did the player 12004 play?
    What level did player 12004 play the most?
    Show me the summary of all 14 attempts by player 12004 for level 76.

    But what if the response generated by the knowledge base isn’t accurate? In those cases, you can add additional context the knowledge base can use to provide more accurate responses. For example, try asking the following question:

    How many total players are there?

    In this case, the response generated by the knowledge base doesn’t match the actual player count in Amazon Redshift. The knowledge base reported about 13,589 players and generated the following query to get the player count:

    SELECT COUNT(DISTINCT player_id) AS "Number of Players" FROM games.game_attempts;

    The following screenshot shows this question and result.

    Test preview showing AI response to player count query with citation

    The knowledge base should have used the players table in Amazon Redshift to find the unique players. The correct response is 10,816 players.

    AWS Redshift Query Editor showing COUNT query result of 10,816 players

    To help the knowledge base, add a curated query for it to use the players table instead of the attempts table to find the total player count. Follow these steps:

    1. On the Amazon Bedrock console in the left navigation pane, choose Knowledge Bases and select your Redshift Knowledge Base.
    2. Choose your query engine and choose Edit, as shown in the following screenshot.
      Amazon Bedrock Query Engine configuration page showing Redshift serverless connection details
    3. Expand the Curated queries section and enter the following:
    4. In the Questions field, enter How many total players are there?.
    5. In the Equivalent SQL query field, enter SELECT count(*) FROM “dev”,“games”,“players”;.
    6. Choose Submit, as shown in the following screenshot.
      Edit query engine page showing curated query example for player count
    7. Navigate back to your knowledge base and query engine. Choose Sync to sync the knowledge base. This starts the metadata ingestion process so that data can be retrieved. The metadata allows Amazon Bedrock Knowledge Bases to translate user prompts into a query for the connected database. Refer to Sync your structured data store with your Amazon Bedrock knowledge base for more details.
    8. Return to Test Knowledge Base with Amazon Nova Pro and repeat the question about how many total players there are, as shown in the following screenshot. Now, the response generated by the knowledge base matches the data in player table in Amazon Redshift, and the query generated by the knowledge base uses the curated query with the player table instead of the attempts table to determine the player count.
      Test results showing total player count query with SQL source details

    Cleanup

    For the walkthrough section, we used serverless services, and your cost will be based on your usage of these services. If you’re using provisioned Amazon Redshift as a knowledge base, follow these steps to stop incurring charges:

    1. Delete the knowledge base in Amazon Bedrock.
    2. Shut down and delete your Redshift cluster.

    Conclusion

    In this post, we discussed how you can use Amazon Redshift as a knowledge base to provide additional context to your LLM. We identified best practices and explained how you can improve the accuracy of responses from the knowledge base by following these best practices.


    About the authors

    Narendra Gupta

    Narendra Gupta

    Narendra is a Specialist Solutions Architect at AWS, helping customers on their cloud journey with a focus on AWS analytics services. Outside of work, Narendra enjoys learning new technologies, watching movies, and visiting new places.

    Satesh Sonti

    Satesh Sonti

    Satesh is a Principal Analytics Specialist Solutions Architect based out of Atlanta, specializing in building enterprise data platforms, data warehousing, and analytics solutions. He has over 19 years of experience in building data assets and leading complex data platform programs for banking and insurance clients across the globe.