Tag Archives: Advanced (300)

Access a VPC-hosted Amazon OpenSearch Service domain with SAML authentication using AWS Client VPN

Post Syndicated from Jan Michael Go Tan original https://aws.amazon.com/blogs/big-data/access-a-vpc-hosted-amazon-opensearch-service-domain-with-saml-authentication-using-aws-client-vpn/

Customers often want to deploy Amazon OpenSearch Service domains in virtual private clouds (VPC) and use single sign-on (SSO) with SAML for access control to enhance security. However, setting this up can be challenging.

In this post, we explore different OpenSearch Service authentication methods and network topology considerations. Then we show how to build an architecture to access an OpenSearch Service domain hosted in a VPC using AWS Client VPN, AWS Transit Gateway, and AWS IAM Identity Center.

Solution overview

The following diagram illustrates the solution architecture.

High-level network diagram

The end-user authenticates with IAM Identity Center and connects to the AWS environment from their browser through Client VPN. The traffic is routed from the VPN VPC to the database VPC where the OpenSearch service endpoints are deployed. The user then authenticates to OpenSearch Service through IAM Identity Center. This architecture provides a scalable, enterprise-grade solution that avoids using bastion hosts while making sure only authorized users can access your OpenSearch Service domains through a secure VPN connection. In the following sections, we walk through the steps to set up IAM Identity Center, configure Transit Gateway to facilitate communication between VPCs, and configure SAML-based authentication using IAM Identity Center for both OpenSearch Service and VPN access. Prior experience setting up Client VPN, IAM Identity Center, and Transit Gateway would be beneficial but is not necessary to follow along with this post.

OpenSearch Service authentication methods and SAML

OpenSearch Service supports multiple authentication methods. You can use AWS Identity and Access Management (IAM) to call the OpenSearch Service configuration API (for details, see Making and signing OpenSearch Service requests). However, this doesn’t give you access to the visual dashboard. To access the visual dashboard and call the OpenSearch Service configuration API, you can use the OpenSearch Service built-in internal user database or Amazon Cognito for authentication and user management features. However, these options use separate user pools, which adds additional security and management overhead when adding and removing users.

Therefore, many customers choose to use SAML federation to integrate OpenSearch Service authentication with their existing identity providers like Entra ID, Okta, or JumpCloud. For this post, we use the IAM Identity Center directory as our identity source. One limitation of this approach is that it only supports identity provider-initiated authentication. This means that users must log in through the IAM Identity Center portal and then access their OpenSearch Service dashboard from there.

Private network topology options for OpenSearch Service

When deploying OpenSearch Service domains in a private VPC, organizations must establish secure and reliable network connectivity to access their OpenSearch Service domains. AWS offers several networking solutions that can be implemented individually or in combination to meet specific access requirements. These options include Transit Gateway for centralized network management, AWS Direct Connect or AWS Site-to-Site VPN for on-premises connectivity, and Client VPN for secure remote access. Each solution provides unique benefits and can be combined to meet different organizational needs, security requirements, and performance expectations.

AWS Transit Gateway

Transit Gateway functions as a cloud router that simplifies network connectivity by acting as a central hub for connecting VPCs and on-premises networks. Implementing Transit Gateway with OpenSearch Service enables consolidated access to your OpenSearch Service domain across multiple VPCs and AWS accounts. Through Transit Gateway route tables, you can precisely control traffic flow between attached networks. It supports transitive routing between VPCs and on-premises networks, significantly reducing the number of peering connections needed to access your OpenSearch Service domain. This centralized approach is a common pattern used by customers, which makes network management scalable as your infrastructure grows.

AWS Client VPN

With Client VPN, you can securely access your private OpenSearch Service domain through a managed OpenVPN-based solution. Using Client VPN removes the need to use a bastion host or proxy server to access an OpenSearch Service domain, reducing your management burden and improving security. Client VPN supports both certificate-based and SAML-based authentication. Client VPN endpoints can be associated with multiple subnets to provide high availability. The service includes comprehensive security features such as connection logging and security group controls.

For more information on VPC connectivity options, refer to the AWS Direct Connect whitepaper.

Combining Client VPN with Transit Gateway provides a scalable and flexible way to access an OpenSearch Service domain in a private VPC. In the subsequent sections, we walk you through how to integrate the various services.

Prerequisites

If you haven’t yet set up IAM Identity Center, refer to Enable IAM Identity Center to enable it. Both organization instances and account instances will work. The Identity Center instance must be deployed in the same AWS Region as your OpenSearch Service domain.

After you set up IAM Identity Center, complete the following steps to create an IAM Identity Center group:

  1. On the IAM Identity Center console, choose Groups in the navigation pane.
  2. Choose Create group and create a group (for this example, we name the group vpn_users.
  3. After you create the group, choose the group name to open its details page.
  4. Locate the group ID under General information. Save this in a text editor.
    IAM Identity Center Group ID
  5. Create a user (or multiple users) and assign them to the vpn_users group. This can be done directly through the user creation flow or after creating the user.

Set up the initial network topology

For this post, we use the network topology shown in the following diagram. One VPC hosts the client VPN endpoint with CIDR range 10.0.0.0/16 and a separate VPC with CIDR range 10.1.0.0/16 that hosts our OpenSearch Service nodes. The two VPCs are connected with Transit Gateway. The CIDR ranges in your environment may vary. The only requirement is that they can’t overlap.

Network topology

Complete the following steps to create the two VPCs using Amazon Virtual Private Cloud (Amazon VPC):

  1. On the Amazon VPC console, choose Create VPC.
  2. Choose VPC and more.
  3. For this post, name the VPC VPN-VPC and use 10.0.0.0/16 for the IPv4 CIDR block.
  4. Choose 3 for the number of Availability Zones.
  5. Choose 0 for the number of public subnets.
  6. Choose 3 for the number of private subnets.
  7. Choose None for the number of NAT gateways.
  8. Choose None for the number of VPC endpoints.
    Initial VPC configuration
  9. Repeat these steps to create the second VPC for the OpenSearch Service domain. Keep the same configuration settings except for the following:
    1. Name: Database-VPC
    2. IPv4 CIDR Block: 10.1.0.0/16

Configure Transit Gateway

Follow the instructions in Create an AWS Transit Gateway using the Amazon VPC Console to create a transit gateway and attach your VPCs to it.

Next, you must update each VPC route table to facilitate connectivity to the OpenSearch Service domain.

  1. On the Amazon VPC console, choose Route tables in the navigation pane.
  2. For VPN-VPC, add routes on the subnets where the Client VPN endpoints are attached. The route is 10.1.0.0/16 using Transit Gateway. This route allows VPN users to reach Database-VPC.
    Route table
  3. For Database-VPC, add routes on the subnets of the OpenSearch Service domain endpoint. The route is 10.0.0.0/16 using Transit Gateway. This route allows responses from Database-VPC back to reach the VPN users.
    OpenSearch Route Table

    Next, you must update the Transit Gateway Security Group Referencing support configuration. This allows the OpenSearch Service domain’s security group to open port 443 to only the Client VPN security group. This makes applying least privilege simpler.

  4. On the Transit Gateway console, select the transit gateway you’re using.
  5. On the Actions menu, choose Modify transit gateway.
    Modify TGW
  6. Select Security Group Referencing support and choose Modify transit gateway.
    TGW Security Group Configuration

Configure Client VPN authentication

Client VPN can be associated to multiple VPC subnets for high availability. Client VPN supports multiple client authentication methods. For this post, we use SAML-based authentication with IAM Identity Center.

To set up SAML-based authentication with IAM Identity Center, follow the instructions in the following sections. For more details, refer to Authenticate AWS Client VPN users with AWS IAM Identity Center. Deploy and associate the Client VPN endpoint with VPN-VPC.

Configure Client VPN access to database VPC

During the initial setup of the Client VPN endpoint, you defined authorization rules that authorized the VPN_users group to access the VPN-VPC network, which is 10.0.0.0/16.Complete the following steps to add connectivity to database-VPC:

  1. On the Amazon VPC console, choose Client VPC endpoints in the navigation pane.
  2. Select the endpoint you created.
  3. In the Authorization rules section, choose Add authorization rules.
    ClientVPN Auth Rules
  4. For Destination network to enable access, enter 10.1.0.0/16 (this is the database VPC).
  5. For Grant access to, select Allow access to all users.
  6. Choose Add authorization rule.
    ClientVPN Add Auth Rule

    After you create the authorization rule, the user now has access to that CIDR range. Next, you add an entry in the Client VPN endpoint’s route table to provide reachability from a network perspective.

  7. On the Client VPN endpoints page, select the endpoint you just created.
  8. In the Route table section, choose Create route.
    ClientVPN Route
  9. For Route destination, enter the CIDR range for Database-VPC (10.1.0.0/16).
  10. For Subnet ID for target network association, choose a subnet ID.
  11. Choose Create route.
    ClientVPN Create Route

You should see the new route in the “Creating” state. After it has reached the “Active” state, VPN users will have a network path to the database VPC to be able to reach the OpenSearch Service domain.

ClientVPN Route Creating State

Configure Client VPN application on your client

Complete the following steps to configure the Client VPN application to your client:

  1. Download the relevant installer for Client VPN for Desktop and install Client VPN.
  2. Download and prepare the Client VPN endpoint file.
  3. Open the Client VPN application.
  4. Choose Manage Profile, then choose Add Profile.
  5. Enter a display name and upload the VPN configuration file.
  6. Choose Add Profile.

Set up federation with IAM Identity Center with OpenSearch Service

Complete the following steps to set up federation with IAM Identity Center with OpenSearch Service:

  1. Create an OpenSearch Service domain in the database VPC.
  2. Set up the SAML integration between OpenSearch Service and IAM Identity Center. Assign the same groups that you assigned to the VPN custom application to the OpenSearch Service custom application.
  3. Modify the security group associated with the OpenSearch Service domain to allow access from the Client VPN subnet.
  4. Modify the security group of Client VPN and add the following entry:
    1. Type: HTTPS
    2. Source: Use Custom and reference the security group of the OpenSearch Service domain

Test the end-to-end flow

Now you can test the entire flow end-to-end:

  1. Run Client VPN on your local machine. Use the profile that you previously configured.
    The client will prompt you to authenticate with IAM Identity Center. After authentication, you will see the message “Authentication details received, processing details. You may close this window at any time.”
  2. Access your IAM Identity Center access portal URL (this can be found on the IAM Identity Center console, under Dashboard). Sign in as a user that has been assigned to the OpenSearch Service custom application in the previous step.
  3. After authentication, choose the Applications tab in AWS Access Portal and choose the OpenSearch Service application.

This should redirect you to the OpenSearch Service Dashboards page with the role that you assigned.

IAM Identity Center - App List

Clean up

After you test the solution, delete the resources you created to avoid incurring future charges:

  1. Delete the OpenSearch Service domain and the SAML application, users, and groups in IAM Identity Center.
  2. Delete the client VPN endpoints that you created and remove the routing rules from Transit Gateway.

Conclusion

In this post, we discussed the networking options for securely accessing an OpenSearch Service domain deployed in a private VPC through services like Transit Gateway, Client VPN, and Site-to-Site VPN. We also discussed how to use IAM Identity Center for authentication and authorization, helping you simplify identity management for OpenSearch Service. If you have feedback about this post, provide it in the comments section.


About the authors

Jan Michael Go Tan

Jan Michael Go Tan

Jan Michael is a Principal Solutions Architect for Amazon Web Services. He helps customers design scalable and innovative solutions with the AWS Cloud.

Kevin Low

Kevin Low

Kevin is a Security Solutions Architect at AWS who helps the largest customers across ASEAN build securely. He specializes in threat detection and incident response and is passionate about integrating resilience and security. Outside of work, he loves spending time with his wife and dog, a poodle called Noodle.

Apache Spark 4.0.1 preview now available on Amazon EMR Serverless

Post Syndicated from Al MS original https://aws.amazon.com/blogs/big-data/apache-spark-4-0-1-preview-now-available-on-amazon-emr-serverless/

Amazon EMR Serverless now supports Apache Spark 4.0.1 in preview, making analytics accessible to more users, simplifying data engineering workflows, and strengthening governance capabilities. The release introduces ANSI SQL compliance, VARIANT data types support for JSON handling, Apache Iceberg v3 table format support, and enhanced streaming capabilities. This preview is available in all regions where EMR Serverless is available.

In this post, we explore key benefits, technical capabilities, and considerations for getting started with Spark 4.0.1 on Amazon EMR Serverless—a serverless deployment option that simplifies running open-source big data frameworks, without requiring managing clusters. With the emr-spark-8.0-preview release label, you can evaluate new SQL capabilities, Python API improvements, and streaming enhancements in your existing EMR Serverless environment.

Benefits

Spark 4.0.1 helps you solve data engineering problems with specific improvements. This section shows how new capabilities help with real-world scenarios.

Make analytics accessible to more users

Simplify Extract Transform Load (ETL) development with SQL scripting. Data engineers often switch between SQL and Python to build complex ETL logic with control flow. SQL scripting in Spark 4.0.1 enables loops, conditionals, and session variables directly in SQL, reducing context-switching and simplifying pipeline development. Use pipe syntax (|>) to chain operations for more readable, maintainable queries.

Improve data quality with ANSI SQL mode. Silent type conversion failures can introduce data quality issues. ANSI SQL mode (now default) enforces standard SQL behavior, raising errors for invalid operations instead of producing unexpected results. Important: ANSI SQL mode is now enabled by default. Test your queries thoroughly during this preview evaluation.

Simplify data engineering workflows

Process JSON data efficiently with VARIANT. Teams working with semi-structured data often see slow performance from repeated JSON parsing. The VARIANT data type stores JSON in an optimized binary format, eliminating parsing overhead. You can efficiently store and query JSON data in data lakes without schema rigidity.

Build Python data sources without Scala. Integrating custom data sources previously required Scala expertise. The Python data Source API lets you build connectors entirely in Python, using existing Python skills and libraries without learning a new language.

Debug streaming applications with queryable state. Troubleshooting stateful streaming applications has historically required indirect methods. The new state data source reader shows streaming state as queryable DataFrames. You can inspect state during debugging, test state values in unit tests, and diagnose production incidents.

Strengthen governance capabilities

Establish comprehensive audit trails with Apache Iceberg v3. The Apache Iceberg v3 table format provides transaction guarantees and tracks data changes over time, giving you the audit trails needed for regulatory compliance. When combined with VARIANT data type support, you can maintain governance controls while handling semi-structured data efficiently in data lakes.

Key capabilities

Spark 4.0.1 Preview on EMR Serverless introduces four major capability areas:

  1. SQL enhancements – ANSI mode, pipe syntax, VARIANT type, SQL scripting, user-defined functions (UDFs)
  2. Python API advances – custom data sources, UDF profiling
  3. Streaming improvements – stateful processing API v2, queryable state
  4. Table format support – Amazon S3 Tables, AWS Lake Formation integration

The following sections provide technical details and code examples for each capability.

SQL enhancements

Spark 4.0.1 introduces new SQL capabilities including ANSI mode compliance, SQL UDFs, pipe syntax for readable queries, VARIANT type for JSON handling, and SQL scripting with control flow.

ANSI SQL mode by default

ANSI SQL mode is now enabled by default, enforcing standard SQL behavior for data integrity. Silent casting of out-of-range values now raises errors rather than producing unexpected results. Existing queries may behave differently, particularly around null handling, string casting, and timestamp operations. Use spark.sql.ansi.enabled=false if you need legacy behavior during migration.

SQL pipe syntax

You can now chain SQL operations using the |> operator for improved readability. The following example shows how you can replace nested subqueries with a more maintainable pipeline:

FROM customer
|> LEFT OUTER JOIN orders ON c_custkey = o_custkey
|> AGGREGATE COUNT(o_orderkey) c_count GROUP BY c_custkey
|> AGGREGATE COUNT(*) AS custdist GROUP BY c_count
|> ORDER BY custdist DESC

This replaces nested subqueries, making complex transformations easier to understand and maintain.

VARIANT data type

The VARIANT type handles semi-structured JSON/XML data efficiently without repeated parsing. It uses an optimized binary representation internally while maintaining schema-less flexibility. Previously, JSON expressions required repeated parsing, degrading performance. VARIANT eliminates this overhead. The following snippet shows how to parse JSON into the VARIANT type:

df = spark.sql("SELECT parse_json('{\"name\":\"Alice\",\"age\":30}') as data")

Spark 4.0.1 on EMR Serverless supports Apache Iceberg v3, enabling the VARIANT data type with Iceberg tables. This combination provides efficient storage and querying of semi-structured JSON data in your data lake. Store VARIANT columns in Iceberg tables and use Iceberg’s schema evolution and time travel capabilities alongside Spark’s optimized JSON processing. The following example shows how to create an Iceberg table with a VARIANT column:

CREATE TABLE catalog.db.events (
  event_id BIGINT,
  event_data VARIANT,
  timestamp TIMESTAMP
) USING iceberg;

INSERT INTO catalog.db.events SELECT 1, parse_json('{"user":"alice","action":"login"}'), current_timestamp();

SQL scripting with session variables

Manage state and control flow directly in SQL using session variables, and IF/WHILE/FOR statements. The following example demonstrates a loop that populates a results table:

BEGIN
  DECLARE counter INT = 10;
  WHILE counter > 0 DO
    INSERT INTO results VALUES (counter);
    SET counter = counter - 1;
  END WHILE;
END

This enables complex ETL logic entirely in SQL without switching to Python.

SQL user-defined functions

Define custom functions directly in SQL. Functions can be temporary (session-scoped) or permanent (catalog-stored). The following example shows how to register and use a simple UDF:

CREATE FUNCTION plusOne(x INT) RETURNS INT RETURN x + 1;
SELECT plusOne(5);

Python API advances

This section covers new Python capabilities including custom data sources and UDF profiling tools.

Python data source API

You can now build custom data sources in Python without Scala knowledge. The following example shows how to create a simple data source that returns sample data:

from pyspark.sql.datasource import DataSource, DataSourceReader
from pyspark.sql.types import StructType, StructField, StringType, IntegerType

class SampleDataSource(DataSource):
    def schema(self):
        return StructType([
            StructField("name", StringType()),
            StructField("age", IntegerType())
        ])
    
    def reader(self, schema):
        return SampleReader()

class SampleReader(DataSourceReader):
    def read(self, partition):
        yield ("Alice", 30)
        yield ("Bob", 25)

# Register and use
spark.dataSource.register(SampleDataSource)
spark.read.format("SampleDataSource").load().show()

Unified UDF profiling

Profile Python and Pandas UDFs for performance and memory insights. The following code enables performance profiling:spark.conf.set("spark.sql.pyspark.udf.profiler", "perf") # or "memory"

Structured streaming enhancements

This section covers improvements to stateful stream processing, including queryable state and enhanced state management APIs.

Arbitrary stateful processing API v2

The transformWithState operator provides robust state management with timer and TTL support for automatic cleanup, schema evolution capabilities, and initial state support for pre-populating state from batch DataFrames.

State data source reader

Query streaming state as a DataFrame for debugging and monitoring. Previously, state data was internal to streaming queries. Now you can verify state values in unit tests, diagnose production incidents, detect state corruption, and optimize performance. Note: This feature is experimental. Source options and behavior may change in future releases.

State store improvements

Upgraded changelog checkpointing for RocksDB removes performance bottlenecks. Enhanced checkpoint coordination and improved sorted string table (SST) file reuse management optimizes streaming operations.

Table format support

This section covers support for AWS S3 Tables and full table access (FTA) with AWS Lake Formation.

AWS S3 Tables

Use Spark 4.0.1 with AWS S3 Tables, a storage solution that provides managed Apache Iceberg tables with automatic optimization and maintenance. S3 Tables simplify data lake operations by handling compaction, snapshot management, and metadata cleanup automatically.

Full table access with Lake Formation

FTA is supported for Apache Iceberg, Delta Lake, and Apache Hive tables when using AWS Lake Formation, a managed service that simplifies data access control. FTA provides coarse-grained access control at the table level. Note that fine-grained access control (FGAC) with column-level or row-level permissions is not available in this preview.

Getting started

Follow these steps to create an EMR Serverless application, run sample code to test new features, and provide feedback on the preview.

Prerequisites

Before you begin, confirm you have the following:

Note: EMR Studio Notebooks and SageMaker Unified Studio are not supported during this preview. Use the AWS CLI or AWS SDK to submit jobs.

Step 1: Create your EMR Serverless application

Create or update your application with the emr-spark-8.0-preview release label. The following command creates a new application:

aws emr-serverless create-application --type spark \
  --release-label emr-spark-8.0-preview \
  --region us-east-1 --name spark4-test

Step 2: Test sample code

Run this PySpark job to verify setup and test Spark 4.0.1 features:

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("Spark 4.0.1 Test").getOrCreate()
print(f"Spark Version: {spark.version}")

# Create sample data
data = [("Alice", 34, "Engineering"), ("Bob", 45, "Sales"),
        ("Charlie", 28, "Engineering"), ("Diana", 52, "Marketing")]
df = spark.createDataFrame(data, ["name", "age", "department"])
df.createOrReplaceTempView("employees")

# Test SQL PIPE syntax
try:
    result = spark.sql("""
        FROM employees
        |> WHERE age > 30
        |> SELECT name, age, department
        |> ORDER BY age DESC
    """)
    result.show()
    print("✓ SQL pipe syntax test passed")
except Exception as e:
    print(f"✗ SQL pipe syntax test failed: {e}")

# Test VARIANT data type
try:
    json_data = spark.sql("""
        SELECT parse_json('{"name":"Alice","skills":["Python","Spark","SQL"]}') as data
    """)
    json_data.show(truncate=False)
    print("✓ VARIANT data type test passed")
except Exception as e:
    print(f"✗ VARIANT data type test failed: {e}")

Submit the job with the following command:

aws emr-serverless start-job-run \
    --application-id <your-application-id> \
    --execution-role-arn <your-execution-role-arn> \
    --job-driver '{
        "sparkSubmit": {
            "entryPoint": "s3://<your-bucket>/spark_4_test.py",
            "sparkSubmitParameters": "--conf spark.executor.cores=4 --conf spark.executor.memory=16g"
        }
    }'

Step 3: Test your workloads

Review the Spark SQL Migration Guide and PySpark Migration Guide, then test production workloads in non-production environments. Focus on queries affected by ANSI SQL mode and benchmark performance.

Step 4: Clean up resources

After testing, delete all resources created during this evaluation to avoid ongoing charges:

# Delete the EMR Serverless application
aws emr-serverless delete-application \
    --application-id spark4-test \
    --region us-east-1
# Remove the test script from S3
aws s3 rm s3://<your-bucket>/spark_4_test.py

Migration considerations

Before evaluating Spark 4.0.1, review the updated runtime requirements and behavioral changes that may affect your existing code.

Runtime requirements

  • Scala: Version 2.13.16 required (2.12 support dropped)
  • Java: JDK 17 or higher required (JDK 8 and 11 support removed)
  • Python: Version 3.9+ required, continued support for 3.11 and newly added 3.12 (3.8 support removed)
  • Pandas: Minimum version 2.0.0 (previously 1.0.5)
  • SparkR: Deprecated; migrate to PySpark

Behavioral changes

With ANSI SQL mode enforcement, you may see different behavior in:

  • Null handling: Stricter null propagation in expressions
  • String casting: Invalid casts now raise errors instead of returning null
  • Map key operations: Duplicate keys now raise errors
  • Timestamp conversions: Overflow returns null instead of wrapped values
  • CREATE TABLE statements: Now respect the spark.sql.sources.defaultconfiguration instead of defaulting to Hive format when USING or STORED AS clauses are omitted

You can control many of these behaviors via legacy configuration flags. Consult the official migration guides for details refer: Spark SQL Migration Guide: 3.5 to 4.0 and PySpark Migration Guide: 3.5 to 4.0.

Preview limitations

The following capabilities are not available in this preview:

  • Fine-grained access control: Fine-grained access control (FGAC) with row-level or column-level filtering is not supported in this preview. Jobs with spark.emr-serverless.lakeformation.enabled=true will fail.
  • Spark Connect: Not supported in this preview. Use standard Spark job submission with the StartJobRun API.
  • Open Table Format limitations: Hudi is not supported in this preview. Delta 4.0.0 does not support Flink connectors (deprecated in Delta 4.0.0). Delta Universal Format is not supported in this preview.
  • Connectors: spark-sql-kinesis, emr-dynamodb, and spark-redshift are unavailable.
  • Interactive applications: Livy and JupyterEnterpriseGateway are not included. Also, SageMaker Unified Studio and EMR Studio are not supported.
  • EMR features: Serverless Storage and Materialized Views are not supported.

This preview lets you evaluate Spark 4.0.1’s core capabilities on EMR Serverless, including SQL enhancements, Python API improvements, and streaming state management. Test your migration path, assess performance improvements, and provide feedback to shape the general availability release.

Conclusion

This post showed you how to get started with the Apache Spark 4.0.1 preview release on Amazon EMR Serverless. You explored how the VARIANT data type works with Iceberg v3 to process JSON data efficiently, how SQL scripting and pipe syntax eliminate context-switching for ETL development, and how queryable streaming state simplifies debugging stateful applications. You also learned about the preview limitations, runtime requirements, and behavioral changes to consider during evaluation.

Test the Spark 4.0.1 preview on EMR Serverless and provide feedback through AWS Support to help shape the general availability release.

To learn more about Apache Spark 4.0.1 features, see the Spark 4.0.1 Release Notes. For EMR Serverless documentation, see the EMR Release Guide.

Resources

Apache Spark Documentation

Amazon EMR Resources


About the authors

Al MS

Al MS

Al is a product manager for Amazon EMR at Amazon Web Services.

Emilie Faracci

Emilie Faracci

Emilie is a Software Development Engineer Amazon Web Services, working on Amazon EMR. She focuses on Spark development and has contributed to open-source Apache Spark v4.0.1.

Karthik Prabhakar

Karthik Prabhakar

Karthik is a Data Processing Engines Architect for Amazon EMR at AWS. He specializes in distributed systems architecture and query optimization, working with customers to solve complex performance challenges in large-scale data processing workloads. His focus spans engine internals, cost optimization strategies, and architectural patterns that enable customers to run petabyte-scale analytics efficiently.

Exploring common centralized and decentralized approaches to secrets management

Post Syndicated from Brendan Paul original https://aws.amazon.com/blogs/security/exploring-common-centralized-and-decentralized-approaches-to-secrets-management/

One of the most common questions about secrets management strategies on Amazon Web Services (AWS) is whether an organization should centralize its secrets. Though this question is often focused on whether secrets should be centrally stored, there are four aspects of centralizing the secrets management process that need to be considered: creation, storage, rotation, and monitoring. In this post, we discuss the advantages and tradeoffs of centralizing or decentralizing each of these aspects of secrets management.

Centralized creation of secrets

When deciding whether to centralize secrets creation, you should consider how you already deploy infrastructure in the cloud. Modern DevOps practices have driven some organizations toward developer portals and internal developer platforms that use golden paths for infrastructure deployment. By using tools that use golden paths, developers can deploy infrastructure in a self-service model through infrastructure as code (IaC) while adhering to organizational standards.

A central function maintains these golden paths, such as a platform engineering team. Examples of services that can be used to maintain and define golden paths might include AWS services such as AWS Service Catalog or popular open source projects such as Backstage.io. Using this approach, developers can focus on application code while platform engineers focus on infrastructure deployment, security controls, and developer tooling. An example of a golden path might be a templatized implementation for a microservice that writes to a database.

For example, a golden path could define that a service or application must be built using the AWS Cloud Development Kit (AWS CDK), running on Amazon Elastic Container Service (Amazon ECS), and use AWS Secrets Manager to retrieve database credentials. The platform team could also build checks to help ensure that the secret’s resource policy only allows access to the role being used by the microservice and is encrypted with a customer managed key. This pattern abstracts deployments away from developers and facilitates resource deployment across accounts. This is one example of a centralized creation pattern, shown in Figure 1.

Figure 1: Architecture diagram highlighting the developer portal deployment pattern for centralized creation

Figure 1: Architecture diagram highlighting the developer portal deployment pattern for centralized creation

The advantages of this approach are:

  • Consistent naming tagging, and access control: When secrets are created centrally, you can enforce a standard naming convention based on the account, workload, service, or data classification. This simplifies implementing scalable patterns like attribute-based access control (ABAC).
  • Least privilege checks in CI/CD pipelines: When you create secrets within the confines of IaC pipelines, you can use APIs such as the AWS IAM Access Analyzer check-no-new-access API. Deployment pipelines can be templatized, so individual teams can take advantage of organizational standards while still owning deployment pipelines.
  • Create mechanisms for collaboration between platform engineering and security teams: Often, the shift towards golden paths and service catalogs is driven by a desire for a better developer experience and reduced operational overhead. A byproduct of this move is that security teams can partner with platform engineering teams to build security by default into these paths.

The tradeoffs of this approach are:

  • It takes time and effort to make this shift. You might not have the resources to invest in full-time platform engineering or DevOps teams. To centrally provision software and infrastructure like this, you must maintain libraries of golden paths that are appropriate for the use cases of your organization. Depending on the size of your organization, this might not be feasible.
  • Golden paths must keep up with the features of the services they support: If you’re using this pattern, and the service you’re relying on releases a new feature, your developers must wait for the features to be added to the affected golden paths.

If you want to learn more about the internal developer platform pattern, check out the re:Invent 2024 talk Elevating the developer experience with Backstage on AWS.

Decentralized creation of secrets

In a decentralized model, application teams own the IaC templates and deployment mechanisms in their own accounts. Here, each team is operating independently, which can make it more difficult to enforce standards as code. We’ll refer to this pattern, shown in Figure 2, as a decentralized creation pattern.

Figure 2: Decentralized creation of secrets

Figure 2: Decentralized creation of secrets

The advantages of this approach are:

  • Speed: Developers can move quickly and have more autonomy because they own the creation process. Individual teams don’t have a dependency on a central function.
  • Flexibility: You can still use features such as the IAM Access Analyzer check-no-new-access API, but it’s up to each team to implement this in their pipeline.

The tradeoffs of this approach are:

  • Lack of standardization: It can become more difficult to enforce naming and tagging conventions, because it’s not templatized and applied through central creation mechanisms. Access controls and resource policies might not be consistent across teams.
  • Developer attention: Developers must manage more of the underlying infrastructure and deployment pipelines.

Centralized storage of secrets

Some customers choose to store their secrets in a central account, and others choose to store secrets in the accounts in which their workloads live. Figure 3 shows the architecture for centralized storage of secrets.

Figure 3: Centralized storage of secrets

Figure 3: Centralized storage of secrets

The advantages of centralizing the storage of secrets are:

  • Simplified monitoring and observability: Monitoring secrets can be simplified by keeping them in a single account and with a centralized team controlling them.

Some tradeoffs of centralizing the storage of secrets are:

  • Additional operational overhead: When sharing secrets across accounts, you must configure resource policies on each secret that is shared.
  • Additional cost of AWS KMS Customer Managed Keys: You must use AWS Key Management Service (AWS KMS) customer managed keys when sharing secrets across accounts. While this gives you an additional layer of access control over secret access, it will increase cost under the AWS KMS pricing. It will also add another policy that needs to be created and maintained.
  • High concentration of sensitive data: Having secrets in a central account can increase the number of resources affected in the event of inadvertent access or misconfiguration.
  • Account quotas: Before deciding on a centralized secret account, review the AWS service quotas to ensure you won’t hit quotas in your production environment.
  • Service managed secrets: When services such as Amazon Relational Database Service (Amazon RDS) or Amazon Redshift manage secrets on your behalf, these secrets are placed in the same account as the resource with which the secret is associated. To maintain a centralized storage of secrets while using service managed secrets, the resources would also have to be centralized.

Though there are advantages to centralizing secrets for monitoring and observability, many customers already rely on services such as AWS Security Hub, IAM Access Analyzer, AWS Config, and Amazon CloudWatch for cross-account observability. These services make it easier to create centralized views of secrets in a multi-account environment.

Decentralized storage of secrets

In a decentralized approach to storage, shown in in Figure 4, secrets live in the same accounts as the workload that needs access to them.

Figure 4: Decentralized storage of secrets

Figure 4: Decentralized storage of secrets

The advantages of decentralizing the storage of secrets are:

  • Account boundaries and logical segmentation: Account boundaries provide a natural segmentation between workloads in AWS. When operating in a distributed multi-account environment, you cannot access secrets from another account by default, and all cross-account access must be allowed by both a resource policy in the source account and an IAM policy in the destination account. You can use resource control polices to prevent the sharing of secrets across accounts.
  • AWS KMS key choice: If your secrets aren’t shared across accounts, then you have the choice to use AWS KMS customer managed keys or AWS managed keys to encrypt your secrets.
  • Delegate permissions management to application owners: When secrets are stored in accounts with the applications that need to consume them, application owners define fine-grained permissions in secrets resource policies.

There are a few tradeoffs to consider for this architecture:

  • Auditing and monitoring require cross-account deployments: Tools that are used to monitor the compliance and status of secrets need to operate across multiple accounts and present information in a single place. This is simplified by AWS native tools, which are described later in this post.
  • Automated remediation workflows: You can have detective controls in place to alert on any misconfiguration or security risks related to your secrets. For example, you can surface an alert when a secret is shared outside of your organizational boundary through a resource policy. These workflows can be more complex in a multi-account environment. However, we have samples that can help, such as the Automated Security Response on AWS solution.

Centralized rotation

Like the creation and storage of secrets, organizations take different approaches to centralizing the lifecycle management and rotation of secrets.

When you centralize lifecycle management, as shown in Figure 5, a central team manages and owns AWS Lambda functions for rotation. The advantages of centralizing the lifecycle management of secrets are:

  • Developers can reuse rotation functions: In this pattern, a centralized team maintains a common library of rotation functions for different use cases. An example of this can be seen in this AWS re:Inforce session. Using this method, application teams don’t have to build their own custom rotation functions and can benefit from common architectural decisions regarding databases and third-party software as a service (SaaS) applications.
  • Logging: When storing and accessing rotation function logs, the centralized pattern can simplify managing logs from a single place.
Figure 5: Centralized rotation of secrets

Figure 5: Centralized rotation of secrets

There are some tradeoffs in centralizing the lifecycle management and rotation of secrets:

  • Additional cross-account access scenarios: When centralizing lifecycle management, the Lambda functions in central accounts require permissions to create, update, delete and read secrets in the application accounts. This increases the operational overhead required to enable secret rotation.
  • Service quotas: When you centralize a function at scale, service quotas can come into play. Check the Lambda service quotas to verify that you won’t hit quotas in your production environments.

Decentralized rotation

Decentralizing the lifecycle management of secrets is a more common choice, where the rotation functions live in the same account as the associated secret, as shown in Figure 6.

Figure 6: Decentralized rotation of secrets

Figure 6: Decentralized rotation of secrets

The advantages of decentralizing the lifecycle management of secrets are:

  • Templatization and customization: Developers can reuse rotation templates, but tweak the functions as needed to meet their needs and use cases
  • No cross-account access: Decentralized rotation of secrets happens all in one account and doesn’t require cross-account access.

The primary tradeoff of decentralizing rotation is that you will need to provide either centralized or federated access to logs for rotation functions in different accounts. By default, Lambda automatically captures logs for all function invocations and sends them to CloudWatch Logs. CloudWatch Logs offers a few different ways that you can centralize your logs, with the tradeoffs of each described in the documentation.

Centralized auditing and monitoring of secrets

Regardless of the model chosen for creation, storage, and rotation of secrets, centralize the compliance and auditing aspect when operating in a multi-account environment. You can use AWS Security Hub CSPM through its integration with AWS Organizations to centralize:

In this scenario, shown in Figure 7, centralized functions get visibility across the organization and individual teams can view their posture at an account level with no need to look at the state of the entire organization.

Use AWS CloudTrail organizational trails to send all API calls for Secrets Manager to a centralized delegated admin account.

Figure 7: Centralized monitoring and auditing

Figure 7: Centralized monitoring and auditing

Decentralized auditing and monitoring of secrets

For organizations that don’t require centralized auditing and monitoring of secrets, you can configure access so that individual teams can determine which logs are collected, alerts are enabled, and checks are in place in relation to your secrets. The advantages of this approach are:

  • Flexibility: Development teams have the freedom to choose what monitoring, auditing, and logging tools work best for them.
  • Reduced dependencies: Development teams don’t have to rely on centralized functions for alerting and monitoring capabilities.

The tradeoffs of this approach are:

  • Operational overhead: This can create redundant work for teams looking to accomplish similar goals.
  • Difficulty aggregating logs in cross-account investigations: If logs, alerts, and monitoring capabilities are decentralized, it can increase the difficulty of investigating events that affect multiple accounts.

Putting it all together

Most organizations choose a combination of these approaches to meet their needs. An example is a financial services company that has a central security team, operates across hundreds of AWS accounts, and has hundreds of applications that are isolated at the account level. This customer could:

  • Centralize the creation process, enforcing organizational standards for naming, tagging, and access control
  • Decentralize storage of secrets, using the AWS account as a natural boundary for access and storing the secret in the account where the workload is operating, delegating control to application owners
  • Decentralize lifecycle management so that application owners can manage their own rotation functions
  • Centralize auditing, using tools like AWS Config, Security Hub, and IAM Access Analyzer to give the central security team insight into the posture of their secrets while letting application owners retain control

Conclusion

In this post, we’ve examined the architectural decisions organizations face when implementing secrets management on AWS: creation, storage, rotation, and monitoring. Each approach—whether centralized or decentralized—offers distinct advantages and tradeoffs that should align with your organization’s security requirements, operational model, and scale. The important points include:

  • Choose your secrets management architecture based on your organization’s specific requirements and capabilities. There’s no one solution that will fit every situation.
  • Use automation and IaC to enforce consistent security controls, regardless of your approach.
  • Implement comprehensive monitoring and auditing capabilities through AWS services to maintain visibility across your environment.

Resources

To learn more about AWS Secrets Manager, check out some of these resources:

Brendan Paul
Brendan Paul

Brendan is a Senior Security Solutions Architect at AWS and has been at AWS for more than 6 years. He spends most of his time at work helping customers solve problems in the data protection and workload identity domains. Outside of work, he’s pursuing his master’s degree in data science from UC Berkeley.
Eduardo Patroncinio
Eduardo Patroncinio

Eduardo is a distinguished Principal Solutions Architect on the AWS Strategic Accounts team, bringing unparalleled expertise to the forefront of cloud technology. With an impressive career spanning more than 25 years, Eduardo has been a driving force in designing and delivering innovative customer solutions within the dynamic realms of cloud and service management.

Modernize your data warehouse by migrating Oracle Database to Amazon Redshift with Oracle GoldenGate

Post Syndicated from Sachin Murkar original https://aws.amazon.com/blogs/big-data/modernize-your-data-warehouse-by-migrating-oracle-database-to-amazon-redshift-with-oracle-goldengate/

In this post, we show how to migrate an Oracle data warehouse to Amazon Redshift using Oracle GoldenGate and DMS Schema Conversion, a feature of AWS Database Migration Service (AWS DMS). This approach facilitates minimal business disruption through continuous replication. Amazon Redshift is a fast, fully managed, petabyte-scale data warehouse service that makes it simple and cost-effective to efficiently analyze your data using your existing business intelligence tools.

Solution overview

Our migration approach combines DMS Schema Conversion for schema migration and Oracle GoldenGate for data replication. The migration process consists of four main steps:

  1. Schema conversion using DMS Schema Conversion.
  2. Initial data load using Oracle GoldenGate.
  3. Change data capture (CDC) for ongoing replication.
  4. Final cutover to Amazon Redshift.

The following diagram shows the migration workflow architecture from Oracle to Amazon Redshift, where DMS Schema Conversion handles schema migration and Oracle GoldenGate manages both initial data load and continuous replication through Extract and Replicat processes running on Amazon Elastic Compute Cloud (Amazon EC2) instances. The solution facilitates minimal downtime by maintaining real-time data synchronization until the final cutover.

The solution comprises the following key migration components:

In the following sections, we walk through how to migrate an Oracle data warehouse to Amazon Redshift. For demonstration purposes, we use an Oracle data warehouse consisting of four tables:

dim_customer
dim_product
dim_date
fact_sales

Prerequisites

We recommend reviewing the licensing requirements for Oracle GoldenGate. For more information, refer to Oracle GoldenGate Licensing Information.

Run schema conversion using DMS Schema Conversion

DMS Schema Conversion automatically converts your Oracle database schemas and code objects to Amazon Redshift-compatible formats. This includes tables, views, stored procedures, functions, and data types.

Set up network for DMS Schema Conversion

DMS Schema Conversion requires network connectivity to both your source and target databases. To set up this connectivity, complete the following steps:

  1. Specify a virtual private cloud (VPC) and subnet where DMS Schema Conversion will run.
  2. Configure security group rules to allow traffic between the following:
    1. DMS Schema Conversion and your source Oracle database
    2. DMS Schema Conversion and your target Redshift cluster
  3. For on-premises databases, set up either:
    1. AWS Site-to-Site VPN
    2. AWS Direct Connect

For comprehensive information about network configurations, refer to Setting up a network for DMS Schema Conversion.

Store database credentials in AWS Secrets Manager

DMS Schema Conversion uses secrets stored in AWS Secrets Manager to connect to your database. For instructions to add source and target credentials to Secrets Manager, refer to Store database credentials in AWS Secrets Manager.

Create S3 bucket

DMS Schema Conversion saves items such as assessment reports, converted SQL code, and information about database schema objects in an S3 bucket. For instructions to create an S3 bucket, refer to Create an S3 bucket.

Create IAM policies and roles

To set up DMS Schema Conversion, you must create appropriate IAM policies and roles. This process makes sure AWS DMS has the necessary permissions to access your source and target databases, as well as other AWS services required for the migration.

Prepare DMS Schema Conversion

In this section, we go through the steps to configure DMS Schema Conversion.

Set up instance profile

An instance profile specifies the network, security, and Amazon S3 settings for DMS Schema Conversion to use. Create an instance profile with the following steps:

  1. On the AWS DMS console, choose Instance profiles in the navigation pane.
  2. Choose Create instance profile.
  3. For Name, enter a name (for example, sc-instance).
  4. For Network type, we use IPv4. DMS Schema Conversion also offers Dual-stack mode for both IPv4 and IPv6.
  5. For Virtual private cloud (VPC) for IPv4, choose Default VPC.
  6. For Subnet group, choose your subnet group (for this post, default).
  7. For VPC security groups, choose your security groups. As previously stated, the instance profile’s VPC security group must have access to both the source and target databases.
  8. For S3 bucket, specify a bucket to store schema conversion metadata.
  9. Choose Create instance profile.

Add data providers

Data providers store database types and information about source and target databases for DMS Schema Conversion to connect to. Configure data providers for the source and target databases with the following steps:

  1. On the AWS DMS console, choose Data providers in the navigation pane.
  2. Choose Create data provider.
  3. To create your target, for Name, enter a name (for example, redshift-target).
  4. For Engine type, choose Amazon Redshift.
  5. For Engine configuration, select Choose from Redshift.
  6. For Redshift cluster, choose the target Redshift cluster.
  7. For Port, enter the port number.
  8. For Database name, enter the name of your database.
  9. Choose Create data provider.
  10. Repeat similar steps to create your source data provider.

Create migration project

The DMS Schema Conversion migration project defines migration entities, including instance profiles, source and target data providers, and migration rules. Create a migration project with the following steps:

  1. On the AWS DMS console, choose Migration projects in the navigation pane.
  2. Choose Create migration project.
  3. For Name, enter a name to identify your migration project (for example, oracle-redshift-commercewh).
  4. For Instance profile, choose the instance profile you created.

  1. In the Data providers section, enter the source and target data providers, Secrets Manager secret, and IAM roles.

  1. In the Schema conversion settings section, enter the S3 URL and choose the applicable IAM role.

  1. Choose Create migration project.

Use DMS Schema Conversion to transform Oracle database objects

Complete the following steps to convert source database objects:

  1. On the AWS DMS console, choose Migration projects in the navigation pane.
  2. Choose the migration project you created.
  3. On the Schema conversion tab, choose Launch schema conversion.

The schema conversion project will be ready when the launch is complete. The left navigation tree represents the source database, and the right navigation tree represents the target database.

  1. Generate and view the assessment report.
  2. Select the objects you want to convert and then choose Convert on the Actions menu to convert the source objects to the target database.

The conversion process might take some time depending on the number and complexity of the selected objects.

You can save the converted code to the S3 bucket that you created earlier in the prerequisite steps.

  1. To save the SQL scripts, select the object in the target database tree and choose Save as SQL on the Actions menu.
  2. After you finalize the scripts, run them manually in the target database.
  3. Alternatively, you can apply the scripts directly to the database using DMS Schema Conversion. Select the specific schema in the target database, and on the Actions menu, choose Apply changes.

This will apply the automatically converted code to the target database.

If some objects require action items, DMS Schema conversion flags them and provides details of action items. For the items that require resolution, perform manual changes and apply the converted changes directly to the target database.

Perform data migration

The migration from Oracle Database to Amazon Redshift using Oracle GoldenGate begins with an initial load process, where Oracle GoldenGate’s Extract process captures the existing data from the Oracle source tables and sends this data to the Replicat process, which loads it into Redshift target tables through the appropriate database connectivity. Simultaneously, Oracle GoldenGate’s CDC mechanism tracks the ongoing changes (inserts, updates, and deletes) in the source Oracle database by reading the redo logs. These captured changes are then synchronized to Amazon Redshift in near real time through the Extract-Pump-Replicat process, facilitating data consistency between the source and target systems throughout the migration process.

Prepare source Oracle database for GoldenGate

Prepare your database for Oracle GoldenGate, including configuring connections and logging, enabling Oracle GoldenGate in your database, setting up the flashback query, and managing server resources.

Oracle GoldenGate for BigData only supports uncompressed UPDATE records when replicating to Amazon Redshift. When UPDATE records contain missing columns, those columns are set to null in the target.

To handle this situation, configure Extract to generate trail records with the column values (enable trandata for the columns). Alternatively, you can disable this check by setting gg.abend.on.missing.columns=false, which may result in unintended NULLs on the target database.When gg.abend.on.missing.columns=true, Replicat process on Oracle GoldenGate for BigData fails and returns the following error for compressed update records:

ERROR OGG-15051 Java or JNI exception: java.lang.IllegalStateException: The UPDATE operation record in the trail at pos[0/XXXXXXX] for table [SCHEMA.TABLENAME] has missing columns.

Install Oracle GoldenGate software on Amazon EC2

You must run Oracle GoldenGate on EC2 instances. The instances must have adequate CPU, memory, and storage to handle the anticipated replication volume. For more details, refer to Operating System Requirements. After you determine the CPU and memory requirements, select a current generation EC2 instance type for Oracle GoldenGate.

When the EC2 instance is up and running, download the following Oracle GoldenGate software from the Oracle GoldenGate Downloads page:

  • Oracle GoldenGate for Oracle 21.3.0.0
  • Oracle GoldenGate for Big Data 21c

For installation, refer to Install, Patch, and Upgrade and Installing and Upgrading Oracle GoldenGate for Big Data.

Configure Oracle GoldenGate for initial load

The initial load configuration transfers existing data from Oracle Database to Amazon Redshift. Complete the following configuration steps:

  1. Create an initial load extract parameter file for the source Oracle database using GoldenGate for Oracle. The following code is the sample file content:
    # Extract initial load configuration (INITLE11)
    
    EXTRACT INITLE11
    SETENV ORACLE_HOME=/u01/app/oracle/product/19.3.0/dbhome_1
    USERID ******************:1521/ORCL, PASSWORD ogg_password
    RMTHOST ec2-xx-xx-xx-xx.compute-1.amazonaws.com, MGRPORT 9809, COMPRESS
    RMTTASK REPLICAT, GROUP INITLR11
    TABLE commerce_wh.dim_customer;
    TABLE commerce_wh.dim_product;
    TABLE commerce_wh.dim_date;
    TABLE commerce_wh.fact_sales;

  2. Add the EXTRACT on the GoldenGate for Oracle prompt by running the following command:
    ADD EXTRACT INITLE11, SOURCEISTABLE
    
    GGSCI (ip-**-**-**-**.us-west-2.compute.internal) 1> info INITLE11
    
    Extract    INITLE11  Initialized  2025-07-08 03:44   Status STOPPED
    Checkpoint Lag       Not Available
    Log Read Checkpoint  Not Available
                         First Record         Record 0
    Task                 SOURCEISTABLE

  3. Create a Replicat parameter file for the target Redshift database for the initial load using GoldenGate for Big Data. The following code is the sample file content:
    # Replicate initial load configuration (INITLR11)
    
    REPLICAT INITLR11
    TARGETDB LIBFILE libggjava.so SET property=/home/ec2-user/ogg_bd/dirprm/rs.props
    MAP commerce_wh.dim_customer, TARGET commerce_wh.dim_customer;
    MAP commerce_wh.dim_product, TARGET commerce_wh.dim_product;
    MAP commerce_wh.dim_date, TARGET commerce_wh.dim_date;
    MAP commerce_wh.fact_sales, TARGET commerce_wh.fact_sales;
    ```

  4. Add the REPLICAT on the GoldenGate for Big Data prompt by running the following command:
    ADD REPLICAT INITLR11, SPECIALRUN
    
    GGSCI (ip-**-**-**-**.us-west-2.compute.internal) 2> info INITLR11
    
    Replicat   INITLR11  Initialized  2025-07-08 03:47   Status STOPPED
    Checkpoint Lag       00:00:00 (updated 00:00:05 ago)
    Log Read Checkpoint  Not Available
    Task                 SPECIALRUN

Configure Oracle GoldenGate for CDC and Amazon Redshift handler

In this section, we walk through the steps to configure Oracle GoldenGate for CDC and the Amazon Redshift handler.

Configure Oracle GoldenGate for extracting from source

For continuous replication, set up the Extract, Pump, and Replicat processes:

  1. Create an Extract parameter file for the source Oracle database for CDC using GoldenGate for Oracle. The following code is the sample file content:
    # Extract configuration (EXTPRD)
    
    EXTRACT EXTPRD
    SETENV ORACLE_HOME=/u01/app/oracle/product/19.3.0/dbhome_1
    USERID ********@oracledb:1521/ORCL, PASSWORD ogg_password
    *************************************************/dirdat/ep
    CHECKPOINTSECS 1
    TABLE commerce_wh.dim_customer;
    TABLE commerce_wh.dim_product;
    TABLE commerce_wh.dim_date;
    TABLE commerce_wh.fact_sales;
    TRANLOGOPTIONS ALTARCHIVELOGDEST /u01/app/oracle/fast_recovery_area/ORCL/archivelog

  2. Add the Extract process and register it:
    # Add Extract and Register (EXTPRD)
    
    ADD EXTRACT EXTPRD, INTEGRATED TRANLOG, BEGIN NOW
    
    REGISTER EXTRACT EXTPRD DATABASE
    
    ADD EXTTRAIL /u01/app/oracle/product/21.3.0/oggcore_1/dirdat/ep, EXTRACT 
    EXTPRD
    
    GGSCI (ip-**-**-**-**.us-west-2.compute.internal) 3>  info EXTPRD
    
    Extract    EXTPRD    Initialized  2025-07-08 03:50   Status STOPPED
    Checkpoint Lag       00:00:00 (updated 00:00:36 ago)
    Log Read Checkpoint  Oracle Integrated Redo Logs
                         2025-07-08 03:50:33

  3. Create an Extract Pump parameter file for the source Oracle database to send the trail files to the target Redshift database. The following code is the sample file content:
    # Pump process configuration (PMPPRD)
    
    EXTRACT PMPPRD
    PASSTHRU
    RMTHOST ec2-xx-xx-xx-xx.compute-1.amazonaws.com, MGRPORT 9809, COMPRESS
    RMTTRAIL /home/********/ogg_bd/dirdat/pt
    TABLE commerce_wh.dim_customer;
    TABLE commerce_wh.dim_product;
    TABLE commerce_wh.dim_date;
    TABLE commerce_wh.fact_sales;

  4. Add the Pump process:
    # Pump process addition
    
    ADD EXTRACT PMPPRD, EXTTRAILSOURCE /u01/app/oracle/product/21.3.0/oggcore_1/dirdat/ep
    
    ADD RMTTRAIL /home/ec2-user/ogg_bd/dirdat/pt, EXTRACT PMPPRD
    
    GGSCI (ip-**-**-**-**.us-west-2.compute.internal) 4> info PMPPRD
    
    Extract    PMPPRD    Initialized  2025-07-08 03:51   Status STOPPED
    Checkpoint Lag       00:00:00 (updated 00:00:09 ago)
    Log Read Checkpoint  File /u01/app/oracle/product/21.3.0/oggcore_1/dirdat/ep000000000
                         First Record  RBA 0

Configure Oracle GoldenGate Redshift handler to apply changes to target

To configure an Oracle GoldenGate Replicat to send data to a Redshift cluster, you must set up a Redshift properties file and a Replicat parameter file that defines how data is migrated to Amazon Redshift. Complete the following steps:

  1. Configure the Replicat properties file (rs.props), which consists of an S3 event handler and Redshift event handler. The following is an example Replicat properties file configured to connect to Amazon Redshift:
    gg.target=redshift
    
    # S3 Event Handler
    gg.eventhandler.s3.region=us-west-2
    gg.eventhandler.s3.bucketMappingTemplate=your-s3-bucket-name
    
    # Redshift Event Handler
    gg.eventhandler.redshift.connectionURL=jdbc:redshift://your-cluster.region.redshift.amazonaws.com:5439/dev
    gg.eventhandler.redshift.userName=your_redshift_username
    gg.eventhandler.redshift.Password=your_redshift_password
    gg.classpath=/path/to/aws-sdk-java/*:/path/to/redshift-jdbc-driver.jar
    jvm.bootoptions=-Xmx8g -Xms8g
    
    gg.eventhandler.redshift.AwsIamRole=arn:aws:iam::your-account-id:role/your-redshift-role
    
    gg.abend.on.missing.columns=false

    To authenticate Oracle GoldenGate’s access to the Redshift cluster for data load operations, you have two options. The recommended and more secure method is to use IAM role authentication by configuring the gg.eventhandler.redshift.AwsIamRole property in the properties file. This approach provides more secure, role-based access. Alternatively, you can use access key authentication by setting the environment variables AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY. For more information, refer to the Oracle GoldenGate for BigData documentation.

  2. Create a Replicat parameter file for the target Redshift database using Oracle GoldenGate for BigData. The following code is the sample file content:
    # Replicat process configuration (RSPRD)
    
    REPLICAT RSPRD
    TARGETDB LIBFILE libggjava.so SET property=/home/********/ogg_bd/dirprm/rs.props
    REPORTCOUNT EVERY 1 MINUTES, RATE
    GROUPTRANSOPS 1000
    MAP commerce_wh.dim_customer, TARGET commerce_wh.dim_customer;
    MAP commerce_wh.dim_product, TARGET commerce_wh.dim_product;
    MAP commerce_wh.dim_date, TARGET commerce_wh.dim_date;
    MAP commerce_wh.fact_sales, TARGET commerce_wh.fact_sales;

  3. Add a Replicat process:
    # Add Replicat
    ADD REPLICAT RSPRD, EXTTRAIL /home/ec2-user/ogg_bd/dirdat/pt, BEGIN NOW
    
    GGSCI (ip-**-**-**-**.us-west-2.compute.internal) 3> info RSPRD
    
    Replicat   RSPRD     Initialized  2025-07-08 03:52   Status STOPPED
    Checkpoint Lag       00:00:00 (updated 00:00:07 ago)
    Log Read Checkpoint  File /home/ec2-user/ogg_bd/dirdat/pt000000000
                         2025-07-08 03:52:48.471461

Start initial load and change sync

First start the change sync extract and data pump on the source Oracle database. This will start capturing changes while you perform the initial load.

  1. In the GoldenGate for Oracle GGSCI utility, start EXTPRD and PMPPRD:
    GGSCI (ip-**-**-**-**.us-west-2.compute.internal as ggsuser@ORCL) 13> start EXTPRD
    
    Sending START request to Manager ...
    Extract group EXTPRD starting.
    
    
    GGSCI (ip-**-**-**-**.us-west-2.compute.internal as ggsuser@ORCL) 15> start PMPPRD
    
    Sending START request to Manager ...
    Extract group PMPPRD starting.

    Do not start Replicat at this point.

  2. Record the Source System Change Number (SCN) from the Oracle database, which serves as the starting point for replication on the target system:
    select current_scn from v$database;
    
    CURRENT_SCN
    13940177

  3. Start the initial load Extract process, which will automatically trigger the corresponding initial load Replicat on the target system:
    GGSCI (ip-**-**-**-**.us-west-2.compute.internal as ggsuser@ORCL) 21> start INITLE11
    
    Sending START request to Manager ...
    Extract group INITLE11 starting.

  4. Monitor the initial load completion status by executing the following command on the GoldenGate for BigData GGSCI utility. Make sure the initial load process has completed successfully before proceeding to the next step. The report will indicate the load status and potential errors that need attention.
    VIEW REPORT INITLR11

  5. Start the change synchronization Replicat RSPRD using the previously captured SCN to facilitate continuous data replication:
    GGSCI (ip-**-**-**-**.us-west-2.compute.internal) 17> start RSPRD , aftercsn 13940177
    
    Sending START request to Manager ...
    Replicat group RSPRD starting.

Refer to the Oracle GoldenGate documentation for Amazon Redshift handlers to learn more about its detailed functionality, unsupported operations, and system limitations.

When transitioning from initial load to continuous replication in an Oracle database to Amazon Redshift migration using Oracle GoldenGate, it’s crucial to properly manage data collisions to maintain data integrity. The key is to capture and use an appropriate SCN that marks the exact point where initial load ends and CDC begins. Without proper collision handling, you might encounter duplicate records or missing data during the transition period. Implementing appropriate collision handling mechanisms makes sure duplicate records are properly managed without causing data inconsistencies in the target system. For more information on HANDLECOLLISIONS, refer to the Oracle GoldenGate documentation.

Clean up

When the migration is complete, complete the following steps:

  1. Stop and remove Oracle GoldenGate processes (EXTRACT, PUMP, REPLICAT).
  2. Delete EC2 instances used for Oracle GoldenGate.
  3. Remove IAM roles created for migration.
  4. Delete S3 buckets used for DMS Schema Conversion (if no longer needed).
  5. Update application connection strings to point to the new Redshift cluster.

Conclusion

In this post, we showed how to modernize your data warehouse by migrating to Amazon Redshift using Oracle GoldenGate. This approach facilitates minimal downtime and provides a flexible, reliable method for transitioning your critical data workloads to the cloud. With the complexity involved in database migrations, we highly recommend testing the migration steps in non-production environments prior to making changes in production. By following the best practices outlined in this post, you can achieve a smooth migration process and set the foundation for a scalable, cost-effective data warehousing solution on AWS. Remember to continuously monitor your new Amazon Redshift environment, optimize query performance, and take advantage of the AWS suite of analytics tools to derive maximum value from your modernized data warehouse.


About the authors

Sachin Murkar

Sachin Murkar

Sachin is a Cloud Support Database Engineer at AWS. He is a Subject Matter Expert in RDS PostgreSQL and Aurora PostgreSQL. Based in the Pacific Northwest region, Sachin focuses on helping customers optimize their AWS database solutions, with particular expertise in Amazon RDS and Aurora.

Ravi Teja Bellamkonda

Ravi Teja Bellamkonda

Ravi is a Technical Account Manager (TAM) at AWS and a Subject Matter Expert (SME) for AWS DMS. With nearly 10 years of experience in database technologies, specializing in PostgreSQL and Oracle, he helps customers design and execute seamless database migration strategies to the cloud.

Bipin Nair

Bipin Nair

Bipin is a Cloud Support Database Engineer at AWS and Subject Matter Expert for AWS DMS and Amazon RDS for PostgreSQL. He has over a decade of experience in working with Oracle databases, Replication Services and AWS relational databases.

On-demand and scheduled scaling of Amazon MSK Express based clusters

Post Syndicated from Subham Rakshit original https://aws.amazon.com/blogs/big-data/on-demand-and-scheduled-scaling-of-amazon-msk-express-based-clusters/

Modern streaming workloads are highly dynamic—traffic volumes fluctuate based on time of day, business cycles, or event-driven bursts. Customers need to dynamically scale Apache Kafka clusters up and down to maintain consistent throughput and performance without incurring unnecessary cost. For example, ecommerce platforms see sharp traffic increases during seasonal sales, and financial systems experience load spikes during market hours. Scaling clusters helps teams align cluster capacity with increased ingress throughput in response to these variations, leading to more efficient utilization and a better cost-to-performance ratio.

Amazon Managed Streaming for Apache Kafka (Amazon MSK) Express brokers are a key component to dynamically scaling clusters to meet demand. Express based clusters deliver 3 times higher throughput, 20 times faster scaling capabilities, and 90% faster broker recovery compared to Amazon MSK Provisioned clusters. In addition, Express brokers support intelligent rebalancing for 180 times faster operation performance, so partitions are automatically and consistently well distributed across brokers. This feature is enabled by default for all new Express based clusters and comes at no additional cost to customers. This capability alleviates the need for manual partition management when modifying cluster capacity. Intelligent rebalancing automatically tracks cluster health and triggers partition redistribution when resource imbalances are detected, maintaining performance across brokers.

This post demonstrates how to use the intelligent rebalancing feature and build a custom solution that scales Express based clusters horizontally (adding and removing brokers) dynamically based on Amazon CloudWatch metrics and predefined schedules. The solution provides capacity management while maintaining cluster performance and minimizing overhead.

Overview of Kafka scaling

Scaling Kafka clusters involves adding or removing brokers to the cluster while providing balanced data distribution and uninterrupted service. When new brokers are added, partition reassignment is required to evenly distribute load across the cluster. This process is typically performed manually—either through the Kafka command line tools (kafka-reassign-partitions.sh) or by using automation frameworks such as Cruise Control, which intelligently calculates and executes reassignment plans. During scale-in operations, partitions hosted on the brokers marked for removal must first be migrated to other brokers, leaving the target brokers empty before decommissioning.

Challenges of scaling Kafka dynamically

The complexity of scaling depends heavily on the underlying storage model. In deployments where broker data resides entirely on local storage, scaling involves physical data movement between brokers, which can take considerable time depending on partition size and replication factor. In contrast, environments that use tiered storage shift most of the data to remote object storage such as Amazon Simple Storage Service (Amazon S3), making scaling a largely metadata-driven operation. This significantly reduces data transfer overhead and accelerates both broker addition and removal, enabling more elastic and operationally efficient Kafka clusters.

However, scaling Kafka remains a non-trivial operation due to the interplay between storage, data movement, and broker resource utilization. When partitions are reassigned across brokers, large volumes of data must be copied over the network, often leading to network bandwidth saturation, storage bandwidth exhaustion, and elevated CPU utilization. Depending on data volume and replication factor, partition rebalancing can take several hours, during which time cluster performance and throughput might temporarily degrade and often require additional configuration to throttle the data movement. Although tools like Cruise Control automate this process, they introduce another layer of complexity: selecting the right combination of rebalancing goals (such as disk capacity, network load, or replica distribution) requires a deep understanding of Kafka internals and trade-offs between speed, balance, and stability. As a result, efficient scaling is an optimization problem, demanding careful orchestration of storage, compute, and network resources.

How Express brokers simplify scaling

Express brokers manage Kafka scaling through their decoupled compute and storage architecture. This innovative design enables unlimited storage without pre-provisioning, significantly simplifying cluster sizing and management. The separation of compute and storage resources allows Express brokers to scale faster than standard MSK brokers, enabling rapid cluster expansion within minutes. With Express brokers, administrators can adjust capacity both vertically and horizontally as needed, alleviating the need for over-provisioning. The architecture provides sustained broker throughput during scaling operations, with Express brokers capable of handling 500 MBps ingress and 1000 MBps egress on m7g.16xl instances. For more information about how the scaling process works in Express based clusters, see Express brokers for Amazon MSK: Turbo-charged Kafka scaling with up to 20 times faster performance.

Added to this faster scaling capability, when you add or remove brokers from your Express based clusters, intelligent rebalancing automatically redistributes partitions to balance resource utilization across the brokers. This makes sure the cluster continues to operate at peak performance, making scaling in and out possible with a single update operation. Intelligent rebalancing is enabled by default on new Express broker clusters and continuously monitors cluster health for resource imbalances or hotspots. For example, if certain brokers become overloaded due to uneven distribution of partitions or skewed traffic patterns, intelligent rebalancing will automatically move partitions to less utilized brokers to restore balance.

Finally, Express based clusters automate client configuration of broker bootstrap connection strings to allow clients to connect to clusters seamlessly as brokers are added and removed. Express based clusters provide three connection strings, one per Availability Zone, which are independent of the brokers in the cluster. This means clients only need to configure these connection strings to maintain consistent connections as brokers are added or removed. These key capabilities of Express based clusters—rapid scaling, intelligent rebalancing, and dynamic broker bootstrapping—are critical to enabling dynamic scaling in Kafka clusters. In the following section, we explore how we use these capabilities to automate the scaling process of Express based clusters.

On-demand and scheduled scaling

Leveraging fast scaling capabilities of Express brokers together with intelligent rebalancing, you can build a flexible and dynamic scaling solution to optimize your Kafka cluster resources. There are two primary approaches for automatic scaling that balance performance needs with cost efficiency: on-demand and scheduled scaling.

On-demand scaling

On-demand scaling tracks cluster performance and responds to capacity demands. This approach addresses scenarios where workload patterns experience traffic spikes. On-demand scaling tracks Amazon MSK performance indicators as CPU utilization and network ingress and egress throughput per broker. Beyond these infrastructure metrics, the solution also supports using CloudWatch metrics to enable business-logic-driven scaling decisions.

The solution evaluates the performance metrics continuously against configurable thresholds to determine when scaling actions are necessary. When brokers operate above capacity thresholds consistently over a period of time, it invokes an Amazon MSK API to increase the broker count of the cluster. The solution in this post currently supports horizontal scaling (adding and removing brokers) only. Intelligent rebalancing will then automatically redistribute the partitions to spread the load across the new brokers that are added. Similarly, when utilization drops below thresholds, the solution invokes an Amazon MSK API to remove brokers. The rebalancing process automatically moves partitions from the broker marked for removal to other brokers in the cluster. This solution requires topics to have sufficient partitions to support rebalancing to new brokers as brokers are added.

The following diagram illustrates the on-demand scaling workflow.

This diagram illustrates the automated scaling and rebalancing workflow for Amazon Managed Streaming for Apache Kafka (MSK). The process consists of four sequential stages that ensure optimal cluster performance through intelligent monitoring and automated actions.

Scheduled scaling

Scheduled scaling adjusts cluster capacity using time-based triggers. This approach is useful for applications with traffic patterns that correlate with business hours or schedules. For example, ecommerce platforms benefit from scheduled scaling during peak sale periods when customer activity peaks. Scheduled scaling is also useful for customers who want to avoid cluster modification operations during business hours. This solution uses a configurable schedule to scale out the cluster capacity before business hours to handle the anticipated traffic and scale in after business hours to reduce costs. This particular solution currently supports horizontal scaling (adding/removing brokers) only. With scheduled scaling, you can handle specific scenarios such as weekday business hours, weekend maintenance windows, or specific dates. You can also specify the desired number of brokers at scale-out and scale-in.

The following diagram illustrates the scheduled scaling workflow.

This horizontal process flow diagram illustrates the automated scaling and rebalancing workflow for Amazon Managed Streaming for Apache Kafka (MSK). The diagram demonstrates how MSK clusters continuously monitor performance, evaluate scaling requirements, execute scaling operations, and automatically rebalance partitions to maintain optimal performance without manual intervention.

Solution overview

This solution provides scaling automation for Express brokers through two approaches:

  • On-demand scaling – Tracks built-in cluster performance metrics or custom CloudWatch metrics and adjusts broker capacity when thresholds are crossed
  • Scheduled scaling – Scales clusters based on specific schedules

In the following sections, we provide the implementation details for both scaling methods.

Prerequisites

Complete the following steps as prerequisites:

  1. Create an Express cluster with intelligent rebalancing enabled. The intelligent rebalancing feature is required for this solution to work. Note the Amazon Resource Name (ARN) of the cluster.
  2. Install Python 3.11 or higher on Amazon Elastic Compute Cloud (Amazon EC2).
  3. Install the AWS Command Line Interface (AWS CLI) and configure it with your AWS credentials.
  4. Install the AWS CDK CLI.

On-demand scaling solution

The solution uses an AWS Lambda function that is triggered by an Amazon EventBridge scheduler periodically. The Lambda function checks the cluster state and time since the last broker addition or removal was done. This is done to determine if the cluster is ready to scale. If the cluster is ready for scaling, the function collects the CloudWatch metrics that need to be evaluated to make the scaling decision. Based on the scaling configuration and using the metrics in CloudWatch, the function evaluates the scaling logic and executes the scaling decision. The scaling decision can lead to addition or removal of brokers to the cluster. In both cases, intelligent rebalancing handles partition distribution across brokers without manual intervention. You can find more details of the scaling logic in the GitHub repo.

The following diagram illustrates the architecture of the on-demand scaling solution.

This AWS architecture diagram illustrates a serverless event-driven workflow that uses Amazon EventBridge Scheduler to trigger AWS Lambda functions that interact with Amazon MSK Express brokers, with monitoring provided by Amazon CloudWatch Metrics. The diagram demonstrates a fully managed, scalable architecture for time-based or event-based Apache Kafka operations.

Deploy on-demand scaling solution

Follow these steps to deploy the on-demand scaling infrastructure. For this post, we demonstrate the on-demand scale-out functionality.

  1. Run the following commands to set the project up:
    git clone https://github.com/aws-samples/sample-msk-express-brokers-scaling.git
    cd sample-msk-express-brokers-scaling/scaling/cdk
    python -m venv .venv && source .venv/bin/activate
    pip install -r requirements.txt

  2. Modify the thresholds to match your MSK broker instance size and business requirements by editing src/config/on_demand_scaling_config.json. Refer to the configuration documentation for more details of the configuration options available.
    By default, on_demand_scaling_config.json considers the express.m7g.large broker instance size. Therefore the scale-in/scale-out ingress/egress thresholds are configured at 70% of the recommended sustained throughput for the instance size.
  3. Bootstrap your environment for use with the AWS CDK.
  4. Deploy the on-demand scaling AWS CDK application:
    cdk deploy MSKOnDemandScalingStack \
      --app "python3 msk_on_demand_scaling_stack.py" \
      --context cluster_arn="<< ARN of the MSK Cluster >>" \
      --context monitoring_frequency_minutes=1 \
      --context stack_name="MSKOnDemandScalingStack"

The monitoring_frequency_minutes parameter controls how often the EventBridge scheduler invokes the scaling logic Lambda function to evaluate cluster metrics.

The deployment creates the AWS resources required to run the on-demand scaling solution. The details of the resources created are shown in the output of the command.

Test and monitor the on-demand scaling solution

Configure the bootstrap server for your MSK cluster. You can get the bootstrap server from the AWS Management console or using the AWS CLI.

export BOOTSTRAP=<<BOOTSTRAP_SERVER>>

Create a Kafka topic in the cluster. Update the following command for the specific authentication method in Amazon MSK. Refer to the Amazon MSK Labs workshop for more details.

Topics should have a sufficient number of partitions that can be distributed across a larger set of brokers.

export TOPIC_NAME=<<TOPIC_NAME>>

bin/kafka-topics.sh \
--bootstrap-server=$BOOTSTRAP \
--create \
--replication-factor 3 \
--partitions 96 \
--topic $TOPIC_NAME

Generate load on the MSK cluster to trigger and verify the scaling operations. You can use an existing application that drives load to your cluster. You can also use the kafka-producer-perf-test.sh utility that is bundled as part of the Kafka distribution to generate load:

bin/kafka-producer-perf-test.sh \
  --topic $TOPIC_NAME \
  --num-records 1000000000 \
  --record-size 1024 \
  --throughput -1 \
  --producer-props bootstrap.servers=$BOOTSTRAP

Monitor the scaling operations by tailing the Lambda function logs:

aws logs tail /aws/lambda/MSKOnDemandScalingStack-MSKScalingFunction  \
--follow --format short

In the logs, look for the following messages to identify the exact times when scaling operations occurred. The log statements above these messages show the rationale behind the scaling decision:

[INFO] Calling MSK UpdateBrokerCount API...
 [INFO] Successfully initiated broker count update operation

The solution also creates a CloudWatch dashboard that provides visibility into scaling operations and many other broker metrics. The link to the dashboard is shown in the output of the cdk deploy command.

The following figure shows a cluster that started with three brokers. After the 09:15 mark, it received consistent inbound traffic, which exceeded the thresholds set in the solution. The solution added three more brokers that came into service at around the 09:45 mark. Intelligent rebalancing reassigned some of the partitions to the newly added brokers and the incoming traffic was split across six brokers. The solution continued adding more brokers until the cluster had 12 brokers and the intelligent rebalancing feature continued distributing the partitions across the newly added brokers.

Amazon MSK Broker Network Throughput Performance Chart: Bytes In Per Second Maximum by Broker This time-series line chart visualizes the maximum inbound network throughput performance across 25 individual Apache Kafka brokers in an Amazon Managed Streaming for Apache Kafka (MSK) cluster over a 3-hour time period from 09:00 to 11:45. The chart demonstrates broker-level network ingestion rates, scaling operations, and performance variations during active workload processing.

The following figure shows the times when partition rebalancing was active (value=1). In the context of this solution, that typically occurs after new brokers are added or removed and the scaling operations are complete.

Amazon MSK Intelligent Rebalancing Status Timeline Chart This binary state timeline chart visualizes the activation and deactivation cycles of Amazon Managed Streaming for Apache Kafka (MSK) Intelligent Rebalancing feature over a 2 hour and 45 minute observation period from 09:00 to 11:45. The chart displays discrete on/off status indicators showing when the automated partition rebalancing feature was actively running versus inactive.

The following figure shows the number of brokers added (positive values) or removed (negative values) from the cluster. This helps visualize and track the size of the cluster as it goes through scaling operations.

Amazon MSK Broker Count Change Timeline Chart This time-series chart visualizes broker count changes in an Amazon Managed Streaming for Apache Kafka (MSK) cluster over a 2 hour and 45 minute period from 09:00 to 11:45 UTC on November 12, 2025. The chart tracks incremental additions and removals of Kafka brokers, demonstrating MSK's dynamic scaling capabilities in response to workload demands.

Scheduled scaling solution

The scheduled scaling implementation supports timing patterns through an EventBridge schedule. You can configure timing to trigger an action using cron expressions. Based on the cron expression, the EventBridge Scheduler triggers a Lambda function at the specified time to scale out or scale in. The Lambda function performs checks if the cluster is ready for a scaling operation and performs the requested scaling operation by invoking the Amazon MSK control plane API. The service allows removing only three brokers at a time from a cluster. The solution handles this scenario by repeatedly removing the brokers in counts of three until the desired number of brokers are reached.

The following diagram illustrates the architecture of the scheduled scaling solution.

This AWS architecture diagram illustrates an event-driven, time-based auto-scaling workflow where two Amazon EventBridge Scheduler instances trigger an AWS Lambda function to execute scale-up and scale-down operations on an Amazon MSK Express broker. The diagram demonstrates serverless capacity management for Apache Kafka infrastructure using scheduled automation.

Configuration parameters

EventBridge schedules support cron expressions for precise timing control, so you can fine-tune scaling operations for specific times of day and days of the week. For example, you can configure scaling to occur at 8:00 AM on weekdays using the cron expression cron(0 8 ? * MON-FRI *). To scale in at 6:00 PM on the same days, use cron(0 18 ? * MON-FRI *). For more patterns, refer to Setting a schedule pattern for scheduled rules (legacy) in Amazon EventBridge. You can also configure the desired broker count to be reached during scale-out and scale-in operations.

Deploy scheduled scaling solution

Follow these steps to deploy the scheduled scaling solution:

  1. Run the following commands to set the project up:
    cd scaling/cdk
    python3 -m venv .venv && source .venv/bin/activate
    pip install -r requirements.txt

  2. Modify the scaling schedule by editing scaling/cdk/src/config/scheduled_scaling_config.json. Refer to the configuration documentation for more details of the configuration options available.
  3. Deploy the scheduled scaling AWS CDK application:
    cdk deploy MSKScheduledScalingStack \
        --app "python3 msk_scheduled_scaling_stack.py" \
        --context cluster_arn="<< ARN of the MSK Cluster >>" \
        --context stack_name="MSKScheduledScalingStack"

Test and monitor the scheduled scaling solution

The scheduled scaling is triggered as specified in the EventBridge Scheduler cron. However, if you want to test the scale-out operations, run the following command to manually invoke the Lambda function:

aws lambda invoke \
  --function-name MSKScheduledScalingStack-MSKScheduledScalingFunction \
  --payload '{"source":"aws.scheduler.scale-out","detail":{"action":"scale_out","schedule_name":"MSKScheduledScaleOut"}}' \
  --cli-binary-format raw-in-base64-out \
  response.json

Similarly, you can manually start a scale-in operation by running the following command:

aws lambda invoke \
  --function-name MSKScheduledScalingStack-MSKScheduledScalingFunction \
  --payload '{"source":"aws.scheduler.scale-in","detail":{"action":"scale_in","schedule_name":"MSKScheduledScaleIn"}}' \
  --cli-binary-format raw-in-base64-out \
  response.json

Monitor the scaling operations by tailing the Lambda function logs:

aws logs tail /aws/lambda/MSKScheduledScalingStack-MSKScheduledScalingFunction  \
--follow --format short

You can monitor scheduled scaling using the CloudWatch dashboard as described in the on-demand scaling section.

Review scaling configuration parameters

The configuration parameters for both on-demand and scheduled scaling are documented in Configuration Options. These configurations give you flexibility to change how and when the scaling happens. It is important to go through the configuration parameters and make sure they meet your business requirement. For on-demand scaling, you can scale the cluster based on built-in performance metrics or custom metrics (for example MessagesInPerSec).

Considerations

Keep in mind the following considerations when deploying either solution:

  • EventBridge notifications for scaling failures – Both on-demand and scheduled scaling solutions publish EventBridge notifications when scaling operations fail. Create EventBridge rules to route these failure events to your monitoring and alerting system to detect failures in scaling and respond to them. For details on event sources, types, and payloads, refer to the EventBridge notifications section in the GitHub repo.
  • Cool-down period management – Properly configure cool-down periods to prevent scaling oscillations where the cluster repeatedly scales out and scales in rapidly. Oscillations typically occur when traffic patterns have short-term spikes that don’t represent sustained demand. Oscillations can also happen when thresholds are set too close to normal operating levels. Set cool-down periods based on your workload characteristics and the scaling completion times. Also consider different cool-down periods for scale-out vs. scale-in operations by setting longer cool-down periods for scale-in operations (scale_in_cooldown_minutes) compared to scaling out (scale_out_cooldown_minutes). Test cool-down settings under realistic load patterns before production deployment to achieve optimal performance.
  • Cost control through monitoring frequency – The solution incurs costs for services like Lambda functions, EventBridge schedules, CloudWatch metrics, and logs that are used in the solution. Both on-demand and scheduled scaling solutions work by running periodically to check the cluster health status and if a scaling operation needs to be performed. The default 1-minute monitoring frequency provides responsive scaling but increases other costs associated with the solution. Consider increasing the monitoring interval based on your workload characteristics to balance scaling responsiveness and the cost incurred by the solution. You can change the monitoring frequency by changing the monitoring_frequency_minutes when you deploy the solution.
  • Solution isolation – The on-demand and scheduled scaling solutions were designed and tested in isolation to support predictable behavior and optimal performance. You can deploy either solution, but avoid running both solutions simultaneously on the same cluster. Using both approaches together can cause unpredictable scaling behavior where the solutions might conflict with each other’s scaling decisions, leading to resource contention and potential scaling oscillations. Choose the approach that best matches your workload patterns and deploy only one scaling solution per cluster.

Clean up

Follow these steps to delete the resources created by the solution. Make sure all the scaling operations that are in flight are completed before you run the cleanup.Delete the on-demand scaling solution with the following code:

cdk destroy MSKOnDemandScalingStack --app "python3 msk_on_demand_scaling_stack.py" --context cluster_arn="<MSK_CLUSTER_ARN>"

Delete the scheduled scaling solution with the following code:

cdk destroy MSKScheduledScalingStack --app "python3 msk_scheduled_scaling_stack.py" --context cluster_arn="<MSK_CLUSTER_ARN>"

Summary

In this post, we showed how to use intelligent rebalancing to scale your Express based cluster based on your business requirements without requiring manual partition rebalancing. You can extend the solution to use the specific CloudWatch metrics that your business depends on to dynamically scale your Kafka cluster. Similarly, you can adjust the scheduled scaling solution to scale out and scale in your cluster when you anticipate significant change in traffic to your cluster at specific times.To learn more about the services used in this solution, refer to the following resources:


About the authors

Subham Rakshit

Subham Rakshit

Subham is a Senior Streaming Solutions Architect for Analytics at AWS based in the UK. He works with customers to design and build streaming architectures so they can get value from analysing their streaming data. His two little daughters keep him occupied most of the time outside work, and he loves solving jigsaw puzzles with them.

Rakshith Rao

Rakshith Rao

Rakshith is a Senior Solutions Architect at AWS. He works with AWS’s strategic customers to build and operate their key workloads on AWS.

Power up your analytics with Amazon SageMaker Unified Studio integration with Tableau, Power BI, and more

Post Syndicated from Narendra Gupta original https://aws.amazon.com/blogs/big-data/power-up-your-analytics-with-amazon-sagemaker-unified-studio-integration-with-tableau-power-bi-and-more/

Organizations face challenges in accessing and analyzing governed data across multiple sources through their preferred business intelligence (BI) and analytics tools while maintaining security and governance. They need a seamless way to connect their familiar tools (like Tableau, Power BI, Excel) to Amazon SageMaker‘s data assets without compromising data governance and security protocols.

Amazon SageMaker supports authentication through the Amazon Athena JDBC driver, allowing data users to query their subscribed data lake assets via popular BI and analytics tools like Tableau, Power BI, Excel, SQL Workbench, DBeaver, and more. This integration empowers data users to access and analyze governed data within Amazon SageMaker using familiar tools, boosting both productivity and flexibility.

Customers use Amazon SageMaker Unified Studio to streamline data access and governance by enabling data users to locate and subscribe to data from multiple sources within a single project. Amazon SageMaker Unified Studio natively integrates with Amazon-specific options like Amazon Athena, Amazon Redshift, and Amazon SageMaker AI, allowing users to analyze their project governed data. With this launch of JDBC connectivity, Amazon SageMaker Unified Studio expands its support for data users, including analysts and scientists, allowing them to work in their preferred tools, whether it’s SQL Workbench, Domino, or Amazon-native solutions like Amazon Athena, while ensuring secure, governed access within Amazon SageMaker Unified Studio.

Getting Started

To get started, download and install the latest Athena JDBC driver for your tool of choice. After installation, copy the JDBC connection string from the Amazon SageMaker Unified Studio portal into the JDBC connection configuration to establish a connection from your tool. This directs you to authenticate using single sign-on (SSO) with your corporate credentials. After connecting, you can query, visualize, and share data—governed by Amazon SageMaker Unified Studio–within the tools you already know and trust.

In this post, we guide you through connecting various analytics tools to Amazon SageMaker Unified Studio using the Athena JDBC driver, enabling seamless access to your subscribed data within your Amazon SageMaker Unified Studio projects.

Solution overview

To demonstrate these capabilities, consider a use case where your marketing team wants to analyze sales data to understand patterns in sales by stores and sales representatives. To achieve this, your marketing team needs access to sales_performance_by_store, and sales_performance_by_rep data owned by the sales team. The sales team, acting as the data producer, publishes the necessary data assets to Amazon SageMaker Unified Studio, allowing the marketing team, as a consumer, to discover and subscribe to these assets.

After the subscription is approved, the data assets become available within the marketing team’s project environment in Amazon SageMaker Unified Studio. The marketing team can then use their preferred tool to perform data exploration. An example architecture of how this is done using DBeaver is shown in the following image:

SageMaker Unified Studio project architecture diagram showing data collaboration between Sales and Marketing teams with Amazon S3 storage and Athena integration

Prerequisites

To follow along with this post, you need the following prerequisites in place:

  1. AWS account – If you don’t have an active AWS account, see How do I create and activate a new AWS account?.
  2. Amazon SageMaker resources – You need a domain for Amazon SageMaker, and two Amazon SageMaker project.
  3. Publish data assets – As the data producer from the sales team, you can now ingest individual data assets into Amazon SageMaker Unified Studio. For this use case, create a data source and import the technical metadata of two data assets – sales_performance_by_store, and sales_performance_by_rep – from AWS Glue Data Catalog. Ensure the data assets are enriched with business descriptions and published to the catalog.
    Note: Here we are using tables which are in the Glue catalog but with Sagemaker Lakehouse you have the option to bring assets from other sources.
  4. Subscribe data assets – As a data analyst from the marketing team, you can now discover and subscribe to the data assets. The data producer from the retail team reviews and approves your subscription. Upon successful fulfillment, the data assets are added to your SageMaker Unified project.

For detailed instructions for publishing and subscribing, see the Amazon SageMaker Unified Studio User Guide.

The following figure shows the subscribed assets added to the subscribed assets section in your marketing project catalog.

SageMaker Unified Studio Assets page displaying subscribed data assets with accessibility status indicators

In the following sections, we walk you through the steps to configure DBeaver to consume the subscribed assets from Amazon SageMaker Unified Studio.

Configuring DBeaver to access subscribed data assets

In this section, you configure DBeaver to access the subscribed assets from the Marketing project

To configure DBeaver:

  1. Connect with JDBC: In the Amazon SageMaker Unified Studio, (1) open the Marketing project, (2) on the Project overview screen, (3) choose JDBC connection details tab.
    SageMaker Unified Studio Project overview page showing JDBC connection parameters for external application integration
  2. Copy the JDBC connection URL into a text editor. The URL should have the following parameters needed for configuring the database connection in DBeaver – Domain ID, Environment ID, Region, and IDC Issuer URL.
    JDBC connection details configuration panel with IDC authentication parameters and copy functionality
  3. Download and install the latest Athena driver:
    • If DBeaver has the Athena driver pre-installed, it might be the older (v2) version. To ensure compatibility with Amazon SageMaker Unified Studio, you need the latest driver (v3), which includes the necessary authentication features.
    • Download the latest JDBC driver—version 3.x.
    • To install the latest driver:
      • Go to Database and then to Driver Manager in DBeaver.
      • Select the Athena driver and choose Edit.
      • Visit the Libraries tab.
      • Choose Download/Update to fetch the latest driver version.
      • If prompted, select the appropriate version and confirm the download.
  4. In the DBeaver SQL client, create a new database connection and select the Athena driver.
    DBeaver database connection dialog showing Amazon Athena driver selection among available database options
  5. Switch to the Driver Properties tab, enter the values of the following properties that are available in the JDBC connection URL you copied from Amazon SageMaker Unified Studio. If any of these properties are not already available, you can add them and provide their respective values.
    • CredentialsProvider: The credentials provider to authenticate requests to AWS
    • DataZoneDomainId: The ID of your Amazon DataZone domain
    • DataZoneDomainRegion: The AWS Region where your domain is hosted
    • DataZoneEnvironmentId: The ID of your DefaultDataLake environment
    • IdentityCenterIssuerUrl: The issuer URL used by AWS Identity and Access Management (IAM) Identity Center for token issuance
    • OutputLocation: Amazon S3 path for storing query results
    • Region: The Region where the environment is created
    • Workgroup: Amazon Athena workgroup of the environment
    • ListenPort: Pick any four digits port number. This is the port number that listens for the IAM Identity Center response

    DBeaver connection configuration dialog for Amazon Athena with driver properties and authentication settings

  6. Choose Test Connection….
  7. You are redirected to the IAM Identity Center sign-in portal. Sign in with Marketing user credentials. If you’re already signed in through single sign-on (SSO), this step can be skipped.
    AWS authentication sign-in page with username input field
  8. After you sign in, if you are prompted to authorize the DataZoneAuthPlugin. Choose Allow access to authorize access to Amazon DataZone from DBeaver.
    AWS DataZone authorization dialog requesting user permission for application access
  9. After sign in completes, you see the following message. You can close the window and go to the DBeaver.
    Amazon DataZone session completion confirmation message
  10. After the connection is established, the following success message appears.
    DBeaver connection test dialog showing successful Amazon Athena connection with performance metrics
  11. You can now view and query all subscribed assets directly within DBeaver.
    DBeaver SQL query interface displaying sales performance data from Amazon Athena database

These steps might also apply to other analytics tools and clients that support JDBC connections. If you’re using a different tool, you might need to adapt these instructions accordingly to ensure proper configuration and access to Amazon SageMaker Unified Studio data assets.

Integration with other applications

You can use similar steps for other BI and analytics tools that support standard database connections.

Connect to Tableau Desktop

Use the Athena JDBC driver to connect Tableau to Amazon SageMaker Unified Studio and visualize your subscribed data.To connect to Tableau Desktop:

  1. Make sure that you’re using the latest Athena JDBC 3.x driver.
  2. Copy the JDBC driver file and place it in the appropriate folders for your operating system
    • For Mac OS: ~/Library/Tableau/Drivers
    • For Windows: C:\Program Files\Tableau\Drivers
  3. Open Tableau Desktop. From the To a Server connection menu, select Other Databases (JDBC) to connect to Amazon SageMaker Unified Studio.
    Tableau start page showing connection options with Other Databases JDBC option highlighted
  4. Paste the JDBC connection URL you copied from the SageMaker Unified Studio portal into the URL. Leave other fields such as Dialect, Username, and Password blank and choose Sign in.
    If you get a port is occupied error – add “;ListenPort=8055” to the URL to change the port. You can use any port number.

    Tableau Other Databases JDBC connection dialog with PostgreSQL dialect configuration

  5. This redirects you to authenticate with IAM Identity Center. Enter the credentials of the Identity Center user that you used to sign in to the SageMaker Unified Studio portal. Authorize the DataZoneAuthPlugin to access Amazon DataZone from Tableau. Once the connection is established with the success message, you can view your project’s subscribed data directly within Tableau and build dashboards.
    Data analytics interface showing sales_performance_by_store table with 283 rows and 15 fields

Connect to Microsoft Power BI

Now, we look at connecting Amazon SageMaker Unified Studio with Microsoft Power BI on Windows.While Amazon Athena provides a native ODBC driver for connecting to ODBC-compatible tools like Microsoft Power BI, it currently doesn’t support Amazon SageMaker Unified Studio authentication. Therefore, in this post, we use an ODBC-JDBC bridge to connect Amazon SageMaker Unified Studio with Microsoft Power BI using the Athena JDBC driver, which supports SageMaker Unified Studio authentication.

In this post, we’re using the ZappySys driver as the ODBC-JDBC bridge. This is a third-party solution that requires a separate licensing fee, which isn’t included in the AWS solution. You can choose to use any other solution for ODBC-JDBC bridge.To connect to Power BI:

  1. Make sure that you have administrator privileges to run the ODBC Data Source Administrator.
  2. From the Windows Start menu, run the ODBC Data Source Administrator (the 64-bit version) using run as Administrator.
  3. Create a New Data Source with the ZappySys JDBC Bridge Driver. You are prompted to enter your connection details.
    Windows ODBC Data Source Administrator dialog showing ZappySys JDBC Bridge Driver selection
  4. Paste the JDBC URL you copied from the SageMaker Unified Studio portal in the Connection String, along with the driver class and JDBC driver file. Make sure that you’re using the latest Athena JDBC 3.x driver.
  5. Choose Test Connection. A new dialog window pops up after the connection is successful.
    Test Connection using ZappySys JDBC Bridge Driver
  6. This redirects you to authenticate with IAM Identity Center. Enter the credentials of the Identity Center user that you used to sign in to the SageMaker Unified Studio portal. Authorize the DataZoneAuthPlugin.
  7. Choose Preview tab on ZappySys JDBC Bridge Driver window and choose one of the subscribed tables to access data.
    ZappySys JDBC Bridge Driver configuration interface showing SQL query preview with sales performance results
  8. After configuring the data source, launch Power BI. Create a blank report or use an existing report to integrate the new visuals. Choose Get Data and select the name of the data source you created. This opens a new browser window to authenticate your credentials. Allow access to authorize the DataZone Auth plugin. After authorization is complete, you can build your reports in Microsoft Power BI with the subscribed data assets.
    Database connection profile selection dialog with PostgreSQL group highlighted

Connect to SQL Workbench

Discover how SQL Workbench can connect to Amazon SageMaker Unified Studio for users who prefer a SQL interface to query data lake tables and views subscribed through projects in Amazon SageMaker Unified Studio.

To connect to SQL Workbench:

  1. Make sure that you’re using the latest Athena JDBC 3.x driver.
  2. Open SQL Workbench/J and choose Manage Drivers.
    Database driver management interface showing SMUSAthenajDBC driver configuration details
  3. Select the option to add a new driver. Enter a name for it, such as SMUSAthenaJDBC, and import the driver you downloaded in the previous steps.
    Database driver management dialog showing SMUSAthenaJDBC driver configuration with library path and class name
  4. Create a new connection profile and enter a name it, such as smus-profile. In the Driver dropdown, select the driver you configured. For the URL, enter the string jdbc:athena://region=us-east-1; (In the example, the Virginia Region is being used). Choose Extended Properties.
    PostgreSQL connection profile configuration dialog with Amazon Athena JDBC driver settings and authentication options
  5. Under Extended Properties, add the following parameters that you copied from the SageMaker Unified Studio portal. You can also include these parameters in the JDBC (URL) connection string. Choose OK.
    • Workgroup
    • OutputLocation
    • DataZoneDomainId
    • IdentityCenterIssuerURL
    • CredentialsProvider
    • DatazoneEnvironmentId
    • DataZoneDomainRegain

    Alos add “ListenPort” with any port number.

    Extended properties configuration dialog showing AWS DataZone connection parameters including domain ID, environment ID, and listen port 8067

  6. This redirects you to authenticate with IAM Identity Center. Enter the credentials of the Identity Center user that you used to sign in to the SageMaker Unified Studio portal. Authorize the DataZoneAuthPlugin.
  7. After successful connection, in SQL Workbench/J, under Database Explorer, select the database from the marketing project of SageMaker unified studio. Choose a subscribed table. Select the Data tab to see the data in the table.
    SQL Workbench showing sales performance data query results from AWS Athena database with 283 customer transaction records

Cleanup

To ensure no additional charges are incurred after testing, be sure to delete the Amazon SageMaker Unified Studio domain. See Delete domains for instructions.

Conclusion

Amazon SageMaker Unified Studio continues to expand its offerings, providing you with more flexibility to access, analyze, and visualize your subscribed data. With support for the Athena JDBC driver, you can now use a wide range of popular BI and analytics tools, making data accessed through Amazon SageMaker Unified Studio more accessible than ever before. Whether you’re using Tableau, Power BI, or other familiar tools, the integration with Amazon SageMaker Unified Studio ensures that your data remains secure and accessible to authorized users.

The feature is supported in all AWS commercial Regions where Amazon SageMaker Unified Studio is currently available. Get started with our technical documentation.


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.

Durga Mishra

Durga Mishra

Durga is a solutions architect at AWS. Outside of work, Durga enjoys spending time with family and loves to hike on Appalachian trails and spend time in nature.

Ramesh Singh

Ramesh Singh

Ramesh is a Senior Product Manager Technical (External Services) at AWS in Seattle, Washington, currently with the Amazon SageMaker team. He is passionate about building high-performance ML/AI and analytics products that help enterprise customers achieve their critical goals using cutting-edge technology.

Nishchai JM

Nishchai JM

Nishchai is an Analytics Specialist Solutions Architect at Amazon Web services. He specializes in building Big-data applications and help customer to modernize their applications on Cloud. He thinks Data is new oil and spends most of his time in deriving insights out of the Data.

Accelerate context-aware data analysis and ML workflows with Amazon SageMaker Data Agent

Post Syndicated from Kshitija Dound original https://aws.amazon.com/blogs/big-data/accelerate-context-aware-data-analysis-and-ml-workflows-with-amazon-sagemaker-data-agent/

Accelerating data analysis and machine learning (ML) development requires AI tools that understand your specific data environment, not just generic code generation. General-purpose AI assistants lack context about your specific data environment, creating a gap between AI capabilities and practical implementation. Data practitioners often start by looking for relevant tables, understanding relationships, and writing exploratory code before answering their first business question. Data teams still spend time translating AI-generated suggestions into working code that correctly references their actual data assets, understands their organization’s data relationships, and integrates with their existing workflows.

AWS released Amazon SageMaker Data Agent in November 2025, addressing these challenges by providing an AI assistant that’s deeply integrated within Amazon SageMaker (IAM-based domains only) with notebooks. SageMaker Data Agent has direct access to your AWS data context, including AWS Glue Data Catalog metadata, Amazon DataZone business data catalog, and your current notebook state. This helps it generate environment-aware code that works directly with your petabyte-scale data through serverless compute resources, helping you analyze massive datasets without infrastructure management overhead. With this contextual awareness, the agent creates executable analysis plans from natural language prompts that specifically reference your actual tables, data types, and analytical needs, while maintaining reasoning throughout multi-step analyses. Importantly, the agent performs these operations securely within the AWS environment, using built-in governance controls, Amazon Identity and Access Management (IAM) policies, and data security features to make sure your data doesn’t leave your organizational boundaries. By operating within your Amazon SageMaker Unified Studio interface, it reduces context-switching between AI assistants and your development environment, improving how you interact with your analytics and ML workflows.

In this post, we demonstrate the capabilities of SageMaker Data Agent, discuss the challenges it addresses, and explore a real-world example analyzing New York City taxi trip data to see the agent in action.

Challenges in data workflows

General AI tools can generate code snippets, but you still face three key challenges when applying these to your specific data environments:

  • Contextual disconnect – Standard AI assistants generate generic code referencing hypothetical tables like customers rather than your actual tables like customer_activity_prod, forcing extensive modifications to work with your data environment.
  • Complex data environment – Many enterprises work with complex data environments containing numerous tables and large-scale data stores, making it extremely difficult to locate relevant data assets for analysis. You must navigate complex catalog structures, understand table relationships, and determine which subset of data is relevant for your specific analytical needs before you can begin actual analysis.
  • Language and syntax barriers – You must work across multiple programming languages and query syntaxes during analysis workflows. Some might excel in SQL but struggle with Python, while others might be Python experts but have limited PySpark knowledge.

Additionally, you face challenges around data quality validation, data governance, and performance optimization. SageMaker Data Agent addresses these fundamental workflow challenges while adapting to your requirements.

Solution overview

SageMaker Data Agent addresses these key challenges through its context-aware architecture and deep AWS integration. In this section, we discuss how it works.

Context-aware understanding

SageMaker Data Agent builds a detailed understanding of your specific data environment and references your actual tables through two parallel processes. SageMaker Data Agent is embedded within your AWS data environment, allowing it to understand what you’re asking, what data you have available, how it’s structured, and how it relates to your analytical objectives. The following are the two ways the agent achieves this contextual understanding:

  • Integrated data environment – SageMaker Data Agent exists within the same integrated environment as your data, harnessing the power of your AWS infrastructure. It begins by exploring the AWS Glue Data Catalog and the Amazon DataZone business data catalog, which reveal business metadata, glossaries, and relationships, enabling it to reference your actual tables rather than generic placeholders. This intelligence extends to working directly with your full datasets where they naturally reside, preserving your existing security policies and access controls without requiring data movement. The agent integrates with Amazon Simple Storage Service (Amazon S3), Amazon Athena, and Amazon SageMaker AI to use their respective capabilities for data storage, query processing, and ML while adapting to your data environment. This lets you process petabyte-scale data through serverless compute resources with the agent acting as an intelligent interface to your complete data environment.
  • Notebook context awareness – Simultaneously, the agent examines your current notebook state, including existing dataframes, imported libraries, previous cell results, and ML artifacts. This context awareness makes sure generated code works with your specific environment without extensive modifications.

Language and syntax flexibility

SageMaker Data Agent resolves language and syntax barriers by selecting the optimal language for each analytical task. The agent can switch between SQL for efficient data querying and Python and PySpark for complex transformations and ML operations without requiring practitioners to manually translate between languages. This avoids language barriers, because the agent automatically selects and generates the appropriate code syntax, whether SQL, Python, or PySpark, based on the specific analytical or ML task at hand.

SageMaker Data Agent provides four key capabilities that work together to give you control over complex analyses:

  • When handling complex requests, the agent creates structured analysis plans by breaking them into logical steps with clear reasoning for each operation.
  • At each stage, you have intermediate validation points where you can review and approve each step before proceeding to the next.
  • Throughout multi-step analyses, the agent maintains consistent context, retaining understanding of your data environment and previous steps.
  • Most importantly, you maintain human-in-the-loop control with full oversight and the ability to modify any generated code to match your specific requirements.

Interaction modes

SageMaker Data Agent provides two interaction modes optimized for different analytical tasks: the Agent Panel and in-line assistance.

The Agent Panel supports comprehensive analytical tasks by breaking them down into structured steps, each with generated code that builds on previous results. When you submit a request such as “perform customer segmentation,” the agent identifies relevant tables, understands their relationships, and creates a complete analysis workflow with intermediate review points. The following screenshot illustrates this example.

In-line assistance mode supports direct cell modifications, one-click error fixes, and keyboard shortcuts (Alt+A for Windows/Linux, Opt+A for Mac) that maintain your coding flow. You can quickly enhance existing code or fix errors without leaving your current notebook context, improving productivity during iterative development. You can code directly within notebook cells by using the inline prompt interface, as illustrated in the following screenshot. Use in-line assistance for focused tasks like specific queries or visualizations directly within cells.

Execution and control

Throughout the process, you maintain execution control. You can review generated plans before execution, execute steps individually with intermediate result review, modify code as needed for your specific requirements by providing feedback, and get AI-powered error diagnosis and fixes using the Fix with AI option when issues arise. This human-in-the-loop approach makes sure you maintain oversight while benefiting from AI assistance.

The following screenshots demonstrate how the Fix with AI feature works in practice, showing how the agent diagnoses code errors and provides corrected solutions with explanations.

By bringing together context-aware understanding, reasoning, and interaction modes within your existing AWS environment, SageMaker Data Agent improves how you work. It removes the traditional friction between AI assistance and your actual data environment, providing direct access to petabyte-scale data with no operational overhead. This combination helps you shift your focus from repetitive setup tasks to high-value analysis and decision-making, accelerating insights while maintaining control over the analytical process.

Getting started with SageMaker Data Agent

Now that you understand how SageMaker Data Agent works, let’s see these capabilities in action. Getting started with SageMaker Data Agent is straightforward. For detailed setup instructions, refer to New one-click onboarding and notebooks with a built-in AI agent in Amazon SageMaker Unified Studio. It provides step-by-step guidance on setting up your environment and beginning your journey with SageMaker Data Agent.

To get the most from SageMaker Data Agent, begin by asking clear, specific questions about your data rather than generic requests. Provide context about your analytical goals so the agent can tailor its responses to your specific use case. Always review and validate generated code before execution, using the agent’s built-in explanations to understand the approach. For complex analyses, take advantage of the agent’s reasoning capabilities that can break down multi-step processes and explain the logic behind each recommendation.

NYC taxi trip analysis

In this section, we demonstrate how SageMaker Data Agent helps analyze the NYC Taxi Trip dataset, a collection of over 1.2 billion taxi trips (approximately 63.7 GB) throughout New York City with information on pickup/drop-off locations, timestamps, trip distances, fare amounts, payment types, and passenger counts.

If you’re looking to try a simpler end-to-end flow before diving into this large-scale analysis, SageMaker Unified Studio provides a sample database with pre-loaded customer churn data. You can perform similar analytical workflows on this smaller dataset to quickly familiarize yourself with the agent’s capabilities before working with larger, more complex datasets. To explore this dataset, complete the following steps:

  1. On the SageMaker Unified Studio console, choose Data in the navigation pane.
  2. In the data explorer, under Catalogs, select AwsDataCatalog.
  3. Select sagemaker_sample_db.
  4. Select the churn table from the tables list.

NYC Taxi Trip dataset

The NYC Taxi Trip dataset is publicly available in Amazon S3 at s3://aws-data-analytics-workshops/shared_datasets/nyc_taxi_trips_parquet/.

To replicate this, you can work with this dataset in two ways:

  • Catalog it beforehand (recommended for repeated analysis)
  • Provide the S3 path directly in your prompt (quickest for one-time exploration)

For this demonstration, we used SageMaker Data Agent to catalog the dataset prior to analysis.

Our analysis approach

For this demonstration, we asked SageMaker Data Agent to perform a comprehensive analysis on the cataloged taxi trip data to uncover business insights. We used the following prompt:

Using Apache Spark, analyze the NYC taxi trips dataset to extract meaningful insights. Please provide:
1/ Fare analysis across different NYC boroughs
2/ Trip trends across boroughs and time
Conclude with multi-panel dashboard and an executive summary highlighting the 3-5 most significant findings and their potential business implications.

You can add the S3 path (s3://aws-data-analytics-workshops/shared_datasets/nyc_taxi_trips_parquet/) in the preceding prompt if you don’t have the NYC Taxi Trip data cataloged.

The following video demonstrates how SageMaker Data Agent processes this natural language prompt and creates a complete analytical workflow. The agent constructs a six-step analysis plan, generates executable code for each step, and progressively builds toward actionable insights.

The outputs shown in this demonstration video are specific to this analysis session. Due to the generative nature of AI, your results might vary when running the same prompts.The agent executed each step sequentially, so we can review intermediate results and provide feedback. After loading and cleaning NYC taxi trip records, the agent analyzed fare patterns and trip trends across boroughs and time periods, then created a comprehensive multi-panel dashboard visualizing key insights, as shown in the following screenshots.

Finally, it provided actionable business insights, highlighting the most significant findings and their business recommendations.

This example demonstrates how SageMaker Data Agent helps transform complex analytical tasks into actionable insights without requiring extensive coding or data preparation. The agent’s ability to understand both the data structure and business context allows it to generate meaningful analyses that directly address business objectives.

Security and governance

SageMaker Data Agent follows your AWS security settings. It accesses data you’ve explicitly permitted through your IAM access controls or using AWS Lake Formation, helping maintain your organization’s security policies. To use SageMaker Data Agent, your project role must have permissions to invoke specific Amazon DataZone APIs, including SendMessage, GenerateCode, StartConversation, GetConversation, and ListConversations. For more information, visit Actions, resources, and condition keys for Amazon DataZone.

Guardrails

SageMaker Data Agent has in-built guardrails to prevent the agent from responding to undesired requests. These include but are not limited to requests asking the agent to reveal its system prompt, internal tools, or other technical implementation. These guardrails also prohibit the agent from talking about non-AWS related topics and from generating output in any language except English.

Data storage and privacy

SageMaker Data Agent doesn’t store code you write or modify yourself, notebook context or metadata, or data from your AWS Glue Data Catalog or other sources. The agent only stores your natural language prompts, questions, and generated code/responses in the AWS Region where your SageMaker Unified Studio domain was created. AWS might use stored content (prompts, questions, and generated code/responses) to improve the service, fix issues, or for debugging, but maintains clear boundaries by not using your self-written code, manually modified code, notebook metadata, or actual data sources for service improvement. To opt out of data usage for service improvement, you can configure an AI services opt-out policy for Amazon DataZone in AWS Organizations, which will delete previously collected data and prevent future collection or usage. For more information, refer to Data storage in the SageMaker Data Agent, Service improvement, and AI services opt-out policies.

Conclusion

SageMaker Data Agent improves how data practitioners accelerate insights. By combining context-aware understanding, AWS integration, and flexible interaction modes, it alleviates the traditional friction between AI-assisted development and your actual data environment. The NYC taxi analysis demonstrated this in practice: what might have required manual data exploration, catalog navigation, and code translation instead took minutes through natural language prompts.

The real value extends beyond speed. SageMaker Data Agent preserves your security posture, maintains governance controls, and keeps your data within your AWS environment while supporting petabyte-scale analysis without operational overhead. More importantly, it shifts your team’s focus from repetitive setup to business analysis and decision-making.

Getting started is straightforward. Begin with simple prompts against your existing data catalog, then progressively tackle more complex analytical challenges. Invest time enriching your data catalog with business metadata—this investment directly multiplies the agent’s effectiveness by providing richer context for code generation.

SageMaker Data Agent adapts to your specific analytical needs, such as analyzing customer behavior, working with financial data, or building ML models. Access it today through your IAM-based SageMaker Unified Studio domain, and discover how context-aware AI assistance can accelerate your organization’s data-driven decision-making.


About the authors

Kshitija Dound

Kshitija Dound

Kshitija is a Specialist Solutions Architect at AWS based in New York City, focusing on data and AI. She collaborates with customers to transform their ideas into cloud solutions, using AWS Big Data and AI services. She also engages in public speaking opportunities, sharing her expertise on cloud technologies, industry trends, and career in the cloud. In her spare time, Kshitija enjoys exploring museums, indulging in art, and embracing NYC’s outdoor scene.

Siddharth Gupta

Siddharth Gupta

Siddharth is heading Generative AI within SageMaker’s Unified Experiences. His focus is on driving agentic experiences, where AI systems act autonomously on behalf of users to accomplish complex tasks. An alumnus of the University of Illinois at Urbana-Champaign, he brings extensive experience from his roles at Yahoo, Glassdoor, and Twitch.

Mohan Gandhi

Mohan Gandhi

Mohan is a Principal Software Engineer at AWS. He has been with AWS for the last 10 years and has worked on various AWS services like Amazon EMR, Amazon EFA, and Amazon RDS. Currently, he is focused on improving the Amazon SageMaker inference experience. In his spare time, he enjoys hiking and marathons.

Ishneet Kaur

Ishneet Kaur

Ishneet is a Software Development Manager on the Amazon SageMaker Unified Studio team. She leads the engineering team to design and build generative AI capabilities in SageMaker Unified Studio.

Shubham Mehta

Shubham Mehta

Shubham is a Senior Product Manager at AWS Analytics. He leads generative AI feature development across services such as AWS Glue, Amazon EMR, and Amazon MWAA, using AI/ML to simplify and enhance the experience of data practitioners building data applications on AWS.

Vikramank Singh

Vikramank Singh

Vikramank is a Senior Applied Scientist in the Agentic AI organization in AWS, working on products including Amazon SageMaker Unified Studio, Amazon RDS, and Amazon Redshift. His research interest lies at the intersection of AI, control systems, and RL, particularly using them to build systems for real-world applications that can autonomously perceive environments, model them, and take optimal decisions at scale.

Murali Narayanaswamy

Murali Narayanaswamy

Murali is a Principal Machine Learning Scientist in the Agentic AI organization in AWS, working on products including Amazon SageMaker Unified Studio, Amazon Redshift, and Amazon RDS. His research interests lie at the intersection of AI, optimization, learning, and inference, particularly using them to understand, model, and combat noise and uncertainty in real-world applications and reinforcement learning in practice and at scale.

Amit Sinha

Amit Sinha

Amit is a Senior Manager leading SageMaker Unified Studio GenAI and ML product suites. He has over a decade of experience in AI/ML products, infrastructure management, and AWS Big Data processing services. An alumnus of Columbia University, in his free time Amit enjoys hiking and binge-watching documentaries on American history.

Streamline large binary object migrations: A Kafka-based solution for Oracle to Amazon Aurora PostgreSQL and Amazon S3

Post Syndicated from Naresh Dhiman original https://aws.amazon.com/blogs/big-data/streamline-large-binary-object-migrations-a-kafka-based-solution-for-oracle-to-amazon-aurora-postgresql-and-amazon-s3/

Customers migrating from on-premises Oracle databases to AWS face a challenge: efficiently relocating large object data types (LOBs) to object storage while maintaining data integrity and performance. This challenge originates from the traditional enterprise database design where LOBs are stored alongside structured data, leading to storage capacity constraints, backup complexity, and performance bottlenecks during data retrieval and processing. LOBs, which can include images, videos, and other large files, often cause traditional data migrations to suffer from slow speeds and LOB truncation issues. These issues are particularly problematic for long-running migrations that can span several years.

In this post, we present a scalable solution that uses Amazon Managed Streaming for Apache Kafka (Amazon MSK), Amazon Aurora PostgreSQL-Compatible Edition, and Amazon MSK Connect. The data streaming enables data replication where modifications are sent and received in a continuous flow, allowing the target database to access and apply the changes in real time. This solution generates events for database actions such as insert, update, and delete, triggering AWS Lambda functions to download LOBs from the source Oracle database and upload them to Amazon Simple Storage Service (Amazon S3) buckets. Simultaneously, the streaming events migrate the structured data from the Oracle database to the target database while maintaining proper linking with their respective LOBs.

The complete implementation is available on GitHub, including AWS Cloud Development Kit (AWS CDK) deployment code, configuration files, and setup instructions.

Solution overview

Although traditional Oracle database migrations handle structured data effectively, they struggle with LOBs that can include images, videos, and documents. These migrations often fail due to size limitations and truncation issues, creating significant business risks, including data loss, extended downtime, and project delays that can force you to delay your cloud transformation initiatives. The problem becomes more acute during long-running migrations spanning several years, where maintaining operational continuity is critical. This solution addresses the key challenges of LOB migration, enabling continuous, long-term operations without compromising performance or reliability.

By removing the size limitations associated with traditional migration technologies, our solution provides a robust framework that helps you seamlessly relocate LOBs while facilitating data integrity throughout the process.

Our approach uses a modern streaming architecture to alleviate the traditional constraints of Oracle LOB migration. The solution includes the following core components:

  • Amazon MSK – Provides the streaming infrastructure.
  • Amazon MSK Connect – Using two connectors:
    • Debezium Connector for Oracle as a source connector to capture row-level changes that occur in Oracle database. The connector emits change events and publishes to a Kafka source topic.
    • Debezium Connector for JDBC as a sink connector to consume events from Kafka source topic and then write those events to Aurora PostgreSQL-Compatible by using a JDBC driver.
  • Lambda function – Triggered by an event source mapping to Amazon MSK. The function processes events from the Kafka source topic, extracting the Oracle row primary key from each event payload. It uses this key to download the corresponding BLOB data from the source Oracle database and uploads it to Amazon S3, organizing files by primary key folders to maintain simple linking with the relational database records.
  • Amazon RDS for Oracle – Amazon Relational Database Service (Amazon RDS) for Oracle is used as the source database to simulate an on-premises Oracle database.
  • Aurora PostgreSQL-Compatible – Used as the target database for migrated data.
  • Amazon S3 – Used as object storage for storing the BLOB data from source database.

The following diagram shows the Oracle LOB data migration architecture solution.

Message flow

When data changes occur in the source Amazon RDS for Oracle database, the solution executes the following sequence, moving through event detection and publication, BLOB processing with Lambda, and structured data processing:

  1. The Oracle source connector captures the change data capture (CDC) events, including the change to BLOB data column. This connector configures the BLOB data column to exclude from the Kafka event to optimize the Kafka payload.
  2. The connector publishes this event to an MSK topic.
    1. The MSK event triggers the BLOB Downloader Lambda function for the CDC events.
      1. The Lambda function examines two key conditions: the Debezium event code (specifically checking for create (c) or update(u)) and the configured list of Oracle BLOB table names along with their column names. When a Kafka message matches both the configured table list and valid Debezium events, the Lambda function initiates the BLOB data download from the Oracle source using the primary key and table name; otherwise, the function bypasses the BLOB download process. This selective approach makes sure the Lambda function only executes SQL queries when processing Kafka messages for tables containing BLOB data, optimizing database interactions.
      2. The Lambda function uploads the BLOB to Amazon S3, organizing by primary key folders with unique object names, which enables linking between structured database records and their corresponding BLOB data in Amazon S3.
    2. The PostgreSQL sink connector receives the event from the MSK topic.
      1. The connector applies these changes to the Aurora PostgreSQL database for the Oracle database changes except the BLOB data column. The BLOB data column is excluded by the Oracle source connector.

Key benefits

The solution offers the following key advantages:

  • Cost optimization and licensing – Our approach offers significant cost optimization benefits by reducing the overall size of your database and alleviating your need for expensive licenses associated with traditional databases and replication technologies. By decoupling LOB storage from the database and using Amazon S3, you can reduce your overall database footprint and reduce costs associated with traditional licensing and replication technologies. The streaming architecture also minimizes your infrastructure overhead during long-running migrations.
  • Avoids size constraints and migration failures – Traditional migration tools often impose size limitations on LOB transfers, leading to truncation issues and failed migrations. This solution removes those constraints entirely, so you can migrate LOBs of different sizes while maintaining data integrity. The event-driven architecture enables near real-time data replication, allowing your source systems to remain operational during migration.
  • Business continuity and operational excellence – Changes flow continuously to your target environment, allowing for business continuity. The solution preserves relationships between structured database records and their corresponding LOBs through primary key-based organization in Amazon S3, allowing for referential integrity while providing the flexibility of object storage for large files.
  • Architectural advantages – Storing LOBs in Amazon S3 while maintaining structured data in Aurora PostgreSQL-Compatible creates a clear separation. This architecture simplifies your backup and recovery operations, improves query performance on structured data, and provides flexible access patterns for binary objects through Amazon S3.

Implementation best practices

Consider the following best practices when implementing this solution:

  • Start small and scale gradually – To implement this solution, start with a pilot project using non-production data to validate your approach before committing to full-scale migration. This gives you a chance to work out issues in a controlled environment and refine your configuration without impacting production systems.
  • Monitoring – Set up comprehensive monitoring through Amazon CloudWatch to track key metrics like Kafka lag, Lambda function errors, and replication latency. Establish alerting thresholds early so you can catch and resolve issues quickly before they impact your migration timeline. Size your MSK cluster based on expected CDC volume and configure Lambda reserved concurrency to handle peak loads during initial data synchronization.
  • Security – For security, use encryption in transit and at rest for both structured data and LOBs, and follow the principle of least privilege when setting up AWS Identity and Access Management (IAM) roles and policies for your MSK cluster, Lambda functions, S3 buckets, and database instances. Document your schema mappings between Oracle and Aurora PostgreSQL-Compatible, including how database records link to their corresponding LOBs in Amazon S3.
  • Testing and preparation – Before you go live, test your failover and recovery procedures thoroughly. Validate scenarios like Lambda function failures, MSK cluster issues, and network connectivity problems to ensure you’re prepared for potential issues. Finally, remember that this streaming architecture maintains eventual consistency between your source and target systems, so there might be brief lag times during high-volume periods. Plan your cutover strategy with this in mind.

Limitations and considerations

Although this solution provides a robust approach for migrating Oracle databases with LOBs to AWS, there are several inherent constraints to understand before implementation.

This solution requires network connectivity between your source Oracle database and AWS environment. For on-premises Oracle databases, you must establish AWS Direct Connect or VPN connectivity before deployment. Network bandwidth directly impacts replication speed and overall migration performance, so your connection must be able to handle the expected volume of CDC events and LOB transfers.

The solution uses Debezium Connector for Oracle as the source connector and Debezium Connector for JDBC as the sink connector. This architecture is specifically designed for your Oracle-to-PostgreSQL migrations. Other database combinations require different connector configurations or might not be supported by the current implementation. Migration throughput is also constrained by your MSK cluster capacity and Lambda concurrency limits. You can also exceed AWS service quotas for large-scale migrations and you might need to request quota increases through AWS Enterprise Support.

Conclusion

In this post, we presented a solution that addresses the critical challenge of migrating your large binary objects from Oracle to AWS by using a streaming architecture that separates LOB storage from structured data. This approach avoids size constraints, reduces Oracle licensing costs, and preserves data integrity throughout extended migration periods.

Ready to transform your Oracle migration strategy? Visit the GitHub repository, where you will find the complete AWS CDK deployment code, configuration files, and step-by-step instructions to get started.


About the authors

Naresh Dhiman

Naresh Dhiman

Naresh is a Sr. Solutions Architect at AWS supporting US federal customers. He has over 25 years of experience as a technology leader and is a recognized inventor with six patents. He specializes in containers, machine learning, and generative AI on AWS.

Archana Sharma

Archana Sharma

Archana is a Sr. Database Specialist Solutions Architect, working with Worldwide Public Sector customers. She has years of experience in relational databases, and is passionate about helping customers in their journey to the AWS Cloud with a focus on database migration and modernization.

Ron Kolwitz

Ron Kolwitz

Ron is a Sr. Solutions Architect supporting US Federal Government Sciences customers including NASA and the Department of Energy. He is especially passionate about aerospace and advancing the use of GenAI and quantum-based technologies for scientific research. In his free time, he enjoys spending time with his family of avid water-skiers.

Karan Lakhwani

Karan Lakhwani

Karan is a Sr. Customer Solutions Manager at Amazon Web Services. He specializes in generative AI technologies and is an AWS Golden Jacket recipient. Outside of work, Karan enjoys finding new restaurants and skiing.

How Bazaarvoice modernized their Apache Kafka infrastructure with Amazon MSK

Post Syndicated from Oleh Khoruzhenko original https://aws.amazon.com/blogs/big-data/how-bazaarvoice-modernized-their-apache-kafka-infrastructure-with-amazon-msk/

This is a guest post by Oleh Khoruzhenko, Senior Staff DevOps Engineer at Bazaarvoice, in partnership with AWS.

Bazaarvoice is an Austin-based company powering a world-leading reviews and ratings platform. Our system processes billions of consumer interactions through ratings, reviews, images, and videos, helping brands and retailers build shopper confidence and drive sales by using authentic user-generated content (UGC) across the customer journey. The Bazaarvoice Trust Mark is the gold standard in authenticity.

Apache Kafka is one of the core components of our infrastructure, enabling real-time data streaming for the global review platform. Although Kafka’s distributed architecture met our needs for high-throughput, fault-tolerant streaming, self-managing this complex system diverted critical engineering resources away from our core product development. Each component of our Kafka infrastructure required specialized expertise, ranging from configuring low-level parameters to maintaining the complex distributed systems that our customers rely on. The dynamic nature of our environment demanded continuous care and investment in automation. We found ourselves constantly managing upgrades, applying security patches, implementing fixes, and addressing scaling needs as our data volumes grew.

In this post, we show you the steps we took to migrate our workloads from self-hosted Kafka to Amazon Managed Streaming for Apache Kafka (Amazon MSK). We walk you through our migration process and highlight the improvements we achieved after this transition. We show how we minimized operational overhead, enhanced our security and compliance posture, automated key processes, and built a more resilient platform while maintaining the high performance our global customer base expects.

The need for modernization

As our platform grew to process billions of daily consumer interactions, we needed to find a way to scale our Kafka clusters efficiently while maintaining a small team to manage the infrastructure. The limitations of self-managed Kafka clusters manifested in several key areas:

  • Scaling operations – Although scaling our self-hosted Kafka clusters wasn’t inherently complex, it required careful planning and execution. Each time we needed to add new brokers to handle increased workload, our team faced a multi-step process involving capacity planning, infrastructure provisioning, and configuration updates.
  • Configuration complexity – Kafka offers hundreds of configuration parameters. Although we didn’t actively manage all of these, understanding their impact was important. Key settings like I/O threads, memory buffers, and retention policies needed ongoing attention as we scaled. Even minor adjustments could have significant downstream effects, requiring our team to maintain deep expertise in these parameters and their interactions to ensure optimal performance and stability.
  • Infrastructure management and capacity planning – Self-hosting Kafka required us to manage multiple scaling dimensions, including compute, memory, network throughput, storage throughput, and storage volume. We needed to carefully plan capacity for all these components, often making complex trade-offs. Beyond capacity planning, we were responsible for real-time management of our Kafka infrastructure. This included promptly detecting and addressing component failures and performance issues. Our team needed to be highly responsive to alerts, often requiring immediate action to maintain system stability.
  • Specialized expertise requirements – Operating Kafka at scale demanded deep technical expertise across multiple domains. The team needed to:
    • Monitor and analyze hundreds of performance metrics
    • Conduct complex root cause analysis for performance issues
    • Manage ZooKeeper ensemble coordination
    • Execute rolling updates for zero-downtime upgrades and security patches

These challenges were compounded during peak business periods, such as Black Friday and Cyber Monday, when maintaining optimal performance was essential for Bazaarvoice’s retail customers.

Choosing Amazon MSK

After evaluating various options, we selected Amazon MSK as our modernization solution. The decision was driven by the service’s ability to minimize operational overhead, provide high availability out of the box with its three Availability Zone architecture, and offer seamless integration with our existing AWS infrastructure.

Key capabilities that made Amazon MSK the clear choice:

  • AWS integration – We already used AWS services for data processing and analytics. Amazon MSK connected directly with these services, alleviating the need to build and maintain custom integrations. This meant our existing data pipelines would continue working with minimal changes.
  • Automated operations management – Amazon MSK automated our most time-consuming tasks. We no longer need to manually monitor instances and storage for failures or respond to these issues ourselves.
  • Enterprise-grade reliability – The platform’s architecture matched our reliability requirements out of the box. Multi-AZ distribution and built-in replication gave us the same fault tolerance we’d carefully built into our self-hosted system, now backed by AWS’s service guarantees.
  • Simplified upgrade process – Before Amazon MSK, version upgrades for our Kafka clusters required careful planning and execution. The process was complex, involving multiple steps and risks. Amazon MSK simplified our upgrade operations. We now use automated upgrades for dev and test workloads and maintain control over production environments. This shift reduced the need for extensive planning sessions and multiple engineers. As a result, we stay current with the latest Kafka versions and security patches, improving our system reliability and performance.
  • Enhanced security controls – Our platform required ISO 27001 compliance, which typically involved months of documentation and security controls implementation. Amazon MSK came with this certification built-in, alleviating the need for separate compliance work. Amazon MSK encrypted our data, controlled network access, and integrated with our existing security tools.

With Amazon MSK selected as our target platform, we began planning the complex task of migrating our critical streaming infrastructure without disrupting the billions of consumer interactions flowing through our system.

Bazaarvoice’s migration journey

Moving our complex Kafka infrastructure to Amazon MSK required careful planning and precise execution. Our platform processes data through two main components: an Apache Kafka Streams pipeline that handles data processing and augmentation, and client applications that move this enriched data to downstream systems. With 40 TB of state across 250 internal topics, this migration demanded a methodical approach.

Planning phase

Working with AWS Solutions Architects proved critical for validating our migration strategy. Our platform’s unique characteristics required special consideration:

  • Multi-Region deployment across the US and EU
  • Complex stateful applications with strict data consistency needs
  • Vital business services requiring zero downtime
  • Diverse consumer ecosystem with different migration requirements

Migration challenges

The biggest hurdle was migrating our stateful Kafka Streams applications. Our data processing runs as a directed acyclic graph (DAG) of applications across regions, using static group membership to prevent disruptive rebalancing. It’s important to note that Kafka Streams keeps its state in internal Kafka topics. For applications to recover properly, replicating this state accurately is crucial. This characteristic of Kafka Streams added complexity to our migration process. Initially, we considered MirrorMaker2, the standard tool for Kafka migrations. However, two fundamental limitations made it challenging:

  • Risk of losing state or incorrectly replicating state across our applications.
  • Inability to run two instances of our applications simultaneously, which meant we needed to shut down the main application and wait for it to recover from the state in the MSK cluster. Given the size of our state, this recovery process exceeded our 30-minute SLA for downtime.

Our solution

We decided to deploy a parallel stack of Kafka Streams applications reading and writing data from Amazon MSK. This approach gave us sufficient time for testing and verification, and enabled the applications to hydrate their state before we delivered the output to our data warehouse for analytics. We used MirrorMaker2 for input topic replication, while our solution offered several advantages:

  • Simplified monitoring of the replication process
  • Avoided consistency issues between state stores and internal topics
  • Allowed for gradual, controlled migration of consumers
  • Enabled thorough validation before cutover
  • Required a coordinated transition plan for all consumers, because we couldn’t transfer consumer offsets across clusters

Consumer migration strategy

Each consumer type required a carefully tailored approach:

  • Standard consumers – For applications supporting Kafka Consumer Group protocol, we implemented a four-step migration. This approach risked some duplicate processing, but our applications were designed to handle this scenario. The steps were as follows:
    • Configure consumers with auto.offset.reset: latest.
    • Stop all DAG producers.
    • Wait for existing consumers to process remaining messages.
    • Cut over consumer applications to Amazon MSK.
  • Apache Kafka Connect Sinks – Our sink connectors served two critical databases:
    • A distributed search and analytics engine – Document versioning depended on Kafka record offsets, making direct migration impossible. To address this, we implemented a solution that involved building new search engine clusters from scratch.
    • A document-oriented NoSQL database – This supported direct migration without requiring new database instances, simplifying the process significantly.
  • Apache Spark and Flink applications – These presented unique challenges due to their internal checkpointing mechanisms:
    • Offsets managed outside Kafka’s consumer groups
    • Checkpoints incompatible between source and target clusters
    • Required complete data reprocessing from the beginning

We scheduled these migrations during off-peak hours to minimize impact.

Technical benefits and improvements

Moving to Amazon MSK fundamentally changed how we manage our Kafka infrastructure. The transformation is best illustrated by comparing key operational tasks before and after the migration, summarized in the following table.

Activity Before: Self-Hosted Kafka After: Amazon MSK
Security patching Required dedicated team time for Kafka and OS updates Fully automated
Broker recovery Needed manual monitoring and intervention Fully automated
Client authentication Complex password rotation procedures AWS Identity and Access Management (IAM)
Version upgrades Complex procedure requiring extensive planning Fully automated

The details of the tasks are as follows:

  • Security patching – Previously, our team spent 8 hours monthly applying Kafka and operating system (OS) security patches across our broker fleet. Amazon MSK now handles these updates automatically, maintaining our security posture without engineering intervention.
  • Broker recovery – Although our self-hosted Kafka had automatic recovery capabilities, each incident required careful monitoring and occasional manual intervention. With Amazon MSK, node failures and storage degradation issues such as Amazon Elastic Block Store (Amazon EBS) slowdowns are handled entirely by AWS and resolved within minutes without our involvement.
  • Authentication management – Our self-hosted implementation required password rotations for SASL/SCRAM authentication, a process that took two engineers several days to coordinate. The direct integration between Amazon MSK and AWS Identity and Access Management (IAM) minimized this overhead while strengthening our security controls.
  • Version upgrades – Kafka version upgrades in our self-hosted environment required weeks of planning and testing as well as weekend maintenance windows. Amazon MSK manages these upgrades automatically during off-peak hours, maintaining our SLAs without disruption.

These improvements proved especially valuable during high-traffic periods like Black Friday, when our team previously needed extensive operational readiness plans. Now, the built-in resiliency of Amazon MSK provides us with reliable Kafka clusters that serve as mission-critical infrastructure for our business. The migration made it possible to break our monolithic clusters into smaller, dedicated MSK clusters. This improved our data isolation, provided better resource allocation, and enhanced performance predictability for high-priority workloads.

Lessons learned

Our migration to Amazon MSK revealed several key insights that can help other organizations modernize their Kafka infrastructure:

  • Expert validation – Working with AWS Solutions Architects to validate our migration strategy caught several critical issues early. Although our team knew our applications well, external Kafka experts identified potential problems with state management and consumer offset handling that we hadn’t considered. This validation prevented costly missteps during the migration.
  • Data verification – Comparing data across Kafka clusters proved challenging. We built tools to capture topic snapshots in Parquet format on Amazon Simple Storage Service (Amazon S3), enabling quick comparisons using Amazon Athena queries. This approach gave us confidence that data remained consistent throughout the migration.
  • Start small – Beginning with our smallest data universe in QA helped us refine our process. Each subsequent migration went smoother as we applied lessons from previous iterations. This gradual approach helped us maintain system stability while building team confidence.
  • Detailed planning – We created specific migration plans with each team, considering their unique requirements and constraints. For example, our machine learning pipeline needed special handling due to strict offset management requirements. This granular planning prevented downstream disruptions.
  • Performance optimization – We found that utilizing Amazon MSK provisioned throughput offered clear cost advantages when storage throughput became a bottleneck. This feature made it possible to improve cluster performance without scaling instance sizes or adding brokers, providing a more efficient solution to our throughput challenges.
  • Documentation – Maintaining detailed migration runbooks proved invaluable. When we encountered similar issues across different migrations, having documented solutions saved significant troubleshooting time.

Conclusion

In this post, we showed you how we modernized our Kafka infrastructure by migrating to Amazon MSK. We walked through our decision-making process, challenges faced, and strategies employed. Our journey transformed Kafka operations from a resource-intensive, self-managed infrastructure to a streamlined, managed service, improving operational efficiency, platform reliability, and team productivity. For enterprises managing self-hosted Kafka infrastructure, our experience demonstrates that successful transformation is achievable with proper planning and execution. As data streaming needs grow, modernizing infrastructure becomes a strategic imperative for maintaining competitive advantage.

For more information, visit the Amazon MSK product page, and explore the comprehensive Developer Guide to learn about the features available to help you build scalable and reliable streaming data applications on AWS.

About the authors

Oleh Khoruzhenko

Oleh Khoruzhenko

Oleh is a Senior Staff DevOps Engineer at Bazaarvoice Inc, specializing in architecting and optimizing high-throughput data streaming solutions. He is an expert in the Apache Kafka ecosystem, utilizing Apache Kafka Streams for complex event processing and Apache Spark for large-scale data ingestion and transformation.

Christian Silva

Christian Silva

Christian is a Sr. Solutions Architect at AWS based in Houston, TX. He works with independent software vendor customers, helping them build and optimize their solutions on AWS. With a background in cloud architecture, Christian is passionate about networking and security, guiding customers to implement robust and efficient cloud infrastructures. Outside of work, he enjoys spending time with his kids, playing soccer, and fishing.

Aravind Marthineni

Aravind Marthineni

Aravind is a Technical Account Manager at AWS based in Austin, TX. He supports SaaS customers in ecommerce and social media analytics on migrations and modernizations using cloud-based architectures. He specializes in cloud governance and is passionate about helping customers operate efficiently on the cloud using generative AI. When not working, he loves playing and teaching cricket to his 2-year-old and cooking Indian delicacies.

Enterprise scale in-place migration to Apache Iceberg: Implementation guide

Post Syndicated from Mihir Borkar original https://aws.amazon.com/blogs/big-data/enterprise-scale-in-place-migration-to-apache-iceberg-implementation-guide/

Organizations managing large-scale analytical workloads increasingly face challenges with traditional Apache Parquet-based data lakes with Hive-style partitioning, including slow queries, complex file management, and limited consistency guarantees. Apache Iceberg addresses these pain points by providing ACID transactions, seamless schema evolution, and point-in-time data recovery capabilities that transform how enterprises handle their data infrastructure.

In this post, we demonstrate how you can achieve migration at scale from existing Parquet tables to Apache Iceberg tables. Using Amazon DynamoDB as a central orchestration mechanism, we show how you can implement in-place migrations that are highly configurable, repeatable, and fault-tolerant—unlocking the full potential of modern data lake architectures without extensive data movement or duplication.

Solution overview

When performing in-place migration, Apache Iceberg uses its ability to directly reference existing data files. This capability is only supported for formats such as Parquet, ORC, and Avro, because these formats are self-describing and include consistent schema and metadata information. Unlike raw formats such as CSV or JSON, they enforce structure and support efficient columnar or row-based access, which allows Iceberg to integrate them without rewriting the data.

In this post, we demonstrate how you can migrate an existing Parquet-based data lake that isn’t cataloged in AWS Glue by using two methodologies:

  • Apache Iceberg migrate and register_table approach. Ideal for converting existing Hive-registered Parquet tables into Iceberg-managed tables.
  • Iceberg add_files approach. Best suited for quickly onboarding raw Parquet data into Iceberg without rewriting files.

The solution also incorporates a DynamoDB table that acts as a scalable control plane, so you can perform in-place migration of your data lake from Parquet format to Iceberg format.

The following diagram shows different methodologies that you can use to achieve this in-place migration of your Hive-style partitioned data lake:

AWS data pipeline architecture diagram showing data flow from Amazon DynamoDB through Amazon EMR and AWS Glue to a Data Lake and Apache Iceberg Lakehouse, both using Parquet format, within an AWS Region.

You use DynamoDB to track the migration state, handling retries and recording errors and outcomes. This provides the following benefits:

  • Centralized control over which Amazon Simple Storage Service (Amazon S3) paths need migration.
  • Lifecycle tracking of each dataset through migration stages.
  • Capture and audit errors on a per-path basis.
  • Enable re-runs by updating stateful flags or clearing failure messages.

Prerequisites

Before you begin, you need:

Create sample Parquet dataset as a source

You can create the sample Parquet dataset for testing the different methodologies using the Athena query editor. Replace <amzn-s3-demo-bucket> with an available bucket in your account.

  1. Create an AWS Glue database(test_db), if not present.
    CREATE DATABASE IF NOT EXISTS test_db

  2. Create a sample Parquet table (table1) and add to be used for testing the add_files approach.
    CREATE TABLE table1
    WITH (
      external_location = 's3://<amzn-s3-demo-bucket>/table1/',
      format = 'PARQUET',
      partitioned_by = ARRAY['date', 'hour']
    )
    AS
    SELECT 
      1 as id,
      'John Doe' as name,
      25 as age,
      'Engineer' as job_title,
      current_date as created_date,
      current_date as date,
      hour(current_timestamp) as hour
    UNION ALL
    SELECT 2, 'Jane Smith', 30, 'Manager', current_date, current_date, hour(current_timestamp)
    UNION ALL  
    SELECT 3, 'Bob Johnson', 35, 'Analyst', current_date, current_date, hour(current_timestamp);

  3. Create a sample Parquet table (table2) and add data to be used for testing the migrate and register_table approach. Replace <amzn-s3-demo-bucket> with your bucket name.
    CREATE TABLE table2
    WITH (
      external_location = 's3://<amzn-s3-demo-bucket>/table2/',
      format = 'PARQUET',
      partitioned_by = ARRAY['date', 'hour']
    )
    AS
    SELECT 
      1 as id,
      'John Doe' as name,
      25 as age,
      'Engineer' as job_title,
      current_date as created_date,
      current_date as date,
      hour(current_timestamp) as hour
    UNION ALL
    SELECT 2, 'Jane Smith', 30, 'Manager', current_date, current_date, hour(current_timestamp)
    UNION ALL  
    SELECT 3, 'Bob Johnson', 35, 'Analyst', current_date, current_date, hour(current_timestamp);

  4. Drop the tables from the Data Catalog because you only need Parquet data with the Hive-style partitioning structure.
    DROP TABLE IF EXISTS test_db.table1

Create a DynamoDB control table

Before beginning the migration process, you must create a DynamoDB table that serves as the control plane. This table maps source Amazon S3 paths to their corresponding Iceberg database and table destinations, enabling systematic tracking of the migration process.

To implement this control mechanism, create a table with the following structure:

  • A primary key s3_path that stores the source Parquet data location
  • Two attributes that define the target Iceberg location:
    • target_db_name
    • target_table_name

To create the DynamoDB control table

  1. Create the Amazon DynamoDB table using the following AWS CLI command:
    aws dynamodb create-table \
    --table-name migration-control-table \
    --attribute-definitions \
    AttributeName=s3_path,AttributeType=S \
    --key-schema \
    AttributeName=s3_path,KeyType=HASH \
    --billing-mode PAY_PER_REQUEST \
    --region <REGION>

  2. Verify the table is created successfully. Replace <REGION> with the AWS Region where your data is stored:
    aws dynamodb describe-table --table-name migration-control-table --region <REGION>

  3. Create a migration_data.json file with the following contents.
    In this example:

    • Replace <amzn-s3-demo-bucket> and <TablePrefix>with the name of your S3 bucket and prefix containing the Parquet data
    • Replace <DatabaseName> with the name of your target Iceberg database
    • Replace <TableName> with the name of your target Iceberg table
    {
        "your-migration-table": [
            {
                "PutRequest": {
                    "Item": {
                        "s3_path": {"S": "s3://<amzn-s3-demo-bucket>/table1/"},
                        "target_db_name": {"S": "test_db"},
                        "target_table_name": {"S": "table1"}
                    }
                }
            },
            {
                "PutRequest": {
                    "Item": {
                        "s3_path": {"S": "s3://<amzn-s3-demo-bucket>/table2/"},
                        "target_db_name": {"S": "test_db"},
                        "target_table_name": {"S": "table2"}
                    }
                }
            },
            {
                "PutRequest": {
                    "Item": {
                        "s3_path": {"S": "s3://<amzn-s3-demo-bucket>/<TablePrefix>/"},
                        "target_db_name": {"S": "<DatabaseName>"},
                        "target_table_name": {"S": "<TableName>"}
                    }
                }
            }
        ]
    }

    This file defines the mapping between Amazon S3 paths and their corresponding Iceberg table destinations.

  4. Run the following CLI command to load the DynamoDB control table.
    aws dynamodb batch-write-item \
    --request-items file://migration_data.json \
    --region <REGION;>

Migration methodologies

In this section, you explore two methodologies for migrating your existing Parquet tables to Apache Iceberg format:

  • Apache Iceberg migrate and register_table approach – This approach first converts your Parquet table to Iceberg format using the native migrate procedure, followed by registering it in AWS Glue using the register_table procedure.
  • Apache Iceberg add_files approach – This method creates an empty Iceberg table and uses the add_files procedure to import existing Parquet data files without physically moving them.

Apache Iceberg migrate and register_table procedure

Use the Apache Iceberg Migrate procedure that is used for in-place conversion of an existing Hive or Parquet table into an Iceberg-managed table. Thereafter, you can use the Apache Iceberg RegisterTable procedure to register the respective table in AWS Glue.

AWS workflow diagram showing DynamoDB to Apache Iceberg migration using Amazon EMR with Hive Metastore for migration and Glue Metastore for registration, displaying configuration tables at each stage.

Migrate

  1. In your EMR cluster with Hive as the metastore, create a PySpark session with the following Iceberg Packages:
    pyspark \
    --name "Iceberg Migration" \
    --conf "spark.jars=/usr/share/aws/iceberg/lib/iceberg-spark3-runtime.jar" \
    --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
    --conf spark.sql.catalog.spark_catalog=org.apache.iceberg.spark.SparkSessionCatalog \
    --conf spark.sql.catalog.spark_catalog.type=hive

    This post uses Iceberg v1.9.1 (Amazon EMR build), which is native to Amazon EMR 7.11. Always verify the latest supported version and update package coordinates accordingly.

  2. Next, create your corresponding table in your Hive catalog (you can skip this step if you already have tables created in your hive catalog). Replace <amzn-s3-demo-bucket> with the name of your S3 bucket.
    In the following snippet, change or remove the PARTITIONED BY command based on the partition strategy of your table, the MSCK Repair table command should only be run if your respective table is partitioned.

    #You can automate this for production Scaling with DynamoDB as control table 
    s3_path = "s3://<amzn-s3-demo-bucket>/table1/"
    target_db_name = "test_db"
    target_table_name = "table1"
    # Read data in a dataframe to infer schema
    df = spark.read.parquet(s3_path)
    df.createOrReplaceTempView("temp_view")
    # Get schema as string
    schema = spark.table("temp_view").schema
    schema_string = ", ".join([f"{field.name} {field.dataType.simpleString()}" for field in schema])
    # Create Database If not exists 
    spark.sql(f"CREATE DATABASE IF NOT EXISTS {target_db_name}").show()
    # full_table_name= test_db.table1
    full_table_name = f"{target_db_name}.{target_table_name}"
    # Create table
    spark.sql(f"""
    CREATE TABLE IF NOT EXISTS {full_table_name} (
        {schema_string}
    )
    STORED AS PARQUET
    PARTITIONED BY (date, hour)
    LOCATION '{s3_path}'
    """)
    # Refresh, repair, and validate
    spark.sql(f"REFRESH TABLE {full_table_name}")
    spark.sql(f"MSCK REPAIR TABLE {full_table_name}")

  3. Convert the Parquet table to an Iceberg table in Hive
    # Run migration procedure
    spark.sql(f"CALL spark_catalog.system.migrate('{full_table_name}')")
    # Validate that the table is successfully migrated 
    spark.sql(f"DESCRIBE FORMATTED {full_table_name}").show(truncate=False)

Run the migrate command to convert the Parquet-based table to an Iceberg table, creating the metadata folder and the metadata.json file therein

You can stop at this point if you don’t intend to migrate your existing iceberg table from Hive to the Data Catalog.

Register

  1. Sign in to the AWS Glue as Spark Catalog enabled EMR cluster.
  2. Register the Iceberg table to your Data Catalog.

    Create the session with the respective Iceberg Packages. Replace <amzn-s3-demo-bucket> with your bucket name, and <warehouse> with warehouse directory.

    pyspark \
    --conf "spark.jars=/usr/share/aws/iceberg/lib/iceberg-spark3-runtime.jar" \
    --conf "spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions" \
    --conf "spark.sql.catalog.glue_catalog=org.apache.iceberg.spark.SparkCatalog" \
    --conf "spark.sql.catalog.glue_catalog.warehouse= s3://<amzn-s3-demo-bucket>/<warehouse>/"  \
    --conf "spark.sql.catalog.glue_catalog.catalog-impl=org.apache.iceberg.aws.glue.GlueCatalog" \
    --conf "spark.sql.catalog.glue_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO"

  3. Run the register_table command to make the Iceberg table visible in AWS Glue.
    • register_table registers an existing Iceberg table’s metadata file (metadata.json) with a catalog(glue_catalog) so that Spark (and other engines) can query it.
    • The procedure creates a Data Catalog entry for the table, pointing it to the given metadata location.

    Replace <amzn-s3-demo-bucket> and <metadata-prefix> with the name of your S3 bucket and metadata prefix name.

    Ensure that your EMR Spark Cluster has been configured with appropriate AWS Glue permissions

    # You can automate this for production Scaling with DynamoDB as control table
    metadata_location = "s3://<amzn-s3-demo-bucket>/table1/metadata/<metadata-prefix>.metadata.json"
    target_db_name = "test_db"
    target_table_name = "table1"
    full_table_name = f"{target_db_name}.{target_table_name}"
    # Register existing Iceberg table metadata in Glue Catalog
    spark.sql(f"CALL glue_catalog.system.register_table('{full_table_name}', '{metadata_location}')")
    # Set table properties (example: Iceberg format version 2)
    spark.sql(f"ALTER TABLE glue_catalog.{full_table_name} SET TBLPROPERTIES('format-version'='2')")

  4. Validate that the Iceberg table is now visible in the Data Catalog.
    # Lookout for format as iceberg/parquet
    spark.sql("SHOW TBLPROPERTIES glue_catalog.test_db.table1").show()

Apache Iceberg’s add_files procedure

AWS workflow diagram showing DynamoDB to Apache Iceberg migration using AWS Glue Add_Files procedure, displaying input configuration and output status tables with metadata location and registration confirmation.

Here, you’re going to use Iceberg’s add_files procedure to import raw data files (Parquet, ORC, Avro) into an existing Iceberg table by updating its metadata. This procedure works for both Hive and Data Catalog, it doesn’t physically move or rewrite the files—it only registers them so Iceberg can manage them.

This methodology comprises the following steps:

  1. Create an empty Iceberg table in AWS Glue.
    Because the add_files procedure expects the iceberg table to be already present, you need to create an empty Iceberg table by inferring the table schema.
  2. Register existing data locations to the Iceberg table

Using the add_files procedure in a Glue-backed Iceberg catalog will register the target S3 path along with all its subdirectories to the empty Iceberg table created in the previous step.

You can consolidate both steps into a single Spark job. For the following AWS Glue job, you have specified iceberg as a value for the --datalake-formats job parameter. See the AWS Glue job configuration documentation for more details.

Replace <amzn-s3-demo-bucket> with your S3 bucket name and <warehouse> with warehouse directory.

from pyspark.sql import SparkSession
import logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
target_db_name = "test_db"
target_table_name = "table2"
s3_path = "s3://<amzn-s3-demo-bucket>/table2"
# Set to None or [] for unpartitioned
partitioned_cols = ["date", "hour"]  
spark = SparkSession.builder \
    .appName("Iceberg Add Files") \
    .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
    .config("spark.sql.catalog.glue_catalog", "org.apache.iceberg.spark.SparkCatalog") \
    .config("spark.sql.catalog.glue_catalog.catalog-impl", "org.apache.iceberg.aws.glue.GlueCatalog") \
    .config("spark.sql.catalog.glue_catalog.io-impl", "org.apache.iceberg.aws.s3.S3FileIO") \
    .config("spark.sql.catalog.glue_catalog.warehouse", "s3://<amzn-s3-demo-bucket>/<warehouse>/") \
    .getOrCreate()
full_table_name = f"glue_catalog.{target_db_name}.{target_table_name}"
# Read schema from one file (schema inference)
df = spark.read.parquet(s3_path)
schema = df.schema
# Create empty Iceberg table
empty_df = spark.createDataFrame([], schema)
if partitioned_cols:
    empty_df.writeTo(full_table_name).using("iceberg").partitionedBy(*partitioned_cols). .tableProperty("format-version", "2").create()
else:
    empty_df.writeTo(full_table_name).using("iceberg").tableProperty("format-version", "2").create()
logger.info(f"Created empty Iceberg table: {full_table_name}")
spark.sql(f"""
CALL glue_catalog.system.add_files(
  '{target_db_name}.{target_table_name}',
  'parquet.`{s3_path}`'
)
""")

When working with non-Hive partitioned datasets, a direct migration to Apache Iceberg using add_files might not behave as expected. See Appendix C for more information.

Considerations

Let’s explore two key considerations that you should address when implementing your migration strategy.

State management using DynamoDB control table

Use the following sample code snippet to update the state of DynamoDB table:

def update_dynamodb_record(self, s3_path, metadata_loc=None, error_msg=None):
    # Get current error message
    try:
        response = self.dynamodb.get_item(
            TableName='migration-control-table',
            Key={'s3_path': {'S': s3_path}}
        )
        current_error = response.get('Item', {}).get('error_message', {}).get('S', '')
    except:
        current_error = ""
    if error_msg:
        # Error case
        error_msg = (error_msg or "Unknown error")[:1000]
        update_expr = "SET error_message = :err"
        attr_values = {':err': {'S': error_msg}}
        if current_error:
            update_expr += ", prev_error_message = :prev"
            attr_values[':prev'] = {'S': current_error}
        update_kwargs = {'TableName': 'Iceberg_migration','Key': {'s3_path': {'S': s3_path}},'UpdateExpression': update_expr,'ExpressionAttributeValues': attr_values}
        self.logger.error(f"Set error for {s3_path}: {error_msg}")
    else:
        # Success case
        update_kwargs = {
            'TableName': 'Iceberg_migration',
            'Key': {'s3_path': {'S': s3_path}},
            'UpdateExpression': 'SET #s = :status, #m = :meta, #p = :prev, #e = :err',
            'ExpressionAttributeNames': {'#s': 'status','#p': 'prev_error_message','#e': 'error_message','#m': 'metadata_location'
            },
            'ExpressionAttributeValues': {
                ':status': {'S': 'Iceberg_Metadata_Populated and Registered'},
                ':prev': {'S': current_error},
                ':err': {'S': ''},
                ':meta': {'S': metadata_loc}
            }
        }
        self.logger.info(f"Updated DynamoDB status for {s3_path}: {metadata_loc}")

This ensures that any errors are logged and saved to DynamoDB as error_message. On successive retries, previous errors move to prev_error_message and new errors overwrite error_message. Successful operations clear error_message and archive the last error.

Protecting your data from unintended deletion

To protect your data from unintended deletion, never delete data or metadata files from Amazon S3 directly. Iceberg tables that are registered in AWS Glue or Athena are managed tables and should be deleted using the DROP TABLE command from Spark or Athena. The DROP TABLE command deletes both the table metadata and the underlying data files in S3. See Appendix D for more information.

Clean up

Complete the following steps to clean up your resources:

  1. Delete the DynamoDB control table
  2. Delete the database and tables
  3. Delete the EMR clusters and AWS Glue job used for testing

Conclusion

In this post, we showed you how to modernize your Parquet-based data lake into an Apache Iceberg–powered lakehouse without rewriting or duplicating data. You learned two complementary approaches for this in-place migration:

  • Migrate and register – Ideal for converting existing Hive-registered Parquet tables into Iceberg-managed tables.
  • add_files – Best suited for quickly onboarding raw Parquet data into Iceberg without rewriting files.

Both approaches benefit from DynamoDB centralized state tracking, which enables retries, error auditing, and lifecycle management across multiple datasets.

By combining Apache Iceberg with Amazon EMR, AWS Glue, and Amazon DynamoDB, you can create a production-ready migration pipeline that is observable, automated, and straightforward to extend to future data format upgrades. This pattern forms a solid foundation for building an Iceberg-based lakehouse on AWS, helping you achieve faster analytics, better data governance, and long-term flexibility for evolving workloads.

To get started, try implementing this solution using the sample tables (table1 and table2) that you created using Athena queries. we encourage you to share your migration experiences and questions in the comments.


Appendix A — Creating an EMR cluster for Hive metastore using console and AWS CLI

Console steps:

  1. Open AWS Management Console for Amazon EMR and choose Create cluster.
  2. Select Spark or Hive under applications.
  3. Under AWS Glue Data Catalog settings, make sure the following options are not selected:
    • Use for Hive table metadata
    • Use for Spark table metadata
  4. Configure SSH access (KeyName).
  5. Configure network (VPC, subnets, SGs) to allow access to S3.

AWS CLI steps:

aws emr create-cluster \
  --region us-east-1 \
  --name "IcebergHiveCluster711" \
  --release-label emr-7.11.0 \
  --applications Name=Hive Name=Spark Name=Hadoop \
  --ec2-attributes '{"KeyName":"<key-pair>","SubnetId":"<subnet-id>"}'  \
  --instance-groups '[
    {
      "Name":"Master",
      "InstanceGroupType":"MASTER",
      "InstanceType":"m5.xlarge",
      "InstanceCount":1
    },
    {
      "Name":"Workers",
      "InstanceGroupType":"CORE",
      "InstanceType":"m5.xlarge",
      "InstanceCount":2
    }
  ]' \
  --use-default-roles

Appendix B — EMR cluster with AWS Glue as Spark Metastore

Console steps:

  1. Open the Amazon EMR console, choose Create cluster and then select EMR Serverless or provisioned EMR.
  2. Under Software Configuration, verify that Spark is installed.
  3. Under AWS Glue Data Catalog settings, select Use Glue Data Catalog for Spark metadata.
  4. Configure SSH access (KeyName).
  5. Configure network settings (VPC, subnets, and security groups) to allow access to Amazon S3 and AWS Glue.

AWS CLI (provisioned Amazon EMR):

aws emr create-cluster \
  --region us-east-1 \
  --name "IcebergGlueCluster711" \
  --release-label emr-7.11.0 \
  --applications Name=Spark Name=Hadoop \
  --ec2-attributes '{"KeyName":"<key-pair>","SubnetId":"<subnet-id>"}' \
  --instance-groups '[
    {
      "Name":"Master",
      "InstanceGroupType":"MASTER",
      "InstanceType":"m5.xlarge",
      "InstanceCount":1
    },
    {
      "Name":"Workers",
      "InstanceGroupType":"CORE",
      "InstanceType":"m5.xlarge",
      "InstanceCount":2
    }
  ]' \
 --configurations '[{"Classification":"spark-hive-site","Properties":{"hive.metastore.client.factory.class":"com.amazonaws.glue.catalog.metastore.AWSGlueDataCatalogHiveClientFactory"}}]' \
 --use-default-roles

Appendix C — Non-Hive partitioned datasets and Iceberg add_files

This appendix explains why a direct in-place migration using an add_files-style procedure might not behave as expected for datasets that aren’t Hive-partitioned and shows recommended fixes and examples.

AWS Glue and Athena follow Hive-style partitioning, where partition column values are encoded in the S3 path rather than inside the data files. For example, following the Parquet dataset created in the Create Sample Parquet Dataset as a source section of this post:

s3://amzn-s3-demo-bucket/events/event_date=2024-09-01/hour=5/part-0000.parquet
s3://amzn-s3-demo-bucket/events/event_date=2024-09-02/hour=5/part-0001.parquet
  • Partition columns (event_date, hour) are represented in the folder structure.
  • Non-partition columns (for example, id, name, age) remain inside the Parquet files.
  • Iceberg add_files can correctly map partitions based on the folder path, even if partition columns are missing from the Parquet file itself.

Partition column

Stored in path

Stored in file

Athena or AWS Glue and Iceberg behavior

event_date Yes Yes Partitions inferred correctly
hour Yes No Partitions still inferred from path

Non-Hive partitioning layout (problem case)

s3://amzn-s3-demo-bucket/events/date/part-0000.parquet
s3://amzn-s3-demo-bucket/events/date/part-0001.parquet
  • No partition columns in the path.
  • File might not contain partition columns.

If you try to create an empty Iceberg table and directly load it using add_files on a non-hive layout, the following happens:

  1. Iceberg cannot automatically map partitions, add_files operations fail or register files with incorrect or missing partition metadata.
  2. Queries in Athena or AWS Glue will return unexpected NULLs or incomplete results.
  3. Successive incremental writes using add_files will fail.

Recommended approaches:

Create an AWS Glue table and use the Iceberg snapshot procedure:

  1. Create a table in AWS Glue pointing to your existing Parquet dataset.

You might need to manually provide the schema because glue crawler might fail to automatically infer it for you.

  1. Use Iceberg’ s snapshot procedure to convert and move the AWS Glue table into your target Iceberg table.

This works because Iceberg relies on AWS Glue for schema inference, so this approach ensures correct mapping of columns and partitions without rewriting the data. For more information, see Snapshot procedure.


Appendix D — Understanding table types: Managed compared to external

By default, all non-Iceberg tables created in AWS Glue or Athena are external tables, Athena doesn’t manage the underlying data. If you use CREATE TABLE without the EXTERNAL keyword for non-Iceberg tables, Athena issues an error.

However, when dealing with Iceberg tables, AWS Glue and Athena also manage the underlying data for the respective tables, so these tables are treated as internal tables.

Running DROP TABLE on Iceberg tables will delete the table and the underlying data.

The following table describes how the effect of DELETE and DROP TABLE actions on Iceberg tables in AWS Glue and Athena:

Operation What it does Effect on S3 data
DELETE FROM mydb.products_iceberg WHERE date = 2025-10-06; Creates new snapshot, hides deleted rows Data files stay until cleanup
DROP TABLE test_db.table1; Deletes table and all data Files are permanently removed

About the authors

Mihir Borkar

Mihir Borkar

Mihir is a seasoned AWS Data Architect with nearly a decade of experience designing and implementing enterprise-scale data solutions on AWS. He specializes in modernizing data architectures using AWS data analytical services, designing scalable data lakes and analytics platforms with a focus on efficient, cost-effective solutions. In his free time, Mihir loves to read about emerging cloud technologies and explore latest developments in AI/ML.

Amit Maindola

Amit Maindola

Amit is a Senior Data Architect with AWS ProServe team focused on data engineering, analytics, and AI/ML at Amazon Web Services. He helps customers in their digital transformation journey and enables them to build highly scalable, robust, and secure cloud-based analytical solutions on AWS to gain timely insights and make critical business decisions.

Arghya Banerjee

Arghya Banerjee

Arghya is a Sr. Solutions Architect at AWS in the San Francisco Bay Area, focused on helping customers adopt and use the AWS Cloud. He is focused on big data, data lakes, streaming and batch analytics services, and generative AI technologies.

Using Amazon EMR DeltaStreamer to stream data to multiple Apache Hudi tables

Post Syndicated from Gautam Bhaghavatula original https://aws.amazon.com/blogs/big-data/using-amazon-emr-deltastreamer-to-stream-data-to-multiple-apache-hudi-tables/

In this post, we show you how to implement real-time data ingestion from multiple Kafka topics to Apache Hudi tables using Amazon EMR. This solution streamlines data ingestion by processing multiple Amazon Managed Streaming for Apache Kafka (Amazon MSK) topics in parallel while providing data quality and scalability through change data capture (CDC) and Apache Hudi.

Organizations processing real-time data changes across multiple sources often struggle with maintaining data consistency and managing resource costs. Traditional batch processing requires reprocessing entire datasets, leading to high resource usage and delayed analytics. By implementing CDC with Apache Hudi’s MultiTable DeltaStreamer, you can achieve real-time updates; efficient incremental processing with atomicity, consistency, isolation, durability (ACID) guarantees; and seamless schema evolution while minimizing storage and compute costs.

Using Amazon Simple Storage Service (Amazon S3), Amazon CloudWatch, Amazon EMR, Amazon MSK and AWS Glue Data Catalog, you’ll build a production-ready data pipeline that processes changes from multiple data sources simultaneously. Through this tutorial, you’ll learn to configure CDC pipelines, manage table-specific configurations, implement 15-minute sync intervals, and maintain your streaming pipeline. The result is a robust system that maintains data consistency while enabling real-time analytics and efficient resource utilization.

What is CDC?

Imagine a constantly evolving data stream, a river of information where updates flow continuously. CDC acts like a sophisticated net, capturing only the modifications—the inserts, updates, and deletes—happening within that data stream. Through this targeted approach, you can focus on the new and changed data, significantly improving the efficiency of your data pipelines.There are numerous advantages to embracing CDC:

  • Reduced processing time – Why reprocess the entire dataset when you can focus only on the updates? CDC minimizes processing overhead, saving valuable time and resources.
  • Real-time insights – With CDC, your data pipelines become more responsive. You can react to changes almost instantaneously, enabling real-time analytics and decision-making.
  • Simplified data pipelines – Traditional batch processing can lead to complex pipelines. CDC streamlines the process, making data pipelines more manageable and easier to maintain.

Why Apache Hudi?

Hudi simplifies incremental data processing and data pipeline development. This framework efficiently manages business requirements such as data lifecycle and improves data quality. You can use Hudi to manage data at the record-level in Amazon S3 data lakes to simplify CDC and streaming data ingestion and handle data privacy use cases requiring record-level updates and deletes. Datasets managed by Hudi are stored in Amazon S3 using open storage formats, while integrations with Presto, Apache Hive, Apache Spark, and Data Catalog give you near real time access to updated data. Apache Hudi facilitates incremental data processing for Amazon S3 by:

  • Managing record-level changes – Ideal for update and delete use cases
  • Open formats – Integrates with Presto, Hive, Spark, and Data Catalog
  • Schema evolution – Supports dynamic schema changes
  • HoodieMultiTableDeltaStreamer – Simplifies ingestion into multiple tables using centralized configurations

Hudi MultiTable Delta Streamer

The HoodieMultiTableStreamer offers a streamlined approach to data ingestion from multiple sources into Hudi tables. By processing multiple sources simultaneously through a single DeltaStreamer job, it eliminates the need for separate pipelines while reducing operational complexity. The framework provides flexible configuration options, and you can tailor settings for diverse formats and schemas across different data sources.

One of its key strengths lies in unified data delivery, organizing information in respective Hudi tables for seamless access. The system’s intelligent upsert capabilities efficiently handle both inserts and updates, maintaining data consistency across your pipeline. Additionally, its robust schema evolution support enables your data pipeline to adapt to changing business requirements without disruption, making it an ideal solution for dynamic data environments.

Solution overview

In this section, we show how to stream data to Apache Hudi Table using Amazon MSK. For this example scenario, there are data streams from three distinct sources residing in separate Kafka topics. We aim to implement a streaming pipeline that uses the Hudi DeltaStreamer with multitable support to ingest and process this data at 15-minute intervals.

Mechanism

Using MSK Connect, data from multiple sources flows into MSK topics. These topics are then ingested into Hudi tables using the Hudi MultiTable DeltaStreamer. In this sample implementation, we create three Amazon MSK topics and configure the pipeline to process data in JSON format using JsonKafkaSource, with the flexibility to handle Avro format when needed through the appropriate deserializer configuration

The following diagram illustrates how our solution processes data from multiple source databases through Amazon MSK and Apache Hudi to enable analytics in Amazon Athena. Source databases send their data changes—including inserts, updates, and deletes—to dedicated topics in Amazon MSK, where each data source maintains its own Kafka topic for change events. An Amazon EMR cluster runs the Apache Hudi MultiTable DeltaStreamer, which processes these multiple Kafka topics in parallel, transforming the data and writing it to Apache Hudi tables stored in Amazon S3. Data Catalog maintains the metadata for these tables, enabling seamless integration with analytics tools. Finally, Amazon Athena provides SQL query capabilities on the Hudi tables, allowing analysts to run both snapshot and incremental queries on the latest data. This architecture scales horizontally as new data sources are added, with each source getting its dedicated Kafka topic and Hudi table configuration, while maintaining data consistency and ACID guarantees across the entire pipeline.

To set up the solution, you need to complete the following high-level steps:

  1. Set up Amazon MSK and create Kafka topics
  2. Create the Kafka topics
  3. Create table-specific configurations
  4. Launch Amazon EMR cluster
  5. Invoke the Hudi MultiTable DeltaStreamer
  6. Verify and query data

Prerequisites

To perform the solution, you need to have the following prerequisites. For AWS services and permissions, you need:

  • AWS account:
  • IAM roles:
    • Amazon EMR service role (EMR_DefaultRole) with permissions for Amazon S3, AWS Glue and CloudWatch.
    • Amazon EC2 instance profile (EMR_EC2_DefaultRole) with S3 read/write access.
    • Amazon MSK access role with appropriate permissions.
  • S3 buckets:
    • Configuration bucket for storing properties files and schemas.
    • Output bucket for Hudi tables.
    • Logging bucket (optional but recommended).
  • Network configuration:
  • Development tools:

Set up Amazon MSK and create Kafka topics

In this step, you’ll create an MSK cluster and configure the required Kafka topics for your data streams.

  1. To create an MSK cluster:
aws kafka create-cluster \
    --cluster-name hudi-msk-cluster \
    --broker-node-group-info file://broker-nodes.json \
    --kafka-version "2.8.1" \
    --number-of-broker-nodes 3 \
    --encryption-info file://encryption-info.json \
    --client-authentication file://client-authentication.json
  1. Verify the cluster status:

aws kafka describe-cluster --cluster-arn $CLUSTER_ARN | jq '.ClusterInfo.State'

The command should return ACTIVE when the cluster is ready.

Schema setup

To set up the schema, complete the following steps:

  1. Create your schema files.
    1. input_schema.avsc:
      {
          "type": "record",
          "name": "CustomerSales",
          "fields": [
              {"name": "Id", "type": "string"},
              {"name": "ts", "type": "long"},
              {"name": "amount", "type": "double"},
              {"name": "customer_id", "type": "string"},
              {"name": "transaction_date", "type": "string"}
          ]
      }

    2. output_schema.avsc:
      {
          "type": "record",
          "name": "CustomerSalesProcessed",
          "fields": [
              {"name": "Id", "type": "string"},
              {"name": "ts", "type": "long"},
              {"name": "amount", "type": "double"},
              {"name": "customer_id", "type": "string"},
              {"name": "transaction_date", "type": "string"},
              {"name": "processing_timestamp", "type": "string"}
          ]
      }

  2. Create and upload schemas to your S3 bucket:
    # Create the schema directory
    aws s3 mb s3://hudi-config-bucket-$AWS_ACCOUNT_ID
    aws s3api put-object --bucket hudi-config-bucket-$AWS_ACCOUNT_ID --key HudiProperties/
    # Upload schema files
    aws s3 cp input_schema.avsc s3://hudi-config-bucket-$AWS_ACCOUNT_ID/HudiProperties/
    aws s3 cp output_schema.avsc s3://hudi-config-bucket-$AWS_ACCOUNT_ID/HudiProperties/

Create the Kafka topics

To create the Kafka topics, complete the following steps:

  1. Get the bootstrap broker string:
    # Get bootstrap brokers
    BOOTSTRAP_BROKERS=$(aws kafka get-bootstrap-brokers --cluster-arn $CLUSTER_ARN --query 'BootstrapBrokerString' --output text)

  2. Create the required topics:
    kafka-topics.sh --create \
        --bootstrap-server $BOOTSTRAP_BROKERS \
        --replication-factor 3 \
        --partitions 3 \
        --topic cust_sales_details
    kafka-topics.sh --create \
        --bootstrap-server $BOOTSTRAP_BROKERS \
        --replication-factor 3 \
        --partitions 3 \
        --topic cust_sales_appointment
    kafka-topics.sh --create \
        --bootstrap-server $BOOTSTRAP_BROKERS \
        --replication-factor 3 \
        --partitions 3 \
        --topic cust_info

Configure Apache Hudi

The Hudi MultiTable DeltaStreamer configuration is divided into two major components to streamline and standardize data ingestion:

  • Common configurations – These settings apply across all tables and define the shared properties for ingestion. They include details such as shuffle parallelism, Kafka brokers, and common ingestion configurations for all topics.
  • Table-specific configurations – Each table has unique requirements, such as the record key, schema file paths, and topic names. These configurations tailor each table’s ingestion process to its schema and data structure.

Create common configuration file

Common Config: kafka-hudi config file where we specify kafka broker and common configuration for all topics as below

Create the kafka-hudi-deltastreamer.properties file with the following properties:

# Common parallelism settings
hoodie.upsert.shuffle.parallelism=2
hoodie.insert.shuffle.parallelism=2
hoodie.delete.shuffle.parallelism=2
hoodie.bulkinsert.shuffle.parallelism=2
# Table ingestion configuration
hoodie.deltastreamer.ingestion.tablesToBeIngested=hudi_sales_tables.cust_sales_details,hudi_sales_tables.cust_sales_appointment,hudi_sales_tables.cust_info
# Table-specific config files
hoodie.deltastreamer.ingestion.hudi_sales_tables.cust_sales_details.configFile=s3://hudi-config-bucket-$AWS_ACCOUNT_ID/HudiProperties/tableProperties/cust_sales_details.properties
hoodie.deltastreamer.ingestion.hudi_sales_tables.cust_sales_appointment.configFile=s3://hudi-config-bucket-$AWS_ACCOUNT_ID/HudiProperties/tableProperties/cust_sales_appointment.properties
hoodie.deltastreamer.ingestion.hudi_sales_tables.cust_info.configFile=s3://hudi-config-bucket-$AWS_ACCOUNT_ID/HudiProperties/tableProperties/cust_info.properties
# Source configuration
hoodie.deltastreamer.source.dfs.root=s3://hudi-config-bucket-$AWS_ACCOUNT_ID/HudiProperties/
# MSK configuration
bootstrap.servers=BOOTSTRAP_BROKERS_PLACEHOLDER
auto.offset.reset=earliest
group.id=hudi_delta_streamer
# Security configuration
hoodie.sensitive.config.keys=ssl,tls,sasl,auth,credentials
sasl.mechanism=PLAIN
security.protocol=SASL_SSL
ssl.endpoint.identification.algorithm=
# Deserializer
hoodie.deltastreamer.source.kafka.value.deserializer.class=io.confluent.kafka.serializers.KafkaAvroDeserializer

Create table-specific configurations

For each topic, create its own configuration with a topic name and primary key details. Complete the following steps:

  1. cust_sales_details.properties:
    # Table: cust sales
    hoodie.datasource.write.recordkey.field=Id
    hoodie.deltastreamer.source.kafka.topic=cust_sales_details
    hoodie.deltastreamer.keygen.timebased.timestamp.type=UNIX_TIMESTAMP
    hoodie.deltastreamer.keygen.timebased.input.dateformat=yyyy-MM-dd HH:mm:ss.S
    hoodie.streamer.schemaprovider.registry.schemaconverter=
    hoodie.datasource.write.precombine.field=ts

  2. cust_sales_appointment.properties:
    # Table: cust sales appointment
    hoodie.datasource.write.recordkey.field=Id
    hoodie.deltastreamer.source.kafka.topic=cust_sales_appointment
    hoodie.deltastreamer.keygen.timebased.timestamp.type=UNIX_TIMESTAMP
    hoodie.deltastreamer.keygen.timebased.input.dateformat=yyyy-MM-dd HH:mm:ss.S hoodie.streamer.schemaprovider.registry.schemaconverter=
    hoodie.datasource.write.precombine.field=ts

  3. cust_info.properties:
    # Table: cust info
    hoodie.datasource.write.recordkey.field=Id
    hoodie.deltastreamer.source.kafka.topic=cust_info
    hoodie.deltastreamer.keygen.timebased.timestamp.type=UNIX_TIMESTAMP
    hoodie.deltastreamer.keygen.timebased.input.dateformat= yyyy-MM-dd HH:mm:ss.S
    hoodie.streamer.schemaprovider.registry.schemaconverter=
    hoodie.datasource.write.precombine.field=ts
    hoodie.deltastreamer.schemaprovider.source.schema.file=-$AWS_ACCOUNT_ID/HudiProperties/input_schema.avsc
    hoodie.deltastreamer.schemaprovider.target.schema.file=-$AWS_ACCOUNT_ID/HudiProperties/output_schema.avsc

These configurations form the backbone of Hudi’s ingestion pipeline, enabling efficient data handling and maintaining real-time consistency. Schema configurations define the structure of both source and target data, maintaining seamless data transformation and ingestion. Operational settings control how data is uniquely identified, updated, and processed incrementally.

The following are critical details for setting up Hudi ingestion pipelines:

  • hoodie.deltastreamer.schemaprovider.source.schema.file – The schema of the source record
  • hoodie.deltastreamer.schemaprovider.target.schema.file – The schema for the target record
  • hoodie.deltastreamer.source.kafka.topic – The source MSK topic name
  • bootstap.servers – The Amazon MSK bootstrap server’s private endpoint
  • auto.offset.reset – The consumer’s behavior when there is no committed position or when an offset is out of range

Key operational fields to achieve in-place updates for the generated schema include:

  • hoodie.datasource.write.recordkey.field – The record key field. This is the unique identifier of a record in Hudi.
  • hoodie.datasource.write.precombine.field – When two records have the same record key value, Apache Hudi picks the one with the largest value for the pre-combined field.
  • hoodie.datasource.write.operation – The operation on the Hudi dataset. Possible values include UPSERT, INSERT, and BULK_INSERT.

Launch Amazon EMR cluster

This step creates an EMR cluster with Apache Hudi installed. The cluster will run the MultiTable DeltaStreamer to process data from your Kafka topics. To create the EMR cluster, enter the following:

# Create EMR cluster with Hudi installed
aws emr create-cluster \
    --name "Hudi-CDC-Cluster" \
    --release-label emr-6.15.0 \
    --applications Name=Hadoop Name=Spark Name=Hive Name=Livy \
    --ec2-attributes KeyName=myKey,SubnetId=$SUBNET_ID,InstanceProfile=EMR_EC2_InstanceProfile \
    --service-role EMR_ServiceRole \
    --instance-groups InstanceGroupType=MASTER,InstanceCount=1,InstanceType=m5.xlarge InstanceGroupType=CORE,InstanceCount=2,InstanceType=m5.xlarge \
    --configurations file://emr-configurations.json \
    --bootstrap-actions Name="Install Hudi",Path="s3://hudi-config-bucket-$AWS_ACCOUNT_ID/bootstrap-hudi.sh"

Invoke the Hudi MultiTable DeltaStreamer

This step configures and starts the DeltaStreamer job that will continuously process data from your Kafka topics into Hudi tables. Complete the following steps:

  1. Connect to the Amazon EMR master node:
    # Get master node public DNS
    MASTER_DNS=$(aws emr describe-cluster --cluster-id $CLUSTER_ID --query 'Cluster.MasterPublicDnsName' --output text)
    
    # SSH to master node
    ssh -i myKey.pem hadoop@$MASTER_DNS

  2. Execute the DeltaStreamer job:
    # 
    spark-submit --deploy-mode client \
      --conf "spark.serializer=org.apache.spark.serializer.KryoSerializer" \
      --conf "spark.sql.catalog.spark_catalog=org.apache.spark.sql.hudi.catalog.HoodieCatalog" \
      --conf "spark.sql.extensions=org.apache.spark.sql.hudi.HoodieSparkSessionExtension" \
      --jars "/usr/lib/hudi/hudi-utilities-bundle_2.12-0.14.0-amzn-0.jar,/usr/lib/hudi/hudi-spark-bundle.jar" \
      --class "org.apache.hudi.utilities.deltastreamer.HoodieMultiTableDeltaStreamer" \
      /usr/lib/hudi/hudi-utilities-bundle_2.12-0.14.0-amzn-0.jar \
      --props s3://hudi-config-bucket-$AWS_ACCOUNT_ID/HudiProperties/kafka-hudi-deltastreamer.properties \
      --config-folder s3://hudi-config-bucket-$AWS_ACCOUNT_ID/HudiProperties/tableProperties/ \
      --table-type MERGE_ON_READ \
      --base-path-prefix s3://hudi-data-bucket-$AWS_ACCOUNT_ID/hudi/ \
      --source-class org.apache.hudi.utilities.sources.JsonKafkaSource \
      --schemaprovider-class org.apache.hudi.utilities.schema.FilebasedSchemaProvider \
      --op UPSERT

    For continuous mode, you need to add the following property:

    
    --continuous \
    --min-sync-interval-seconds 900
    

With the job configured and running on Amazon EMR, the Hudi MultiTable DeltaStreamer efficiently manages real-time data ingestion into your Amazon S3 data lake.

Verify and query data

To verify and query the data, complete the following steps:

  1. Register tables in Data Catalog:
    # Start Spark shell
    spark-shell --conf "spark.serializer=org.apache.spark.serializer.KryoSerializer" \
      --conf "spark.sql.catalog.spark_catalog=org.apache.spark.sql.hudi.catalog.HoodieCatalog" \
      --conf "spark.sql.extensions=org.apache.spark.sql.hudi.HoodieSparkSessionExtension" \
      --jars "/usr/lib/hudi/hudi-spark-bundle.jar"
    
    # In Spark shell
    spark.sql("CREATE DATABASE IF NOT EXISTS hudi_sales_tables")
    
    spark.sql("""
    CREATE TABLE hudi_sales_tables.cust_sales_details
    USING hudi
    LOCATION 's3://hudi-data-bucket-$AWS_ACCOUNT_ID/hudi/hudi_sales_tables.cust_sales_details'
    """)
    
    # Repeat for other tables

  2. Query with Athena:
    -- Sample query
    SELECT * FROM hudi_sales_tables.cust_sales_details LIMIT 10;

You can use Amazon CloudWatch alarms to alert you of issues with the EMR job or data processing. To create a CloudWatch alarm to monitor EMR job failures, enter the following:

aws cloudwatch put-metric-alarm \
    --alarm-name EMR-Hudi-Job-Failure \
    --metric-name JobsFailed \
    --namespace AWS/ElasticMapReduce \
    --statistic Sum \
    --period 300 \
    --threshold 1 \
    --comparison-operator GreaterThanOrEqualToThreshold \
    --dimensions Name=JobFlowId,Value=$CLUSTER_ID \
    --evaluation-periods 1 \
    --alarm-actions $SNS_TOPIC_ARN

Real-world impact of Hudi CDC pipelines

With the pipeline configured and running, you can achieve real-time updates to your data lake, enabling faster analytics and decision-making. For instance:

  • Analytics – Up-to-date inventory data maintains accurate dashboards for ecommerce platforms.
  • Monitoring – CloudWatch metrics confirm the pipeline’s health and efficiency.
  • Flexibility – The seamless handling of schema evolution minimizes downtime and data inconsistencies.

Cleanup

To avoid incurring future charges, follow these steps to clean up resources:

  1. Terminate the Amazon EMR cluster
  2. Delete the Amazon MSK cluster
  3. Remove Amazon S3 objects

Conclusion

In this post, we showed how you can build a scalable data ingestion pipeline using Apache Hudi’s MultiTable DeltaStreamer on Amazon EMR to process data from multiple Amazon MSK topics. You learned how to configure CDC with Apache Hudi, set up real-time data processing with 15-minute sync intervals, and maintain data consistency across multiple sources in your Amazon S3 data lake.

To learn more, explore these resources:

By combining CDC with Apache Hudi, you can build efficient, real-time data pipelines. The streamlined ingestion processes simplify management, enhance scalability, and maintain data quality, making this approach a cornerstone of modern data architectures.


About the authors

Radhakant Sahu

Radhakant Sahu

Radhakant is a Senior Data Engineer and Amazon EMR subject matter expert at Amazon Web Services (AWS) with over a decade of experience in the data space. He specializes in big data, graph databases, AI, and DevOps, building robust, scalable data and analytics solutions that help global clients derive actionable insights and drive business outcomes.

Gautam Bhaghavatula

Gautam Bhaghavatula

Gautam is an AWS Senior Partner Solutions Architect with over 10 years of experience in cloud infrastructure architecture. He specializes in designing scalable solutions, with a focus on compute systems, networking, microservices, DevOps, cloud governance, and AI operations. Gautam provides strategic guidance and technical leadership to AWS partners, driving successful cloud migrations and modernization initiatives.

Sucharitha Boinapally

Sucharitha Boinapally

Sucharitha is a Data Engineering Manager with over 15 years of industry experience. She specializes in agentic AI, data engineering, and knowledge graphs, delivering sophisticated data architecture solutions. Sucharitha excels at designing and implementing advanced knowledge mapping systems.

Veera “Bhargav” Nunna

Veera “Bhargav” Nunna

Veera is a Senior Data Engineer and Tech Lead at AWS pioneering Knowledge Graphs for Large Language Models and enterprise-scale data solutions. With over a decade of experience, he specializes in transforming enterprise AI from concept to production by delivering MVPs that demonstrate clear ROI while solving practical challenges like performance optimization and cost control.

Access Snowflake Horizon Catalog data using catalog federation in the AWS Glue Data Catalog

Post Syndicated from Andries Engelbrecht original https://aws.amazon.com/blogs/big-data/access-snowflake-horizon-catalog-data-using-catalog-federation-in-the-aws-glue-data-catalog/

This is a guest post by Andries Engelbrecht, Principal Partner Solutions Engineer at Snowflake, in partnership with AWS.

AWS announced a new catalog federation feature that allows you to directly access data from Snowflake Horizon Catalog through the AWS Glue Data Catalog. This integration enables you to discover and query Horizon Catalog data in Iceberg format through REST endpoints while applying fine-grained access controls using AWS Lake Formation. The new catalog federation combined with Snowflake’s catalog-linked database feature means users can access data stored across AWS and Snowflake from a single point of entry, reducing data movement and associated costs by eliminating the need to duplicate data across platforms.

In this post, we show you how to connect the AWS Glue Data Catalog to Snowflake Horizon Catalog and query the data using AWS analytics services. We cover how to set up catalogs in Horizon Catalog and configure required permissions, create and configure the federation connection in AWS Glue, implement fine-grained access controls using AWS Lake Formation, and finally, query federated tables using Amazon Athena. This step-by-step approach guides you through the complete process of establishing a integration between your Snowflake and AWS data environments.

Business examples and key benefits

Catalog federation enables several critical business scenarios while delivering key operational and strategic benefits.

Common examples

This federation capability addresses several key business scenarios:

  • Governed, cross-platform analytics: Query data across AWS and Snowflake environments to improve data-driven decision making without data movement or duplication
  • Data mesh implementation: Enable secure and federated data discovery while maintaining domain-oriented ownership
  • Compliance management: Implement consistent access controls and auditing across platforms

Key benefits

  • Operational efficiency: Eliminate data duplication and reduce Extract Transform Load (ETL) workloads
  • Enhanced security: Centralize access control through AWS Lake Formation with fine-grained permissions
  • Cost optimization: Minimize data transfer and storage costs across platforms
  • Improved agility: Enable faster time to insights with direct query access
  • Simplified governance: Maintain unified compliance and audit framework

Solution overview

The solution uses catalog federation in the AWS Glue Data Catalog to integrate with Snowflake Horizon Catalog. This integration supports both Snowflake Horizon, where the catalog is internal to Snowflake, and external catalogs such as Apache Polaris, Snowflake Open Catalog (a managed service that hosts Apache Polaris), and others.

The following diagram illustrates how AWS Glue Data Catalog federates with Snowflake Horizon Catalog, enabling customers to directly access Iceberg-format data managed by Snowflake Horizon Catalog through the Glue Data Catalog.

Architecture diagram showing integration between AWS services and Snowflake using federated catalog connections through Apache Iceberg REST API.

The integration works through three main components:

  1. Authentication: Uses OAuth2 credentials of Snowflake principal
  2. Access Control: AWS Lake Formation manages fine-grained permissions
  3. Query Access: AWS Analytics services like Amazon Athena can directly query the federated tables

Now, we walk through the step-by-step process of setting up this integration.

Prerequisites

Before you begin, confirm you have the following:

Configure Snowflake Horizon Catalog for Iceberg external access

Snowflake Horizon Catalog already supports managing Iceberg tables. For this walkthrough, you need to create Snowflake-managed Iceberg tables with data stored in Amazon S3.

Follow these steps in order:

  1. Create an external volume for S3: First, create an external volume that points to your S3 bucket where Iceberg table data is stored. Follow the instructions in Create External Volume(s) for the Iceberg Tables on S3.
  2. Create a database: Create a database to organize your tables. Refer to the Snowflake database creation documentation.
  3. Create a schema: Create a schema within your database following the Snowflake schema creation guide.
  4. Create an Iceberg table: Create your Iceberg table using the external volume. Follow the instructions to Create Iceberg Table.

After completing these steps, your Snowflake-managed Iceberg tables are ready to federate with AWS Glue Data Catalog.

Configure access control and authentication

To enable AWS Glue to access your Snowflake-managed Iceberg tables, you need to configure access control and obtain authentication credentials.

Step 1: Configure access control

Create a dedicated Snowflake role for external engine access to establish clear governance boundaries. Follow the instructions in Configure Access Control for external engines and set up the appropriate permissions for your Iceberg tables.

Step 2: Obtain an access token

Generate an access token for authenticating AWS Glue to Snowflake Horizon Catalog. Snowflake supports three authentication mechanisms:

  • External OAuth
  • Key-pair authentication
  • Programmatic Access Token (PAT)

Choose the authentication method that best fits your security requirements and follow the corresponding Snowflake documentation to generate your credentials.

Catalog Federation supports OAuth or custom authentication. For details on using OAuth refer to Federate to Snowflake Iceberg Catalog.

For this post, we use custom authentication and generate access token using PAT. Replace role_name with the principal role and token_value with the principal’s Programmatic Access Token.

curl --location 'https://<accountidentifier>.snowflakecomputing.com/polaris/api/catalog/v1/oauth/tokens' \
--header 'Content-Type: application/x-www-form-urlencoded' \
--data-urlencode 'grant_type=client_credentials' \
--data-urlencode 'scope=session:role:<role_name>' \
--data-urlencode 'client_secret=<token_value>'

Note down the access token that is generated.

Step 3: Enable catalog federation

With access control configured and authentication credentials in hand, AWS Glue Catalog Federation can now connect to and access Snowflake’s Horizon Catalog.

Optional: Snowflake Open Catalog configuration

If you prefer to use Snowflake Open Catalog for Iceberg external access instead, refer to Sync a Snowflake-managed table with Snowflake Open Catalog for alternative setup instructions.

Setup Glue Catalog federation with Snowflake Horizon Catalog

Create a secret on AWS Secrets Manager

Log in to AWS console using the IAM role that has access to AWS Secrets Manager. Open Secrets Manager:

  • Choose Store a new secret and select Other type of secret for the secret type.
  • Set the key-value pair:
    • Key: BEARER_TOKEN
    • Value: The access token noted earlier
  • Choose Next and provide the secret name as horizon-secret.
  • Complete the setup by choosing Store.

Alternatively, you can use the CLI to create the secret by running the following command.

Replace your-access-token and your-region with your actual values:

aws secretsmanager create-secret \
    --name horizon-secret \
    --description "Snowflake Horizon access token" \
    --secret-string '{
        "BEARER_TOKEN": "your-access-token"
    }' \
    --region your-region

Create IAM role for catalog federation

As the catalog owner of a federated catalog in AWS Glue Data Catalog, you can use Lake Formation to implement comprehensive access controls for your data teams:

Access control options

You can implement access controls at different granularity levels depending on your governance needs:

  • Coarse-grained: Table-level permissions
  • Fine-grained: Column-level, row-level, and cell-level filtering
  • Tag-based: Dynamic access based on data classification tags

Lake Formation requires an IAM role with permissions to access the underlying S3 locations of your external catalog.

Create an IAM role that enables the Glue Connection to access AWS Secrets Manager, VPC configurations (optional) and Lake formation to manage credential vending for S3 bucket/prefix.

Required permissions

  1. Secrets Manager access: The Glue connection requires permissions to retrieve secret values from Secrets Manager for OAuth tokens stored for your Snowflake service connection.
  2. Amazon Virtual Private Cloud (VPC) Access (optional): When using VPC endpoints to restrict connectivity to your Snowflake Open Catalog account, the Glue connection needs permissions to describe and use VPC network interfaces. This configuration ensures secure, controlled access to both your stored credentials and network resources while maintaining proper isolation through VPC endpoints.
  3. S3 bucket and AWS Key Management Service (KMS) key permission: The Glue connection requires S3 permissions to read certificates if used in the connection setup. Additionally, Lake Formation requires read permissions on the bucket/prefix where the remote catalog table data resides. If the data is encrypted using a KMS key, additional KMS permissions are required.

Setup steps:

Run the following command using AWS CLI by replacing the placeholder with your setup information:

Create a JSON file (e.g., trust-policy.json) with the following structure:

{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Effect": "Allow",
            "Principal": {
                "Service": ["glue.amazonaws.com","lakeformation.amazonaws.com"]
            },
            "Action": "sts:AssumeRole"
        }
    ]
}

Use the aws iam create-role command, referencing the trust policy file:

aws iam create-role \
    --role-name LFDataAccessRole \
    --assume-role-policy-document file://<path_file_downloaded>/trust-policy.json 

First, create a JSON file (such as, permissions-policy.json) for the permissions:


{
"Version": "2012-10-17",
    "Statement": [
        {
            "Effect": "Allow",
            "Action": [
                "secretsmanager:GetSecretValue",
                "secretsmanager:DescribeSecret"
            ],
            "Resource": [
                "<secrets manager ARN>"
            ]
        },
        {
            "Effect": "Allow",
            "Action": [
                "ec2:CreateNetworkInterface",
                "ec2:DeleteNetworkInterface",
                "ec2:DescribeNetworkInterfaces"
            ],
            "Resource": "*",
            "Condition": {
                "ArnEquals": {
                    "ec2:Vpc": "arn:aws:ec2:region:account-id:vpc/<vpc-id>", 
                    "ec2:Subnet": [ 
                        "arn:aws:ec2:region:account-id:subnet/<subnet-id>"
                    ]
                }
            }
        },
        {
           # Required when using custom cert to sign requests.
            "Effect": "Allow",
            "Action": [
                "s3:GetObject"
            ],
            "Resource": [
                "arn:aws:s3:::<bucketname>/<certpath>"
            ]
        },
        { # Required when using customer managed encryption key for s3 
            "Effect": "Allow",
            "Action": [
                "kms:decrypt",
                "kms:encrypt"
            ],
            "Resource": [
                "<kmsKey>"
            ]
        }
    ]
}

Then, attach it to the role:

aws iam put-role-policy \
--role-name LFDataAccessRole \
--policy-name myaccesspolicies \
--policy-document file://<path_file_downloaded>/permissions- policy.json

Create federated catalog in Glue Data Catalog

AWS Glue supports the SNOWFLAKEICEBERGRESTCATALOG connection type for connecting Glue Data Catalog with Snowflake Horizon Catalog and Snowflake Open Catalog. This Glue connector supports OAuth2 authentication and includes additional configuration parameters like CASING_TYPE to customize how AWS Glue Data Catalog discovers metadata in the Snowflake Horizon Catalog accounts.

Log in to your AWS console as a data lake admin and open the AWS Lake Formation console.

  1. Choose Catalog in the left navigation pane and select Create catalog.
  2. Choose the data source as Snowflake Horizon Catalog.
    AWS Lake Formation console screenshot showing Step 1 of catalog creation wizard with five federation type options, Snowflake Horizon Catalog selected.
  3. Provide the following information:
    • Name: Name of the federated catalog in Glue Catalog. For this post, we use federated_lakehousedb
    • Catalog name in Snowflake: Catalog name existing in Snowflake Horizon Catalog, this should match exact name in Horizon catalog. For this post, we use LAKEHOUSEDB
    • For Connection details, choose New connection configurations:
      • Connection name: Name for the glue connection. For this post, we use federatedconnection1.
      • Workspace URL: Horizon IRC url (format: https://<account_identifier>.snowflakecomputing.com)
      • Casing type: choose Uppercase only
      • Authentication:
        • Authentication type: choose Custom. Alternatively, you can select OAuth2 authentication. For Custom authentication, an access token is created, refreshed, and managed by the customer’s application or system and stored using AWS Secrets Manager.
        • OAuth Secret: Provide the secret manager ARN that was created in the previous step.
  • If you have AWS PrivateLink setup and/or a proxy setup, you can provide network details under Settings for network configurations (optional).
  • For Register Glue connection with Lake Formation:
    • Choose the IAM role created earlier(LFDataAccessRole) to manage data access using Lake Formation.

To test the connection, choose Run test. After the connection information is validated, it shows as successful.

Green success banner displaying "Connection test successful" with checkmark icon, confirming valid AWS configuration.

You can now create the catalog by selecting Create catalog.

Alternatively, you can use AWS CLI to create connection and catalog using example commands:

aws glue create-connection \
--connection-input '{
"Name": "federatedconnection1",
"ConnectionType": "SNOWFLAKEICEBERGRESTCATALOG",
"ConnectionProperties": {
    "INSTANCE_URL": "<your-snowflake-account-URL>",
    "ROLE_ARN": "< ARN_of_LFDataAccessRole>",
    "CATALOG_CASING_FILTER": "UPPERCASE_ONLY"
},
"AuthenticationConfiguration": {
    "AuthenticationType": "CUSTOM",
    "SecretArn": "arn:aws:secretsmanager:<your-aws-region>:<your-aws-account-id>:secret:horizon-secret"
}
}' \
--region <your-aws-region>
aws lakeformation register-resource \
    --resource-arn <ARN_of_federatedconnection1_connection> \
    --role-arn <ARN_of_LFDataAccessRole> \
    --with-federation \
    --with-privileged-access \
    --region <your-aws-region>
aws glue create-catalog \
    --name federated_lakehousedb \
    --catalog-input '{
    "FederatedCatalog": {
        "Identifier": "LAKEHOUSEDB",
        "ConnectionName": “federatedconnection1 "
    },
    "CreateTableDefaultPermissions": [],
    "CreateDatabaseDefaultPermissions": []
}'

After the catalog is created, the Horizon databases and tables are listed under the federated catalog.

You can implement fine grained access control on the tables by applying row/column filter using Lake Formation.

Query the data using Athena query editor:

Open the Amazon Athena console and run the following query to access the federated Horizon table:

SELECT * FROM "public"."customer" limit 10;

Clean up

To clean up your resources, complete the following steps:

  1. Drop the Snowflake Database with Cascade.
  2. Drop External Volume created for Iceberg Tables on S3.
  3. Drop the resources in Glue Data Catalog and Lake Formation created for this post.
  4. Delete the IAM roles and S3 buckets used for this post.
  5. Delete any VPC, KMS keys if used for this post setup.

Conclusion

In this post, we demonstrated how to establish a secure connection between AWS Analytics services and Snowflake Horizon Catalog, enabling you to access your data from a single connected and governed view. You learned how to:

  • Configure catalog federation between AWS Glue Data Catalog and Snowflake Horizon Catalog
  • Set up OAuth2 authentication for secure access
  • Grant access to Iceberg table in Snowflake Horizon Catalog using AWS Lake Formation
  • Query federated tables using Amazon Athena

You can follow the same steps to establish a secure connection with open-source catalog options such as Snowflake Open Catalog, a managed service for Apache Iceberg. Remember to clean up any resources you created while following this tutorial to avoid ongoing charges.

To further explore this solution in your environment, consider the following resources:

These resources can help you to implement and optimize this integration pattern for your specific use case. As you begin this journey, remember to start small, validate your architecture with test data, and gradually scale your implementation based on your organization’s needs. Stay tuned for future workshops and resources.


About the authors

 

Andries Engelbrecht

Andries Engelbrecht

Andries is a Principal Partner Solutions Engineer at Snowflake working with AWS. He supports product and service integrations, as well the development of joint solutions with AWS. Andries has over 25 years of experience in the field of data and analytics.

Nidhi Gupta

Nidhi Gupta

Nidhi is a Senior Partner Solutions Architect at AWS, specializing in data analytics and AI. She helps customers and partners build and optimize Snowflake workloads on AWS. Nidhi has extensive experience leading development, production releases and deployments, with focus on Data, AI, ML, generative AI, and Advanced Analytics.

Srividya Parthasarathy

Srividya Parthasarathy

Srividya is a Senior Big Data Architect on the AWS Lake Formation team. She works with the product team and customers to build robust features and solutions for their analytical data platform. She enjoys building data mesh solutions and sharing them with the community.

Pratik Das

Pratik Das

Pratik is a Senior Product Manager with AWS Lake Formation. He is passionate about all things data and works with customers to understand their requirements and build delightful experiences. He has a background in building data-driven solutions and machine learning systems.

 

Access Databricks Unity Catalog data using catalog federation in the AWS Glue Data Catalog

Post Syndicated from Srividya Parthasarathy original https://aws.amazon.com/blogs/big-data/access-databricks-unity-catalog-data-using-catalog-federation-in-the-aws-glue-data-catalog/

AWS has launched the catalog federation capability, enabling direct access to Apache Iceberg tables managed in Databricks Unity Catalog through the AWS Glue Data Catalog. With this integration, you can discover and query Unity Catalog data in Iceberg format using an Iceberg REST API endpoint, while maintaining granular access controls through AWS Lake Formation. This approach significantly reduces operational overhead for managing catalog synchronization and associated costs by alleviating the need to replicate or duplicate datasets between platforms.

In this post, we demonstrate how to set up catalog federation between the Glue Data Catalog and Databricks Unity Catalog, enabling data querying using AWS analytics services.

Use cases and key benefits

This federation capability is particularly valuable if you run multiple data platforms, because you can maintain your existing Iceberg catalog investments while using AWS analytics services. Catalog federation supports read operations and provides the following benefits:

  • Interoperability – You can enable interoperability across different data platforms and tools through Iceberg REST APIs while preserving the value of your established technology investments.
  • Cross-platform analytics – You can connect AWS analytics tools (Amazon Athena, Amazon Redshift, Apache Spark) to query Iceberg and UniForm tables stored in Databricks Unity Catalog. It supports Databricks on AWS integration with the AWS Glue Iceberg REST Catalog for metadata retrieval, while using Lake Formation for permission management.
  • Metadata management – The solution avoids manual catalog synchronization by making Databricks Unity Catalog databases and tables discoverable within the Data Catalog. You can implement unified governance through Lake Formation for fine-grained access control across federated catalog resources.

Solution overview

The solution uses catalog federation in the Data Catalog to integrate with Databricks Unity Catalog. The federated catalog created in AWS Glue mirrors the catalog objects in Databricks Unity Catalog and supports OAuth-based authentication. The solution is represented in the following diagram.

The integration involves three high-level steps:

  1. Set up an integration principal in Databricks Unity Catalog and provide required read access on catalog resources to this principal. Enable OAuth-based authentication for the integration principal.
  2. Set up catalog federation to Databricks Unity Catalog in the Glue Data Catalog:
    1. Create a federated catalog in the Data Catalog using an AWS Glue connection.
    2. Create an AWS Glue connection that uses the credentials of the integration principal (in Step 1) to connect to Databricks Unity Catalog. Configure an AWS Identity and Access Management (IAM) role with permission to Amazon Simple Storage Service (Amazon S3) locations where the Iceberg table data resides. In a cross-account scenario, make sure the bucket policy grants required access to this IAM role.
  3. Discover Iceberg tables in federated catalogs using Lake Formation or AWS Glue APIs. During query operations, Lake Formation manages fine-grained permissions on federated resources and credential vending for access to the underlying data.

In the following sections, we walk through the steps to integrate the Glue Data Catalog with Databricks Unity Catalog on AWS.

Prerequisites

To follow along with the solution presented in this post, you must have the following prerequisites:

  • Databricks Workspace (on AWS) with Databricks Unity Catalog configured.
  • An IAM role that is a Lake Formation data lake administrator in your AWS account. A data lake administrator is an IAM principal that can register S3 locations, access the Data Catalog, grant Lake Formation permissions to other users, and view AWS CloudTrail logs. See Create a data lake administrator for more information.

Configure Databricks Unity Catalog for external access

Catalog federation to a Databricks Unity Catalog uses the OAuth2 credentials of a Databricks service principal configured in the workspace admin settings. This authentication mechanism allows the Data Catalog to access the metadata of various objects (such as catalogs, databases, and tables) within Databricks Unity Catalog, based on the privileges associated with the service principal. For proper functionality, grant the service principal with the necessary permissions (read permission on catalog, schema, and tables) to read the metadata of these objects and allow access from external engines.

Next, catalog federation enables discovery and query of Iceberg tables in your Databricks Unity Catalog. For reading delta tables, enable UniForm on a Delta Lake table in Databricks to generate Iceberg metadata. For more information, refer to Read Delta tables with Iceberg clients.

Follow the Databricks tutorial and documentation to create the service principal and associated privileges in your Databricks workspace. For this post, we use a service principal named integrationprincipal that is configured with required permissions (SELECT, USE CATALOG, USE SCHEMA) on Databricks Unity Catalog objects and will be used for authentication to catalog instance.

Catalog federation supports OAuth2 authentication, so enable OAuth for the service principal and note down the client_id and client_secret for later use.

Set up Data Catalog federation with Databricks Unity Catalog

Now that you have service principal access for Databricks Unity Catalog, you can set up catalog federation in the Data Catalog. To do so, you create an AWS Secrets Manager secret and create an IAM role for catalog federation.

Create secret

Complete the following steps to create a secret:

  1. Sign in to the AWS Management Console using an IAM role with access to Secrets Manager.
  2. On the Secrets Manager console, choose Store a new secret and Other type of secret.
  3. Set the key-value pair:
    1. Key: USER_MANAGED_CLIENT_APPLICATION_CLIENT_SECRET
    2. Value: The client secret noted earlier
  4. Choose Next.
  5. Enter a name for your secret (for this post, we use dbx).
  6. Choose Store.

Create IAM role for catalog federation

As the catalog owner of a federated catalog in the Data Catalog, you can use Lake Formation to implement comprehensive access controls, including table filters, column filters, and row filters, as well as tag-based access for your data teams.

Lake Formation requires an IAM role with permissions to access the underlying S3 locations of your external catalog.

In this step, you create an IAM role that enables the AWS Glue connection to access Secrets Manager, optional virtual private cloud (VPC) configurations, and Lake Formation to manage credential vending for the S3 bucket and prefix:

  • Secrets Manager access – The AWS Glue connection requires permissions to retrieve secret values from Secrets Manager for OAuth tokens stored for your Databricks Unity service connection.
  • VPC access (optional) – When using VPC endpoints to restrict connectivity to your Databricks Unity account, the AWS Glue connection needs permissions to describe and utilize VPC network interfaces. This configuration provides secure, controlled access to both your stored credentials and network resources while maintaining proper isolation through VPC endpoints.
  • S3 bucket and AWS KMS key permission – The AWS Glue connection requires Amazon S3 permissions to read certificates if used in the connection setup. Additionally, Lake Formation requires read permissions on the bucket and prefix where the remote catalog table data resides. If the data is encrypted using an AWS Key Management Service (AWS KMS) key, additional AWS KMS permissions are required.

Complete the following steps:

  1. Create an IAM role called LFDataAccessRole with the following policies:
    {
     "Version": "2012-10-17",
         "Statement": [
             {
                 "Effect": "Allow",
                 "Action": [
                     "secretsmanager:GetSecretValue",
                     "secretsmanager:DescribeSecret"
                 ],
                 "Resource": [
                     "<secrets manager ARN>"
                 ]
             },
             {
                 "Effect": "Allow",
                 "Action": [
                     "ec2:CreateNetworkInterface",
                     "ec2:DeleteNetworkInterface",
                     "ec2:DescribeNetworkInterfaces"
                 ],
                 "Resource": "*",
                 "Condition": {
                     "ArnEquals": {
                         "ec2:Vpc": "arn:aws:ec2:region:account-id:vpc/<vpc-id>", 
                         "ec2:Subnet": [ 
                             "arn:aws:ec2:region:account-id:subnet/<subnet-id>" 
                         ]
                     }
                 }
             },
             {
                # Required when using custom cert to sign requests.
                 "Effect": "Allow",
                 "Action": [
                     "s3:GetObject"
                 ],
                 "Resource": [
                     "arn:aws:s3
    :::<bucketname>/<certpath>"
                 ]
             },
             { # Required when using customer managed encryption key for s3 
                 "Effect": "Allow",
                 "Action": [
                     "kms:decrypt",
                     "kms:encrypt"
                 ],
                 "Resource": [
                     "<kmsKey>"
                 ]
             }
         ]
     }

  2. Configure the role with the following trust policy:
    {
          "Version": "2012-10-17",
          "Statement": [
              {
                  "Effect":  "Allow",
                  "Principal": {
                       "Service": ["glue.amazonaws.com","lakeformation.amazonaws.com"]
                  },
                  "Action":  "sts:AssumeRole"
              }
          ]
      }

Create federated catalog in Data Catalog

AWS Glue supports the DATABRICKSICEBERGRESTCATALOG connection type for connecting the Data Catalog with managed Databricks Unity Catalog. This AWS Glue connector supports OAuth2 authentication for discovering metadata in Databricks Unity Catalog.

Complete the following steps to create the federated catalog:

  1. Sign in to the console as a data lake admin.
  2. On the Lake Formation console, choose Catalogs in the navigation pane.
  3. Choose Create catalog.
  4. For Name, enter a name for your catalog.
  5. For Catalog name in Databricks, enter the name of a catalog existing in Databricks Unity Catalog.
  6. For Connection name, enter a name for the AWS Glue connection.
  7. For Workspace URL, enter the Unity Iceberg REST API URL (in format https://<workspace-url>/cloud.databricks.com).
  8. For Authentication, provide the following information:
    1. For Authentication type, choose OAuth2. Alternatively, you can choose Custom authentication. For Custom authentication, an access token is created, refreshed, and managed by the customer’s application or system and stored using Secrets Manager.
    2. For Token URL, enter the token authentication server URL.
    3. For OAuth Client ID, enter the client_id for integrationprincipal.
    4. For OAuth Secret, enter the secret ARN that you created in the previous step. Alternatively, you can provide the client_secret directly.
    5. For Token URL parameter map scope, provide the API scope supported.
  9. If you have AWS PrivateLink set up or a proxy set up, you can provide network details under Settings for network configurations.
  10. For Register Glue connection with Lake Formation, choose the IAM role (LFDataAccessRole) created earlier to manage data access using Lake Formation.

When the setup is done using AWS Command Line Interface (AWS CLI) commands, you have options to create two separate IAM roles:

  • IAM role with policies to access network and secrets, which AWS Glue assumes to manage authentication
  • IAM role with access to the S3 bucket, which Lake Formation assumes to manage credential vending for data access

On the console, this setup is simplified with a single role having combined policies. For more details, refer to Federate to Databricks Unity Catalog.

  1. To test the connection, choose Run test.
  2. You can proceed to create the catalog.

After you create the catalog, you can see the databases and tables in Databricks Unity Catalog listed under the federated catalog. You can implement fine-grained access control on the tables by applying row and column filters using Lake Formation. The following video shows the catalog federation setup with Databricks Unity Catalog.

Discover and query the data using Athena

In this post, we show how to use the Athena query editor to discover and query the Databricks Unity Catalog tables. On the Athena console, run the following query to access the federated table:SELECT * FROM "customerschema"."person" limit 10;The following video demonstrates querying the federated table from Athena.

If you use the Amazon Redshift query engine, you must create a resource link on the federated database and grant permission on the resource link to the user or role. This database resource link is automounted under awsdatacatalog based on the permission granted for the user or role and available for querying. For instructions, refer to Creating resource links.

Clean up

To clean up your resources, complete the following steps:

  1. Delete the catalog and namespace in Databricks Unity Catalog for this post.
  2. Drop the resources in the Data Catalog and Lake Formation created for this post.
  3. Delete the IAM roles and S3 buckets used for this post.
  4. Delete any VPC and KMS keys if used for this post.

Conclusion

In this post, we explored the key elements of catalog federation and its architectural design, illustrating the interaction between the AWS Glue Data Catalog and Databricks Unity Catalog through centralized authorization and credential distribution for protected data access. By removing the requirement for complicated synchronization workflows, catalog federation makes it possible to query Iceberg data on Amazon S3 directly at its source using AWS analytics services with data governance across multi-catalog platforms. Try out the solution for your own use case, and share your feedback and questions in the comments.


About the Authors

Srividya Parthasarathy

Srividya Parthasarathy

Srividya is a Senior Big Data Architect on the AWS Lake Formation team. She works with the product team and customers to build robust features and solutions for their analytical data platform. She enjoys building data mesh solutions and sharing them with the community.

Venkatavaradhan (Venkat) Viswanathan

Venkatavaradhan (Venkat) Viswanathan

Venkat” is a Global Partner Solutions Architect at Amazon Web Services. Venkat is a Technology Strategy Leader in Data, AI, ML, Generative AI, and Advanced Analytics. Venkat is a Global SME for Databricks and helps AWS customers design, build, secure, and optimize Databricks workloads on AWS.

Use Amazon SageMaker custom tags for project resource governance and cost tracking

Post Syndicated from David Victoria original https://aws.amazon.com/blogs/big-data/use-amazon-sagemaker-custom-tags-for-project-resource-governance-and-cost-tracking/

Amazon SageMaker announced a new feature that you can use to add custom tags to resources created through an Amazon SageMaker Unified Studio project. This helps you enforce tagging standards that conform to your organization’s service control policies (SCPs) and helps enable cost tracking reporting practices on resources created across the organization.

As a SageMaker administrator, you can configure a project profile with tag configurations that will be pushed down to projects that currently use or will use that project profile. The project profile is set up to pass either required key and value tag pairings or pass the key of the tag with a default value that can be modified during project creation. All tags passed to the project will result in the resources created by that project being tagged. This provides you with a governance mechanism that enforces that project resources have the expected tags across all projects of the domain.

The first release of custom tags for project resources is supported through an application programming interface (API), through Amazon DataZone SDKs. In this post, we look at use cases for custom tags and how to use the AWS Command Line Interface (AWS CLI) to add tags to project resources.

What we hear from customers

As customers continue to build and collaborate using AWS tools for model development, generative AI, data processing, and SQL analytics, they see the need to bring control and visibility into the resources being created. To support connectivity to these AWS tools from SageMaker Unified Studio projects, many different types of resources across AWS services need to be created. These resources are created through AWS CloudFormation stacks (through project environment deployment) by the Amazon SageMaker service. From customers we hear the following use cases:

  • Customers need to enforce that tagging practices conform to company policies through the use of AWS controls, such as SCPs, for resource creation. These controls block the creation of resources unless specific tags are placed on the resource.
  • Customers can also start with policies to enforce that the correct tags are placed when resources are created with the additional goal of standardizing on resource reporting. By placing identifiable information on resources when created, they enforce consistency and completeness when performing cost attribution reporting and observability.

Customer Swiss Life uses SageMaker as a single solution for cataloging, discovery, sharing, and governance of their enterprise data across business domains. They require all resources have a set of mandatory tags for their finance group to bill organizations across their company for the AWS resources created.

“The launch of project resource tags for Amazon SageMaker allows us to bring visibility to the costs incurred across our accounts. With this capability we are able to meet the resource tagging guidelines of our company and have confidence in attributing costs across our multi-account setup for the resources created by Amazon SageMaker projects.”

– Tim Kopacz, Software Developer at Swiss Life

Prerequisites

To get started with custom tags, you must have the following resources:

  • A SageMaker Unified Studio domain.
  • An AWS Identity and Access Management (IAM) entity with privileges to make AWS CLI calls to the domain.
  • An IAM entity authorized to make changes to the domain IAM provisioning role. If SageMaker created this for you, it will be called AmazonSageMakerProvisioning-<accountId>. The provisioning role provisions and manages resources defined in the selected blueprints in your account.

How to set up project resource tags

The following steps outline how you can configure custom tags for your SageMaker Unified Studio project resources:

  1. (Optional) Update the SageMaker provisioning role to permit specific tag keys.
  2. Create a new project profile with project resource tags configured.
  3. Create a new project with project resource tags.
  4. Update an existing project with project resource tags.
  5. Validate that the resources are tagged.

(Optional) Update a SageMaker provisioning role to permit tag key values

The AmazonSageMakerProvisioning-<accountId> role has an AWS managed policy with condition aws:TagKeys allowing tags to be created by this role only if the tag key begins with AmazonDataZone. For this example, we will change the tag key to begin with different strings. Skip to Create a new project profile with project resource tags configured if you don’t need tag keys to have a different structure (such as begins with, contains, and so on)

  1. Open the AWS Management Console and go to IAM.
  2. In the navigation pane, choose Roles.
  3. In the list, choose AmazonSageMakerProvisioning-<accountId>.
  4. Choose the Permissions tab.
  5. Choose Add permissions, and then choose Create inline policy.
  6. Under Policy editor, select JSON.
  7. Enter the following policy. Add the strings under the condition aws:TagKeys. In this example, tag keys beginning with ACME or tag keys with the exact match of CostCenter will be created by the role.
    {
        "Version": "2012-10-17",
        "Statement": [
            {
                "Sid": "CustomTagsUnTagPermissions",
                "Effect": "Allow",
                "Action": [
                    "codecommit:UntagResource",
                    "iam:UntagRole",
                    "logs:UntagResource",
                    "athena:UntagResource",
                    "redshift-serverless:UntagResource",
                    "scheduler:UntagResource",
                    "bedrock:UntagResource",
                    "neptune-graph:UntagResource",
                    "quicksight:UntagResource",
                    "glue:UntagResource",
                    "airflow:UntagResource",
                    "secretsmanager:UntagResource",
                    "lambda:UntagResource",
                    "emr-serverless:UntagResource",
                    "elasticmapreduce:RemoveTags",
                    "sagemaker:DeleteTags",
                    "ec2:DeleteTags"
                ],
                "Resource": "*",
                "Condition": {
                    "StringEquals": {
                        "aws:ResourceAccount": "${aws:PrincipalAccount}"
                    },
                    "ForAllValues:StringLike": {
                        "aws:TagKeys": [
                            "AmazonDataZone*",
                            "ACME*",
                            "CostCenter"
                        ]
                    },
                    "Null": {
                        "aws:ResourceTag/AmazonDataZoneProject": "false"
                    }
                }
            },
            {
                "Sid": "CustomTagsTaggingPermissions",
                "Effect": "Allow",
                "Action": [
                    "cloudformation:TagResource",
                    "codecommit:TagResource",
                    "iam:TagRole",
                    "glue:TagResource",
                    "athena:TagResource",
                    "lambda:TagResource",
                    "redshift-serverless:TagResource",
                    "logs:TagResource",
                    "secretsmanager:TagResource",
                    "sagemaker:AddTags",
                    "emr-serverless:TagResource",
                    "neptune-graph:TagResource",
                    "bedrock:TagResource",
                    "elasticmapreduce:AddTags",
                    "airflow:TagResource",
                    "scheduler:TagResource",
                    "quicksight:TagResource",
                    "emr-containers:TagResource",
                    "logs:CreateLogGroup",
                    "athena:CreateWorkGroup",
                    "scheduler:CreateScheduleGroup",
                    "cloudformation:CreateStack",
                    "ec2:*"
                ],
                "Resource": "*",
                "Condition": {
                    "ForAnyValue:StringLike": {
                        "aws:TagKeys": [
                            "AmazonDataZone*",
                            "ACME*",
                            "CostCenter"
                        ]
                    },
                    "StringEquals": {
                        "aws:ResourceAccount": "${aws:PrincipalAccount}"
                    }
                }
            }
        ]
    }

It’s possible to scope down the specific AWS service tag and un-tag permissions based on which blueprints or capabilities are being used.

Create a new project profile with project resource tags configured

Use the following steps to create a new SQL Analytics project profile with custom tags. The example uses AWS CLI commands.

  1. Open the AWS CloudShell console.
  2. Create a project profile using the following CLI command.
    1. The project-resource-tags parameter consists of key (tag key), value (tag value), and isValueEditable (boolean indicating if the tag value can be modified during project creation or update).
    2. The allow-custom-project-resource-tags parameter set to true permits the project creator to create additional key-value pairs. The key needs to conform to the inline policy of the AmazonSageMakerProvisioning-<accountId> role.
    3. The project-resource-tags-description parameter is a description field for project resource tags. The max character limit is 2,048. The description needs to be passed in every time create-project-profile or update-project-profile is called.
    aws datazone create-project-profile \
      --name "SQL Analytics with Project Resource Tags" \
      --description "Analyze your data in SageMaker Lakehouse using SQL" \
      --domain-identifier "$DOMAIN_ID" \
      --region "$REGION" \
      --status ENABLED \
      --project-resource-tags '[
        {
            "key": "ACME-Application",
            "value": "SageMaker",
            "isValueEditable": false
        },
        {
            "key": "CostCenter",
            "value": "123",
            "isValueEditable": true
        }
      ]' \
      --allow-custom-project-resource-tags \
      --environment-configurations '[
        {
            "name": "Tooling",
            "description": "Configuration for the Tooling Environment",
            "environmentBlueprintId": "",
            "deploymentMode": "ON_CREATE",
            "deploymentOrder": 0,
            "awsAccount": {
            "awsAccountId": "$ACCOUNT"
        },
        "awsRegion": {
            "regionName": "$REGION"
        },
            "configurationParameters": {
                "parameterOverrides": [
                    {
                        "name": "enableSpaces",
                        "value": "false",
                        "isEditable": false
                    },
                    {
                        "name": "maxEbsVolumeSize",
                        "isEditable": false
                    },
                    {
                        "name": "idleTimeoutInMinutes",
                        "isEditable": false
                    },
                    {
                        "name": "lifecycleManagement",
                        "isEditable": false
                    },
                    {
                        "name": "enableNetworkIsolation",
                        "isEditable": false
                    }
                ]
            }
        },
        {
            "name": "Lakehouse Database",
            "description": "Creates databases in Amazon SageMaker Lakehouse for storing tables in S3 and Amazon Athena resources for your SQL workloads",
            "environmentBlueprintId": "",
            "deploymentMode": "ON_CREATE",
            "deploymentOrder": 1,
            "awsAccount": {
                "awsAccountId": "$ACCOUNT"
            },
            "awsRegion": {
            "regionName": "$REGION"
            },
            "configurationParameters": {
                "parameterOverrides": [
                    {
                        "name": "glueDbName",
                        "value": "glue_db",
                        "isEditable": true
                    }
                ]
            }
        },
        {
            "name": "OnDemand RedshiftServerless",
            "description": "Enables you to create an additional Amazon Redshift Serverless workgroup for your SQL workloads",
            "environmentBlueprintId": "",
            "deploymentMode": "ON_DEMAND",
            "awsAccount": {
            "awsAccountId": "$ACCOUNT"
            },
            "awsRegion": {
                "regionName": "$REGION"
            },
            "configurationParameters": {
                "parameterOverrides": [
                    {
                        "name": "redshiftDbName",
                        "value": "dev",
                        "isEditable": true
                        },
                        {
                        "name": "redshiftMaxCapacity",
                        "value": "512",
                        "isEditable": true
                        },
                        {
                        "name": "redshiftWorkgroupName",
                        "value": "redshift-serverless-workgroup",
                        "isEditable": true
                        },
                        {
                        "name": "redshiftBaseCapacity",
                        "value": "128",
                        "isEditable": true
                        },
                        {
                        "name": "connectionName",
                        "value": "redshift.serverless",
                        "isEditable": true
                        },
                        {
                        "name": "connectToRMSCatalog",
                        "value": "false",
                        "isEditable": false
                        }
                    ]
                }
            },
            {
                "name": "OnDemand Catalog for Redshift Managed Storage",
                "description": "Enables you to create additional catalogs in Amazon SageMaker Lakehouse for storing data in Redshift Managed Storage",
                "environmentBlueprintId": "",
                "deploymentMode": "ON_DEMAND",
                "awsAccount": {
                "awsAccountId": "$ACCOUNT"
                },
                "awsRegion": {
                    "regionName": "$REGION"
                },
                "configurationParameters": {
                    "parameterOverrides": [
                        {
                            "name": "catalogName",
                            "isEditable": true
                        },
                        {
                            "name": "catalogDescription",
                            "value": "RMS catalog",
                            "isEditable": true
                        }
                    ]
                }
            }
      ]'

This project profile will have the tag ACME-Application = SageMaker placed on all projects associated to the project profile and cannot be modified by the project creator. The tag CostCenter = 123 can have the value modified by the project creator because the isValueEditable property is set to true.

Grant permissions for users to use the project profile during project creation. In the Authorization section of the project profile set either Selected users or groups or Allow all users and groups.

The use of the allow-custom-project-resource-tags parameter means the project creator can add their own tags (key-value pair). The key must conform to the condition check in the policy of the provisioning role (AmazonSageMakerProvisioning-<accountId>). If the allow-custom-project-resource-tagsparameter is changed to false after a project created tags, tags created by the project will be removed during the next project update.

Updates to the project profile

Updates to project resource tags are possible through the update-project-profile command. The command will replace all values in the project-resource-tags section so be sure to include the exhaustive set of tags. Updates to the project profile are reflected in projects after running the update-project command or when a new project is created using the project profile. The following example adds a new tag, ACME-BusinessUnit = Retail.

There are three ways to work with the project-resource-tags parameter when updating the project profile.

  • Passing a non-empty list of project resource tags will replace the tags currently configured on the project profile.
  • Passing an empty list of project resource tags will clear out all previously configured tags:
    • --project-resource-tags '[]'
  • Not including the project resource tag parameter will keep previously configured tags as-is.
aws datazone update-project-profile \
  --domain-identifier "$DOMAIN_ID" \
  --identifier "$PROJECT_PROFILE_ID" \
  --region "$REGION" \
  --project-resource-tags '[
    {
        "key": "ACME-Application",
        "value": "SageMaker",
        "isValueEditable": false
    },
    {
        "key": "CostCenter",
        "value": "123",
        "isValueEditable": true
    },
    {
        "key": "ACME-BusinessUnit",
        "value": "Retail",
        "isValueEditable": false
    }
  ]'

Create a new project with project resource tags

The following steps walk you through creating a new project that inherits tags from the project profile and lets the project creator modify one of the tag values.

  1. Create a project using the following example CLI command.
  2. Modify the CostCenter tag value using the --resource-tags parameter. Tags configured on the project profile where the isValueEditable attribute is false will be pushed to the project automatically.
    aws datazone create-project \
      --domain-identifier "$DOMAIN_ID" \
      --region "$REGION" \
      --name "$PROJECT_NAME" \
      --description "New project with tags" \
      --project-profile-id "$PROJECT_PROFILE_ID" \
      --resource-tags '{
            "CostCenter": "456"
        }'

Update existing project with project resource tags

For existing projects associated to the project profile, you must update the project for the new tags to be applied.

  1. Update the project using the following example CLI command.
  2. In this scenario, an editable value needs to be updated and a new tag added. Tag CostCenter will have its default value overwritten as “789” and the new ACME-Department = Finance tag will be added.
    aws datazone update-project \
      --domain-identifier "$DOMAIN_ID" \
      --identifier "$PROJECT_ID" \
      --project-profile-version "latest" \
      --region "$REGION" \
      --resource-tags '{
            "CostCenter": "789",
            "ACME-Department": "Finance"
        }' 

Project level tags (those not configured from the project profile) need to be passed during project update to be preserved. For tags with isValueEditable = true configured from the project profile, any override previously set needs to be applied or the value will revert to the default from the project profile.

Validating resources are tagged

Validate that tags are placed correctly. An example resource that is created by the project is the project IAM role. Viewing the tags for this role should show the tags configured from the project profile.

  1. Open SageMaker Unified Studio to get the project role from the Project details section of the project. The role name begins with datazone_usr_role_.
  2. Open the IAM console.
  3. In the navigation pane, choose Roles.
  4. Search for the project IAM role.
  5. Select the Tags tab.

Conclusion

In this post, we discussed tagging related use cases from customers and walked through getting started with custom tags in Amazon SageMaker to place tags on the resources created by the project. By giving administrators a way to configure project profiles with standardized tag configurations, you can now help ensure consistent tagging practices across all SageMaker Unified Studio projects while maintaining compliance with SCPs. This feature addresses two critical customer needs: enforcing organizational tagging standards through automated governance mechanisms and enabling accurate cost attribution reporting across multi-service deployments.

To learn more, visit Amazon SageMaker, then get started with Project resource tags.


About the authors

David Victoria

David Victoria

David is a Senior Technical Product Manager with Amazon SageMaker at AWS. He focuses on improving administration and governance capabilities needed for customers to support their analytics systems. He is passionate about helping customers realize the most value from their data in a secure, governed manner.

Rohit Srikanta

Rohit Srikanta

Rohit is a Senior Software Engineer at AWS. He works on building and scaling services within Amazon SageMaker. He focuses on developing robust and scalable distributed systems and is passionate about solving complex engineering challenges to deliver maximum customer value.

Ahan Malli

Ahan Malli

Ahan is a Software Development Engineer at AWS. He works on the core data and governance layer behind Amazon SageMaker. He’s passionate about building scalable distributed systems and streamlining developer workflows. When he’s not coding, you can find him traveling or hiking Pacific Northwest trails.

Create AWS Glue Data Catalog views using cross-account definer roles

Post Syndicated from Aarthi Srinivasan original https://aws.amazon.com/blogs/big-data/create-aws-glue-data-catalog-views-using-cross-account-definer-roles/

With AWS Glue Data Catalog views you can create a SQL view in the Data Catalog that references one or more base tables. These multi-dialect views support various SQL query engines, providing consistent access across multiple Amazon Web Services (AWS) services including Amazon Athena, Amazon Redshift Spectrum, and Apache Spark in both Amazon EMR and AWS Glue 5.0.

You can now create Data Catalog views using a cross-account AWS Identity and Access Management (IAM) definer role. A definer role is an IAM role used to create the Data Catalog view and has SELECT permissions on all columns of the underlying base tables. This definer role is assumed by AWS Glue and AWS Lake Formation service principals to vend credentials to the base tables’ data whenever the view is queried. The definer role allows the Data Catalog view to be shared to principals or AWS accounts so that you can share a filtered subset of data without sharing the base tables.

Previously, Data Catalog views required a definer role within the same AWS account as the base tables. The introduction of cross-account definer roles enables Data Catalog view creation in enterprise data mesh architectures. In this setup, database and table metadata is centralized in a governance account, and individual data owner accounts maintain control over table creation and management through their IAM roles. Data owner accounts can now create and manage Data Catalog views in the central governance accounts using their existing continuous integration and continuous delivery (CI/CD) pipeline roles.

In this post, we show you a cross-account scenario involving two AWS accounts: a central governance account containing the tables and hosting the views and a data owner (producer) account with the IAM role used to create and manage views. We provide implementation details for both SPARK dialect using AWS SDK code samples and ATHENA dialect using SQL commands. Using this approach, you can implement sophisticated data governance models at enterprise scale while maintaining operational efficiency across your AWS environment.

Key benefits

Key benefits for cross-account definer roles are as follows:

  • Enhanced data mesh support – Enterprises with multi-account data lakehouse architectures can now maintain their existing operational model where data owner accounts manage table creation and updates using their established IAM roles. These same roles can now create and manage Data Catalog views across account boundaries.
  • Strengthened security controls – By keeping table and view management within data owner account roles:
    • Security posture is enhanced through proper separation of duties.
    • Audit trails become more comprehensive and meaningful.
    • Access controls follow the principle of least privilege.
  • Elimination of data duplication – Data owner accounts can create views in central accounts that:
    • Provide access to specific data subsets without duplicating tables.
    • Reduce storage costs and management overhead.
    • Maintain a single source of truth while enabling targeted data sharing.

Solution overview

An example customer has a database with two transaction tables in their central account, where the catalog and permissions are maintained. With the database shared to the data owner (producer) account, we create a Data Catalog view in the central account on these two tables, using the producer’s definer role. The view from the central account can be shared to additional consumer accounts and queried. We illustrate creating the SPARK dialect using create-table CLI, and add the ATHENA dialect for the same view from the Athena console. We also provide the AWS SDK sample code for CreateTable() and UpdateTable(), with view definition and a sample pySpark script to read and verify the view in AWS Glue.

The following diagram shows the table, view, and definer IAM role placements between a central governance account and data producer account.

Prerequisites

To perform this solution, you need to have the following prerequisites:

  1. Two AWS accounts with AWS Lake Formation set up. For details, refer to Set up AWS Lake Formation. The Lake Formation setup includes registering your IAM admin role as Lake Formation administrator. In the Data Catalog settings, shown in the following screenshot, Default permissions for newly created databases and tables is set to use Lake Formation permissions only. Cross-account version settings is set to Version 4.

  1. Create an IAM role Data-Analyst in the producer account. For the IAM permissions on this role, refer to Data analyst permissions. This role will also be used as the view definer role. Add the permissions to this definer role from the Prerequisites for creating views.

Create database and tables in the central account

In this step, you create two tables in the central governance account and populate them with few rows of data:

  1. Sign in to the central account as admin user. Open the Athena console and set up the Athena query results bucket.
  2. Run the following queries to create two sample Iceberg tables, representing bank customer transactions data:
/* Check if the Database exists, if not create new database. */
CREATE DATABASE IF NOT EXISTS bankdata_icebergdb;

/*Create transaction_table1*/ Replace the bucket name
CREATE TABLE bankdata_icebergdb.transaction_table1 (
  transaction_id string,
  transaction_type string,
  transaction_amount double)
LOCATION 's3://<bucket-name>/bankdata_icebergdb/transaction-table1'
TBLPROPERTIES (
  'table_type'='iceberg',
  'write_compression'='zstd'
);

/*Create transaction_table2*/
CREATE TABLE bankdata_icebergdb.transaction_table2 (
  transaction_id string,
  transaction_location string,
  transaction_date date)
LOCATION 's3://<bucket-name>/bankdata_icebergdb/transaction-table2'
TBLPROPERTIES (
  'table_type'='iceberg',
  'write_compression'='zstd'
);


INSERT INTO bankdata_icebergdb.transaction_table1 (transaction_id, transaction_type, transaction_amount)
VALUES
  ('T001', 'purchase', 50.0),
  ('T002', 'purchase', 120.0),
  ('T003', 'refund', 200.5),
  ('T004', 'purchase', 80.0),
  ('T005', 'withdrawal', 500.0),
  ('T006', 'purchase', 300.0),
  ('T007', 'deposit', 1000.0),
  ('T008', 'refund', 20.0),
  ('T009', 'purchase', 150.0),
  ('T010', 'withdrawal', 75.0);


INSERT INTO bankdata_icebergdb.transaction_table2 (transaction_id, transaction_location, transaction_date)
VALUES
  ('T001', 'Charlotte', DATE '2024-10-01'),
  ('T002', 'Seattle', DATE '2024-10-02'),
  ('T003', 'Chicago', DATE '2024-10-03'),
  ('T004', 'Miami', DATE '2024-10-04'),
  ('T005', 'New York', DATE '2024-10-05'),
  ('T006', 'Austin', DATE '2024-10-06'),
  ('T007', 'Denver', DATE '2024-10-07'),
  ('T008', 'Boston', DATE '2024-10-08'),
  ('T009', 'San Jose', DATE '2024-10-09'),
  ('T010', 'Phoenix', DATE '2024-10-10');
  1. Verify the created tables in Athena query editor by running a preview.

Share the database and tables from central to producer account

In the central governance account, you share the database and the two tables to the producer account and the Data-Analyst role in producer.

  1. Sign in to the Lake Formation console as the Lake Formation admin role.
  2. In the navigation pane, choose Data permissions.
  3. Choose Grant and provide the following information:
    1. For Principals, select External accounts and enter the producer account ID, as shown in the following screenshot.
    2. For Named Data Catalog Resources, select the default catalog and database bankdata_icebergdb, as shown in the following screenshot.
    3. Under Database permissions, select Describe. For Grantable permissions, select Describe.
    4. Choose Grant.
    5. Repeat the preceding steps to grant access to the producer account definer role Data-Analyst on the database bankdata_icebergdb and the two tables transaction_table1 and transaction_table2 as follows.
    6. Under Database permissions, grant Create table and Describe permissions.
    7. Under Table permissions, grant Select and Describe on all columns.

With these steps, the central governance account data admin steward has shared the database and tables to the producer account definer role.

Steps for producer account

Follow these steps for the producer account:

  1. Sign in to the Lake Formation console on the producer account as the Lake Formation administrator.
  2. In the left navigation pane, choose Databases. A blue banner will appear on the console, showing pending invitations from AWS Resource Access Manager (AWS RAM).
  3. Open the AWS RAM console and review the AWS RAM shares under Shared with me. You will see the AWS RAM shares in pending state. Select the pending AWS RAM share from central account and choose Accept resource share. After the resource share request is accepted, the shared database shows up in the producer account.
  4. On the Lake Formation console, select the database. On the Create dropdown list, choose Resource link. Provide a name rl_bank_iceberg and choose Create.
  5. Let’s grant Describe permission on the resource link to the Data-Analyst role in the producer account in the following steps.
    1. In the left navigation pane, choose Data permissions. Choose the Data-Analyst role. Select the resource link rl_bank_iceberg for the database as shown in the following screenshot.
    2. Grant Describe permission on the resource link.

Note: Cross-account Data Catalog views can’t be created using a resource link, although a resource link is needed for the SDK use of SPARK dialect.

Next, add the central account Data Catalog as a Data Source in Athena from producer account:

  1. Open the Athena console.
  2. On the left navigation pane, choose Data sources and catalogs. Choose Create data source.
    1. Select S3-AWS Glue Data Catalog.
    2. Choose AWS – Glue Data Catalog in another account and name the data source as centraladmin.
    3. Choose Next and then create data source.

After the data source is created, navigate to the Query editor and verify the Data source centraladmin appears, as shown in the following screenshot.

The definer role can also now access and query the central catalog database.

Create SPARK dialect view

In this step, you create a view with SPARK dialect, using AWS Glue CLI command create-table:

  1. Sign in to the AWS console in the producer account as Data-Analyst role. Enter the following command in your CLI environment, such as AWS CloudShell, to create a SPARK DIALECT:
aws glue create-table --cli-input-json '{
   "DatabaseName": "rl_bank_iceberg",
   "TableInput": {
     "Name": "mdv_transaction1",
     "StorageDescriptor": {
       "Columns": [
         {
           "Name": "transaction_id",
           "Type": "string"
         },
         {
           "Name": "transaction_type",
           "Type": "string"
         },
         {
           "Name": "transaction_amount",
           "Type": "float"
         },
         {
           "Name": "transaction_location",
           "Type": "string"
         },
         {
           "Name": "transaction_date",
           "Type": "date"
         }
       ],
       "SerdeInfo": {}
     },
     "ViewDefinition": {
       "SubObjects": [
         "arn:aws:glue:<your-region>:<your-central-account-id>:table/bankdata_icebergdb/transaction_table1",
         "arn:aws:glue:<your-region>:<your-central-account-id>:table/bankdata_icebergdb/transaction_table2"
        ],
       "IsProtected": true,
       "Representations": [
         {
           "Dialect": "SPARK",
           "DialectVersion": "1.0",
           "ViewOriginalText": "SELECT t1.transaction_id, t1.transaction_type, t1.transaction_amount, t2.transaction_location, t2.transaction_date FROM transaction_table1 t1 JOIN transaction_table2 t2 ON t1.transaction_id = t2.transaction_id WHERE t1.transaction_amount > 100;",
           "ViewExpandedText": "SELECT t1.transaction_id, t1.transaction_type, t1.transaction_amount, t2.transaction_location, t2.transaction_date FROM transaction_table1 t1 JOIN transaction_table2 t2 ON t1.transaction_id = t2.transaction_id WHERE t1.transaction_amount > 100;"
         }
       ]
     }
   }
 }'
  1. Open the Lake Formation console and verify if the view is created. Verify the dialect of the view on the SQL definitions tab for the view details.

Add ATHENA dialect

To add ATHENA dialect, follow these steps:

  1. On the Athena console, select centraladmin from the Data source.
  2. Enter the following SQL script to create the ATHENA dialect for the same view:
ALTER VIEW mdv_transaction1 FORCE ADD DIALECT
AS
SELECT t1.transaction_id, t1.transaction_type, t1.transaction_amount, t2.transaction_location, t2.transaction_date FROM transaction_table1 t1 JOIN transaction_table2 t2 ON t1.transaction_id = t2.transaction_id WHERE t1.transaction_amount > 100

We can’t use the resource link rl_bank_iceberg in the Athena query editor to create or alter a view in the central account.

  1. Verify the added dialect by running a preview in Athena. For running the query, you can use either the resource link rl_bank_iceberg from the producer account catalog or use the centraladmin catalog.

The following screenshot shows querying using the resource link of the database in the producer account catalog.

The following screenshot shows querying the view from the producer using the connected catalog centraladmin as the data source.

  1. Verify the dialects on the view by inspecting the table in the Lake Formation console.

You can now query the view as the Data-Analyst role in the producer account, using both Athena and Spark. The view will also show in the central account as shown in the following code example, with access to the Lake Formation admin.

You can also create the view with ATHENA dialect and add the SPARK dialect. The SQL syntax to create the view in ATHENA dialect is shown in the following example:

create protected multi dialect view mdv_transaction1
security definer
as
SELECT t1.transaction_id, t1.transaction_type, t1.transaction_amount, t2.transaction_location, t2.transaction_date FROM transaction_table1 t1 
JOIN transaction_table2 t2 
ON t1.transaction_id = t2.transaction_id 
WHERE t1.transaction_amount > 100;

The update-table CLI to add the corresponding SPARK dialect is shown in the following example:

aws glue update-table --cli-input-json '{
    "DatabaseName": "rl_bankdatadb",
    "ViewUpdateAction": "ADD",
    "Force": true,
    "TableInput": {
        "Name": " mdv_transaction1",
        "StorageDescriptor": {
            "Columns": [
                {
                  "Name": "transaction_id",
                  "Type": "string"
                },
                {
                  "Name": "transaction_type",
                  "Type": "string"
                },
                {
                  "Name": "transaction_amount",
                  "Type": "float"
                },
                {
                  "Name": "transaction_location",
                  "Type": "string"
                },
                {
                  "Name": "transaction_date",
                  "Type": "date"
                }
             ],
             "SerdeInfo": {}
         },
         "ViewDefinition": {
         "SubObjects": [
               " "arn:aws:glue:<your-region>:<your-central-account-id>:table/bankdata_icebergdb/transaction_table1",
           "arn:aws:glue:<your-region>:<your-central-account-id>:table/bankdata_icebergdb/transaction_table2" 
],
         "IsProtected": true,
         "Representations": [
             {
                 "Dialect": "SPARK",
                 "DialectVersion": "1.0",
                 "ViewOriginalText": " SELECT t1.transaction_id, t1.transaction_type, t1.transaction_amount, t2.transaction_location, t2.transaction_date FROM transaction_table1 t1 JOIN transaction_table2 t2 ON t1.transaction_id = t2.transaction_id WHERE t1.transaction_amount > 100",
                 "ViewExpandedText": " SELECT t1.transaction_id, t1.transaction_type, t1.transaction_amount, t2.transaction_location, t2.transaction_date FROM transaction_table1 t1 JOIN transaction_table2 t2 ON t1.transaction_id = t2.transaction_id WHERE t1.transaction_amount > 100"
              }
           ]
        }
    }
}'

The following is a sample Python script to create a SPARK dialect view: glueview-createtable.py.

The following code block is a sample AWS Glue extract, transfer, and load (ETL) script to access the Spark dialect of the view from AWS Glue 5.0 from the central account. The AWS Glue job execution role should have Lake Formation SELECT permission on the AWS Glue view:

from pyspark.context import SparkContext
from pyspark.sql import SparkSession

aws_region = "<your-region>"
aws_account_id = "<your-central-account-id>"
local_catalogname = "spark_catalog"
warehouse_path = "s3://<your-bucket-name>/bankdata_icebergdb/transaction-table1"

spark = SparkSession.builder.appName('query_glue_view') \
    .config('spark.sql.extensions','org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions') \
    .config(f'spark.sql.catalog.{local_catalogname}', 'org.apache.iceberg.spark.SparkSessionCatalog') \
    .config(f'spark.sql.catalog.{local_catalogname}.catalog-impl', 'org.apache.iceberg.aws.glue.GlueCatalog') \
    .config(f'spark.sql.catalog.{local_catalogname}.client.region', aws_region) \
    .config(f'spark.sql.catalog.{local_catalogname}.glue.account-id', aws_account_id) \
    .config(f'spark.sql.catalog.{local_catalogname}.io-impl', 'org.apache.iceberg.aws.s3.S3FileIO') \
    .config(f'spark.sql.catalog.{local_catalogname}.warehouse',warehouse_path) \
    .getOrCreate()
spark.sql(f"show databases").show()
spark.sql(f"SHOW TABLES IN {local_catalogname}.bankdata_icebergdb").show()
spark.sql(f"SELECT * FROM {local_catalogname}.bankdata_icebergdb. mdv_transaction1").show()

In the AWS Glue job-details, for Lake Formation managed tables and for Iceberg tables, set additional parameters respectively as follows:

--enable-lakeformation-fine-grained-access = true
--datalake-formats = iceberg

Cleanup

To avoid incurring costs, clean up the resources you used for this post:

  1. Revoke the Lake Formation permissions granted to the Data-Analyst role and Producer account
  2. Drop the Athena tables
  3. Delete the Athena query results from your Amazon Simple Storage Service (Amazon S3) bucket
  4. Delete the Data-Analyst role from IAM

Conclusion

In this post, we demonstrated how to use cross-account IAM definer roles with AWS Glue Data Catalog views. We showed how data owner accounts can create and manage views in a central governance account while maintaining security and control over their data assets. This feature enables enterprises to implement sophisticated data mesh architectures without compromising on security or requiring data duplication.

The ability to use cross-account definer roles with Data Catalog views provides several key advantages:

  • Streamlines view management in multi-account environments
  • Maintains existing CI/CD workflows and automation
  • Enhances security through centralized governance
  • Reduces operational overhead by eliminating the need for data duplication

As organizations continue to build and scale their data lakehouse architectures across multiple AWS accounts, cross-account definer roles for Data Catalog views provide a crucial capability for implementing efficient, secure, and well-governed data sharing patterns.


About the authors

Aarthi Srinivasan

Aarthi Srinivasan

Aarthi is a Senior Big Data Architect at Amazon Web Services (AWS). She works with AWS customers and partners to architect data lake solutions, enhance product features, and establish best practices for data governance.

Sundeep Kumar

Sundeep Kumar

Sundeep is a Sr. Specialist Solutions Architect at Amazon Web Services (AWS), helping customers build data lake and analytics platforms and solutions. When not building and designing data lakes, Sundeep enjoys listening to music and playing guitar.

Simplify multi-warehouse data governance with Amazon Redshift federated permissions

Post Syndicated from Satesh Sonti original https://aws.amazon.com/blogs/big-data/simplify-multi-warehouse-data-governance-with-amazon-redshift-federated-permissions/

Modern data architectures increasingly rely on multi-warehouse deployments to achieve workload isolation, cost optimization, and performance scaling. Amazon Redshift federated permissions simplify permissions management across multiple Redshift warehouses.

With federated permissions, you register Redshift warehouse namespaces with the AWS Glue Data Catalog, creating a unified catalog that spans your entire warehouse fleet in the account. Registered namespaces are automatically mounted in every warehouse, providing data discovery without manual configuration. You can define permissions on database objects using familiar Redshift SQL commands, specifying global identities through AWS Identity and Access Management (IAM) or AWS IAM Identity Center (IDC). These permissions are stored alongside the warehouse data and enforced consistently, regardless of which warehouse runs the query. This provides a unified and secure access control model across your Redshift environment.

In this post, we show you how to define data permissions one time and automatically enforce them across warehouses in your AWS account, removing the need to re-create security policies in each warehouse.

Key capabilities of Amazon Redshift federated permissions

Federated permissions in Amazon Redshift offer the following key capabilities:

  • Global identity integration – Federated permissions use IAM and IAM Identity Center to provide single sign-on (SSO) across all registered warehouses. Users authenticate one time through their existing identity provider (IdP) and receive consistent access based on their global identity, regardless of which warehouse they connect to. This alleviates the need to create and manage separate user accounts in each warehouse, reducing administrative overhead and improving the user experience.
  • Unified catalog with automatic mounting – When you register a Redshift namespace with the Data Catalog using federated permissions, it becomes automatically visible in all warehouses within your account. Analysts using the Amazon Redshift Query Editor v2 or their preferred SQL client can discover and query tables across registered warehouses without manual catalog configuration. This automatic mounting capability simplifies data discovery and enables cross-warehouse analytics.
  • Consistent fine-grained access control – Row-level security (RLS) policies, dynamic data masking (DDM) policies, and column-level security (CLS) defined on warehouses using Amazon Redshift federated permissions automatically enforce when data is queried from consuming warehouses. You can implement advanced access controls—such as AWS Region-based row filtering, role-based masking for sensitive columns like SSN or credit card numbers, and time-based access restrictions—with confidence that these policies apply across warehouses.
  • SQL-based permission management – Federated permissions use familiar Redshift SQL syntax for permission management. You create RLS policies with CREATE RLS POLICY, attach them to tables and roles with ATTACH RLS POLICY, define masking policies with CREATE MASKING POLICY, and grant permissions with standard GRANT statements. This SQL interface enables infrastructure as code (IaC) approaches, supports database administrators to use their existing skills, and integrates naturally with existing extract, transform, and load (ETL) and automation workflows that use IAM or IAM Identity Center authentication.

Multi-warehouse architecture with federated permissions

The multi-warehouse architecture with federated permissions in Amazon Redshift represents a data mesh approach where multiple independent compute resources operate on shared data with unified governance. The following diagram illustrates the Redshift federated permissions setup process with the Data Catalog.

The process consists of the following steps:

  1. Each Redshift warehouse (1,2…N) registers with the Data Catalog. Refer onboarding documentation on registering the warehouse.
  2. After you register your Redshift warehouses with the Data Catalog, you can query data across your warehouses. Registered catalogs are automatically mounted in every warehouse in the account, appearing in the database explorer of Query Editor v2, and SQL clients connected to Amazon Redshift. To query a table in a registered catalog, use the three-part naming convention: database@catalog_name.schema_name.table_name.
  3. When you run a cross-catalog query, Amazon Redshift propagates your global identity (IAM role or IAM Identity Center user) to the remote warehouse. The remote warehouse’s catalog instance validates your permissions against the grants and fine-grained access control policies defined on the queried tables. If you have the necessary permissions, the table metadata and any applicable RLS, DDM, or CLS policies are returned to the consuming warehouse. Your local warehouse’s compute instance integrates these security policies into the query execution plan and runs the query on Redshift Managed Storage (RMS).

The enforcement of fine-grained access controls on remote data is a key differentiator of federated permissions. Traditional Redshift data sharing doesn’t support RLS or DDM policies on shared tables. With federated permissions, the security policies defined on the remote warehouse automatically apply when data is queried from any consumer warehouse. This supports compliance with data governance requirements without requiring administrators to duplicate security policies across warehouses.

The multi-warehouse architecture scales horizontally without increasing governance complexity. When you add a new warehouse to your account and register it with federated permissions, it automatically inherits the appropriate permission model without manual configuration. Analysts connecting to the new warehouse immediately see all databases they have access to across the mesh, and all security policies apply automatically. This alleviates the N-squared problem of managing permissions across N warehouses, reducing the administrative burden from N separate configurations to a single unified governance model.

Query lifecycle

The following diagram illustrates the step-by-step flow of how a user query on Redshift Warehouse 1 accesses objects in Redshift Warehouse N with federated permissions.

Note: Steps 2, 3, and 4 will be skipped if permission details are available in the local cache

The workflow consists of the following steps:

  1. The user connects to Redshift Warehouse 1 and queries a table in Federated Catalog N.
  2. Redshift Warehouse 1 calls the Data Catalog GetTable API. This request includes the user’s token.
  3. The request routes to Redshift Warehouse N.
  4. Redshift Warehouse N verifies the user permissions. If it’s authorized, it returns the table metadata and security policy details such as RLS policies, DDM rules, and CLS settings.
  5. Redshift Warehouse 1 applies the security policies in the query plan and runs the query against Redshift Managed Storage (RMS), where Redshift stores data in an optimized format.
  6. The results are returned to the user.

Solution overview

The example in this post demonstrates how to define RLS and DDM policies on a data warehouse and verify that these policies are enforced when querying from another data warehouse.

We will create a table with credit card data and apply RLS and DDM policies to limit consumer cards data and mask credit card values for non-admin users. These policies will be applied across all the data warehouses consistently and mask the credit card details when non-admin users query the table.

Prerequisites

Create the following IAM roles:

Create table and load data

Run following steps to create a credit_card table and load sample data.

  1. Connect to the first Redshift data warehouse1 using the IAM Aadmin role
  2. Create a credit_cards table
    -- Create table
    CREATE TABLE credit_cards (
      customer_id INT,
      credit_card varchar(16),
      card_type varchar(10)
    );

  3. Insert sample data
    -- Insert sample data
    INSERT INTO credit_cards
    VALUES
      (100, '4532993817514842', 'consumer'),
      (100, '4716002041425888', 'corporate'),
      (102, '5243112427642649', 'consumer'),
      (102, '6011720771834675', 'consumer'),
      (102, '6011378662059710', 'corporate'),
      (103, '373611968625635', 'consumer');

Apply RLS and DDM policies

Run following steps to create and apply RLS and DDM policies.

  1. Create an RLS policy to filter only consumer card types:
    -- Create RLS policy
    CREATE RLS POLICY consumer_cards
    WITH (card_type VARCHAR(10))
    USING (card_type = 'consumer');

  2. Create a DDM policy that masks credit cards:
    -- Create masking policy
    CREATE MASKING POLICY mask_credit_card_full
    WITH (credit_card VARCHAR(256))
    USING ('000000XXXX0000'::TEXT);

  3. Attach RLS and DDM Policies to RedOnly role
    -- Attach RLS and DDM policies to ReadOnly role
    ATTACH RLS POLICY consumer_cards 
    ON credit_cards 
    TO "IAMR:ReadOnly";
    
    ATTACH MASKING POLICY mask_credit_card_full
    ON credit_cards(credit_card)
    TO "IAMR:ReadOnly";

  4. Enable Row Level Security on the table
    ALTER TABLE credit_cards ROW LEVEL SECURITY ON;

  5. Grant select on the table to Readonly role
    GRANT SELECT ON credit_cards TO "IAMR:ReadOnly";

Connect to data warehouse 2 as read-only user

Run following steps on data warehouse 2 to query the data.

  1. Connect to data warehouse 2 as a read-only user and expand the external databases. The following screenshot shows an example using Query Editor V2.
  2. Notice the credit_cards table from data warehouse 1 when you expand the catalog.
  3. Run the following SQL to query the table. Replace rs-demo-dw1 in the following SQL with the catalog name you gave while registering data warehouse 1:
    -- SQL to query credit cards table in data warehouse1. 
    SELECT * FROM "dev@rs-demo-dw1"."public"."credit_cards";

  4. You should see only consumer type credit cards with card details masked in the output. The RLS and DDM policies applied in data warehouse 1 on the IAMR:ReadOnly user are enforced even though you queried the table from a different data warehouse.
    The following screenshot shows an example output.
  5. For auditing, you can run SHOW commands to view the policies applied on the tables for the roles:
    -- Show all RLS policies in the database.
    SHOW RLS POLICIES FROM DATABASE "dev@rs-demo-dw1";
    -- Show all masking policies in the database.
    SHOW MASKING POLICIES FROM DATABASE "dev@rs-demo-dw1";

This example demonstrates the power of federated permissions: security policies defined one time on a warehouse automatically enforce across your warehouses, maintaining compliance without duplicating policy definitions.

Considerations

Keep in mind the following when using federated permissions:

Clean up

To avoid incurring future charges, delete the resources you created, including the Redshift data warehouses and IAM roles.

Conclusion

Amazon Redshift federated permissions transform multi-warehouse data governance into a streamlined, automated process. For organizations operating multiple Redshift warehouses, federated permissions deliver immediate value by reducing administrative time and supporting consistent security enforcement. The familiar SQL interface and backward compatibility with existing Redshift permissions enable rapid adoption without requiring teams to learn new governance models.

The integration with IAM and IAM Identity Center provides enterprise-grade identity management with SSO capabilities, and the automatic mounting of registered catalogs simplifies data discovery and cross-warehouse analytics. If you are currently using Amazon Redshift local permissions, refer to the tool described in Modernize Amazon Redshift authentication by migrating user management to AWS IAM Identity Center.

To learn more and get started, see Amazon Redshift Federated Permissions documentation.


About the authors

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 20 years of experience in building data assets and leading complex data platform programs for banking and insurance clients across the globe.

Sandeep Adwankar

Sandeep Adwankar

Sandeep is a Senior Product Manager with Amazon SageMaker Lakehouse . Based in the California Bay Area, he works with customers around the globe to translate business and technical requirements into products that help customers improve how they manage, secure, and access data.

Abhishek Rai Sharma

Abhishek Rai Sharma

Abhishek is a Senior Software Engineer focused on Amazon Redshift Catalog and Governance. He is passionate about creating reliable, scalable infrastructure solutions for distributed analytics workloads and enterprise data mesh architectures.

Ramchandra Anil Kulkarni

Ramchandra Anil Kulkarni

Anil is a Senior Software Engineer at Amazon Redshift with expertise in the Governance and Query Processing areas. He is passionate about distributed systems and solving impactful problems for AWS customers.

Ning Di

Ning Di

Ning is a Senior Software Development Engineer at Amazon Redshift, driven by a genuine passion for exploring all aspects of technology.

Unifying governance and metadata across Amazon SageMaker Unified Studio and Atlan

Post Syndicated from Karan Singh Thakur, Satabrata Paul original https://aws.amazon.com/blogs/big-data/unifying-governance-and-metadata-across-amazon-sagemaker-unified-studio-and-atlan/

This post was cowritten with Satabrata Paul and Karan Singh Thakur from Atlan

In this post, we show you how to unify governance and metadata across Amazon SageMaker Unified Studio and Atlan through a comprehensive bidirectional integration. You’ll learn how to deploy the necessary Amazon Web Services (AWS) infrastructure, configure secure connections, and set up automated synchronization to maintain consistent metadata across both platforms.

As organizations scale their data and AI programs, teams often work across distributed tools such as governance solutions for business users and analytics or machine learning (ML) environments for technical teams. Without tight integration between these systems, metadata becomes fragmented. A single asset can appear under different names, documentation might drift out of sync, and governance signals can become inconsistent across systems.

To address these challenges, Atlan, a modern data workspace that makes collaboration among diverse users like business, analysts, and engineers easier, increasing efficiency and agility in data projects, and AWS have built a bidirectional integration between Atlan and Amazon SageMaker Unified Studio. This integration creates a continuous connection between both environments so every team within the enterprise can work with a single, trusted, and synchronized view of metadata for their data and AI assets. By bridging the gap between diverse users collaborating in Atlan and technical teams working within Amazon SageMaker Unified Studio for analytics and ML, this integration maintains consistency across both platforms without requiring teams to switch contexts or manually reconcile metadata differences.

Why unified metadata governance matters

Enterprises today operate in hybrid environments. Business users rely on Atlan as an active metadata solution to manage, govern, and collaborate on data assets across the modern data stack. Atlan helps teams find, understand, and trust their data so they can use it effectively to drive business outcomes.

Organizations also use Amazon SageMaker Catalog to simplify the discovery, governance, and collaboration for both business and technical data across structured and unstructured sources. Teams can use the catalog to organize data products, capture context, and apply governance policies consistently within Amazon SageMaker Unified Studio.

This new integration synchronizes metadata between SageMaker Catalog and Atlan, maintaining consistency and keeping content current across both environments. With a unified view, every team within the enterprise can work confidently with a single, trusted representation of their data and AI assets.

Solution overview

The solution follows a phased rollout strategy to provide you with immediate value while progressively expanding toward comprehensive data and AI governance capabilities. The current phase focuses on establishing secure, scalable, and reliable metadata synchronization between Atlan and Amazon SageMaker Unified Studio.

The Phase 1 integration between Amazon SageMaker Catalog and Atlan enables both on-demand and scheduled bidirectional metadata synchronization across the two solutions. It uses the standard APIs of Amazon SageMaker Unified Studio and Atlan to create a scalable and configurable mechanism for metadata exchange. Key capabilities include:

  • Secure connection using IAM roles – The integration is established through a controlled AWS Identity and Access Management (IAM) based handshake. A predefined AWS CloudFormation template automatically provisions the IAM role and policies required to enable a secure, least-privilege connection between Amazon SageMaker Catalog and the Atlan application.
  • On-demand and scheduled synchronization – The integration supports both manual and automated metadata synchronization. API-driven workflows manage the exchange of glossary terms, asset descriptions, and classifications in both directions, keeping metadata consistent across systems.

After you’ve implemented Phase 1, you can perform bidirectional synchronization of glossary terms and descriptions between Amazon SageMaker Unified Studio and Atlan. This keeps your terminology consistent across both platforms, and your teams can maintain a single source of truth for business definitions. The integration also preserves your glossary structures, including parent-child relationships, so your carefully organized taxonomy remains intact during the sync process. Additionally, glossary terms are automatically associated with related data assets, saving you the manual effort of linking terms to the appropriate datasets and reducing the risk of inconsistencies.

Beyond glossary management, Phase 1 enables comprehensive ingestion of assets and metadata from Amazon SageMaker Unified Studio into Atlan. This includes your projects, both published and subscribed assets, domains and data products, glossaries and terms, metadata forms, and column descriptions. By bringing this information into Atlan, you create a unified view of your data landscape that makes it easier for data consumers to discover, understand, and trust the data they’re working with.

Prerequisites

To follow along with this integration setup, you must have the following resources already configured in your environment:

  • An Atlan tenant
  • A Node group IAM role
  • An Amazon SageMaker Unified Studio domain.
  • At least one Amazon SageMaker Unified Studio project with assets created and glossary terms defined.
  • Atlan API Token. You can generate this by navigating to API access under the Atlan’s Admin center.
  • Atlan top-level glossary. You can create this glossary container on Atlan to ingest SageMaker Unified Studio glossaries and terms.

The next section offers a step-by-step walkthrough of the integration, from initial setup to full operation. It demonstrates how you can establish the trust handshake between Amazon SageMaker Unified Studio and Atlan and how bidirectional synchronization functions in practice.

Setup on AWS

To begin the integration, you need Atlan’s Account Node Instance IAM role. This role allows the Atlan SageMaker Unified Studio application to securely assume the IAM role that you will create in your AWS account using an AWS CloudFormation template. The trust relationship between these two roles authorizes Atlan to publish metadata to Amazon SageMaker Catalog and to perform reverse synchronization from AWS back into Atlan.

The IAM policy follows the principle of least privilege, granting Atlan access only to the resources necessary for cataloging and governance. This approach maintains accurate metadata synchronization while preserving your existing cloud security and compliance controls.

Follow AWS best practices when configuring trust relationships. These cross-account access mechanisms require careful management and monitoring, particularly during security incidents. For comprehensive guidance on securing IAM roles and trust policies, refer to the Security best practices in IAM and Require workloads to use temporary credentials with IAM roles to access AWS.

Contact your Atlan administrator to obtain the Amazon Resource Name (ARN) of the Atlan Account Node Instance IAM role. You will need this value when configuring the CloudFormation stack in AWS.

The next step is to create an AWS IAM role using the provided CloudFormation template. This role establishes the trust relationship between your Amazon SageMaker Unified Studio environment and your Atlan tenant. Follow these steps:

  1. Access the CloudFormation template. The CloudFormation template is currently available as a YAML file.
  2. On the AWS Management Console, navigate to CloudFormation and choose Create stack, then choose With new resources (standard), as shown in the following screenshot.
  3. Choose the provided CloudFormation template and choose Next.
  4. Enter a name for the stack and complete the required parameters, as shown in the following screenshot:
    1. AtlanNodeInstanceRoleArn – The ARN of the Atlan node instance role.
    2. SMUSDomainId – The unique identifier for the SageMaker Unified Studio domain.
    3. SMUSProjectsToSync – The project IDs where SageMaker Unified Studio and Atlan synchronization will be enabled. You can choose to either add the project IDs and keep updating this stack every time a Project is added or add the created IAM role to each project as owner.

  5. Select the acknowledgement checkbox and choose Next, as shown in the following screenshot.
  6. Choose Submit to start the stack deployment. When the process is complete, the stack status will update to CREATE_COMPLETE.
  7. Note the IAM role ARN
  8. After the CloudFormation stack has been deployed and the IAM role has been created, copy the IAM Role ARN from the CloudFormation output. You will need this value during the configuration process on the Atlan side to establish the secure connection between your Amazon SageMaker Unified Studio environment and your Atlan tenant.

Setup on Atlan

Now that you’ve deployed the necessary AWS resources, you’ll configure Atlan to establish the connection with Amazon SageMaker Unified Studio. This involves setting up the API token, configuring the IAM role, and creating the glossary container that will receive your synchronized metadata. Follow these steps:

  1. Sign in to your Atlan tenant, as shown in the following screenshot.
  2. On the New dropdown menu, choose New workflow.
  3. On the Marketplace tab, search for and select the AWS SageMaker Unified Studio app, as shown in the following screenshot.
  4. Enter credential details. Use the IAM role or user created by the CloudFormation template before, enter an API token, and choose your AWS Region, as shown in the following screenshot.
  5. Enter connection details. In Connection name, enter a name. Under Connection Admins, choose the plus icon to add members (other users) to the connectors as admins. Assigning admin permissions to the connection allows these users to:
    1. View and edit the assets in the connection.
    2. Edit connection preferences.
    3. Edit persona-based policies for the connection.

  6. Choose metadata filters and preflight checks, as shown in the following screenshot:
    • In the Select Glossary to enrich dropdown menu, choose the glossary container in Atlan to be enriched with glossaries and terms from Atlan.
    • To check for necessary permissions required to run the workflow, select Quick test for necessary permissions before workflow run.
    • To run the workflow, choose Run. To schedule it to run later, choose Schedule & Run.

Synchronization of metadata

Now that you’ve configured the integration between Atlan and Amazon SageMaker Unified Studio, let’s explore how metadata flows bidirectionally between both platforms to maintain consistency and governance across your data landscape.

The Atlan SageMaker Unified Studio connector uses a bidirectional synchronization model that keeps business context and technical metadata consistent across both solutions. The process delivers reliability, traceability, and governance-safe updates, regardless of where changes originate. The following diagram illustrates the solution architecture.

Sequential workflow for the SageMaker Unified Studio Atlan integration

The integration between SageMaker Unified Studio and Atlan follows a carefully orchestrated sequential workflow that enables seamless metadata synchronization across both platforms.

The process begins with connection setup through IAM, where authentication and authorization are configured to establish secure access between the customer’s AWS account and Atlan’s AWS environment. This foundational security layer allows subsequent data exchanges to occur within a trusted framework.

After the connection is established, the metadata sync workflow can be triggered either on a defined schedule or manually by the user, providing flexibility based on organizational needs. When triggered, the Atlan SageMaker Unified Studio app calls the SageMaker Unified Studio APIs to ingest assets and metadata from the source system.

The ingested assets then undergo processing and transformation within Atlan, where they are converted into Atlan’s metadata model. This processing step is crucial because it makes the assets discoverable, searchable, and governable inside the Atlan platform, which means teams can use Atlan’s full governance capabilities.

A key capability of this integration is its real-time reverse sync for metadata updates. When a user modifies metadata for the assets inside Atlan (such as adding tags or updating descriptions), Atlan’s real-time reverse sync pipelines immediately detect these changes and push the updates back to SageMaker Unified Studio. This keeps SageMaker Unified Studio reflecting the most up-to-date metadata entered by users in Atlan, eliminating the risk of metadata drift between systems.

This bidirectional sync creates a continuous loop where metadata flows from SageMaker Unified Studio to Atlan for ingestion and publication, simultaneously flowing back from Atlan to SageMaker Unified Studio through real-time reverse sync. The result is a consistent, bidirectional metadata flow that keeps both platforms synchronized. Teams can work confidently knowing that their metadata governance efforts are reflected across their data.

The following diagram illustrates this complete workflow, showing how metadata moves through each stage of the integration from initial IAM authentication through the continuous bidirectional sync loop that maintains metadata consistency across both platforms.

SageMaker Unified Studio to Atlan: Ingestion of metadata

The Atlan-SageMaker Unified Studio App periodically connects to SageMaker Unified Studio using secure API calls to ingest metadata. This metadata is transformed and mapped into Atlan’s metadata model, then published through the Atlan publish app as new or updated assets.

Each ingestion cycle is fully logged by Atlan’s audit service, which captures timestamps, correlation IDs, and the full change record. These logs support deduplication, troubleshooting, and replay in the event of partial failures.

Atlan to SageMaker Unified Studio: Synchronizing enriched business context

When users enrich assets inside Atlan, for example by updating descriptions or attaching glossary terms, the integration detects these changes and selectively pushes them back to SageMaker Unified Studio.

The reverse sync control plane is a pipeline that automatically detects changes made to assets and then triggers SageMaker Unified Studio Update API calls in the background to keep everything synchronized.

What’s next?

Phase 1 delivers core metadata synchronization and principal catalog selection for immediate consistency across your data governance platforms. Phase 2 will synchronize lineage and data quality, so teams see the same data flows and quality signals in both Atlan and SageMaker Catalog, enabling end-to-end visibility into how data moves through your pipelines and maintaining quality metrics consistently tracked across both systems. Phase 3 will add integrated approval workflows to streamline how access is requested and granted across solutions, reducing friction for data consumers while maintaining robust governance controls. These upcoming phases build toward a fully connected governance experience, keeping metadata, lineage, quality, and access policies aligned across the modern data stack.

Cleanup

If you no longer need the SageMaker Unified Studio connector integration, complete the following steps to clean up your environment and avoid unintended resource usage:

  1. Delete the CloudFormation stack. Navigate to the AWS CloudFormation console, locate the stack deployed for this solution, and choose Delete. This action removes the AWS resources provisioned by the stack, including IAM roles, policies, and supporting components.
  2. Remove the connection in Atlan. Visit Delete a connection to follow the steps outlined in Atlan’s documentation to delete the associated connection.

Cleaning up these components keeps your AWS and Atlan environments streamlined, secure, and cost-efficient.

Conclusion

In this post, you learned how to establish a bidirectional integration between Atlan and Amazon SageMaker Unified Studio that unifies metadata governance across your data and AI environments. You walked through deploying the necessary AWS infrastructure using CloudFormation, configuring the secure IAM based connection, and setting up bidirectional synchronization to keep glossary terms, descriptions, and governance context aligned across both platforms.

Organizations can use this integration to connect business and technical users within a single governance framework, creating a consistent, trusted view of data across the enterprise. With one secure configuration, teams can synchronize metadata between Atlan and Amazon SageMaker Unified Studio, establishing a reliable foundation for innovation, collaboration, and responsible AI at scale.


About the authors

Karan Singh Thakur

Karan is a Senior Product Manager at Atlan, leading the strategy and execution for deep hyperscaler integrations, especially across AWS. Before Atlan, Karan spent over a decade building cloud-based, data-intensive environments, including serving as the founding PM for a fully managed lakehouse engine and leading enterprise analytics, governance, and Kubernetes-based workload systems.

Satabrata Paul

Satabrata Paul

Satabrata is a Senior Software Engineer on Atlan’s Metadata Marketplace team, where he designs and scales backend systems and CI/CD workflows for high-quality metadata connector integrations. Focused on modern data environments, he helps teams streamline asset discovery, lineage, and cataloging across complex environments.

Divij Bhatia

Divij Bhatia

Divij is a Software Development Engineer at Amazon Web Services (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 Amazon Web Services (AWS). He has over a decade of experience in data management, helping customers around the globe address their business and technical needs.

Modernize Apache Spark workflows using Spark Connect on Amazon EMR on Amazon EC2

Post Syndicated from Philippe Wanner original https://aws.amazon.com/blogs/big-data/modernize-apache-spark-workflows-using-spark-connect-on-amazon-emr-on-amazon-ec2/

Apache Spark Connect, introduced in Spark 3.4, enhances the Spark ecosystem by offering a client-server architecture that separates the Spark runtime from the client application. Spark Connect enables more flexible and efficient interactions with Spark clusters, particularly in scenarios where direct access to cluster resources is limited or impractical.

A key use case for Spark Connect on Amazon EMR is to be able to connect directly from your local development environments to Amazon EMR clusters. By using this decoupled approach, you can write and test Spark code on your laptop while using Amazon EMR clusters for execution. This capability reduces development time and simplifies data processing with Spark on Amazon EMR.

In this post, we demonstrate how to implement Apache Spark Connect on Amazon EMR on Amazon Elastic Compute Cloud (Amazon EC2) to build decoupled data processing applications. We show how to set up and configure Spark Connect securely, so you can develop and test Spark applications locally while executing them on remote Amazon EMR clusters.

Solution architecture

The architecture centers on an Amazon EMR cluster with two node types. The primary node hosts both the Spark Connect API endpoint and Spark Core components, serving as the gateway for client connections. The core node provides additional compute capacity for distributed processing. Although this solution demonstrates the architecture with two nodes for simplicity, it scales to support multiple core and task nodes based on workload requirements.

In Apache Spark Connect version 4.x, TLS/SSL network encryption is not inherently supported. We show you how to implement secure communications by deploying an Amazon EMR cluster with Spark Connect on Amazon EC2 using an Application Load Balancer (ALB) with TLS termination as the secure interface. This approach enables encrypted data transmission between Spark Connect clients and Amazon Virtual Private Cloud (Amazon VPC) resources.

The operational flow is as follows:

  1. Bootstrap script – During Amazon EMR initialization, the primary node fetches and executes the start-spark-connect.sh file from Amazon Simple Storage Service (Amazon S3). This script starts the Spark Connect server.
  2. Server availability – When the bootstrap process is complete, the Spark Server enters a waiting state, ready to accept incoming connections. The Spark Connect API endpoint becomes available on the configured port (typically 15002), listening for gRPC connection from remote clients.
  3. Client interaction – Spark Connect clients can establish secure connections to an Application Load Balancer. These clients translate DataFrame operations into unresolved logical query plans, encode these plans using protocol buffers, and send them to the Spark Connect API using gRPC.
  4. Encryption in transit – The Application Load Balancer receives incoming gRPC or HTTPS traffic, performs TLS termination (decrypting the traffic), and forwards the requests to the primary node. The certificate is stored in AWS Certificate Manager (ACM).
  5. Request processing – The Spark Connect API receives the unresolved logical plans, translates them into Spark’s built-in logical plan operators, passes them to Spark Core for optimization and execution, and streams results back to the client as Apache Arrow-encoded row batches.
  6. (Optional) Operational access – Administrators can securely connect to both primary and core nodes through Session Manager, a capability of AWS Systems Manager, enabling troubleshooting and maintenance without exposing SSH ports or managing key pairs.

The following diagram depicts the architecture of this post’s demonstration for submitting Spark unresolved logical plans to EMR clusters using Spark Connect.

Apache Spark Connect on Amazon EMR solution architecture diagram

Apache Spark Connect on Amazon EMR solution architecture diagram

Prerequisites

To proceed with this post, ensure you have the following:

Implementation steps

In this recipe, through AWS CLI commands, you will:

  1. Prepare the bootstrap script, a bash script starting Spark Connect on Amazon EMR.
  2. Set up the permissions for Amazon EMR to provision resources and perform service-level actions with other AWS services.
  3. Create the Amazon EMR cluster with these associated roles and permissions and eventually attach the prepared script as a bootstrap action.
  4. Deploy the Application Load Balancer and certificate with ACM secure data in transit over the internet.
  5. Modify the primary node’s security group to allow Spark Connect clients to connect.
  6. Connect with a test application connecting the client to Spark Connect server.

Prepare the bootstrap script

To prepare the bootstrap script, follow these steps:

  1. Create an Amazon S3 bucket to host the bootstrap bash script:
    REGION=
    BUCKET_NAME=
    aws s3api create-bucket \
       --bucket $BUCKET_NAME \ 
       --region $REGION \
       --create-bucket-configuration LocationConstraint=$REGION

  2. Open your preferred text editor, add the following commands in a new file with a name such start-spark-connect.sh. If the script runs on the primary node, it starts Spark Connect server. If it runs on a task or core node, it does nothing:
    #!/bin/bash
    if grep isMaster /mnt/var/lib/info/instance.json | grep false;
    then
        echo "This is not master node, do nothing."
        exit 0
    fi
    echo "This is master, continuing to execute script"
    SPARK_HOME=/usr/lib/spark
    SPARK_VERSION=$(spark-submit --version 2>&1 | grep "version" | head -1 | awk '{print $NF}' | grep -oE '[0-9]+\.[0-9]+\.[0-9]+')
    SCALA_VERSION=$(spark-submit --version 2>&1 | grep -o "Scala version [0-9.]*" | awk '{print $3}' | grep -oE '[0-9]+\.[0-9]+')
    echo "Spark version ${SPARK_VERSION} is installed under ${SPARK_HOME} running with scala version ${SCALA_VERSION}"
    sudo "${SPARK_HOME}"/sbin/start-connect-server.sh --packages org.apache.spark:spark-connect_"${SCALA_VERSION}:${SPARK_VERSION}"

  3. Upload the script into the bucket created in step 1:
    aws s3 cp start-spark-connect.sh s3://$BUCKET_NAME
    

Set up the permissions

Before creating the cluster, you must create the service role, and instance profile. A service role is an IAM role that Amazon EMR assumes to provision resources and perform service-level actions with other AWS services. An EC2 instance profile for Amazon EMR assigns a role to every EC2 instance in a cluster. The instance profile must specify a role that can access the resources for your bootstrap action.

  1. Create the IAM role:
    aws iam create-role \
    --role-name AmazonEMR-ServiceRole-SparkConnectDemo \
    --assume-role-policy-document '{
    	"Version": "2012-10-17",
    	"Statement": [{
    		"Effect": "Allow",
    		"Principal": {"Service": "elasticmapreduce.amazonaws.com"},
    		"Action": "sts:AssumeRole"
    		}]
    }'
    

  2. Attach the necessary managed policies to the service role to allow Amazon EMR to manage the underlying services Amazon EC2 and Amazon S3 on your behalf and optionally grant an instance to interact with Systems Manager:
    aws iam attach-role-policy \
    --role-name AmazonEMR-ServiceRole-SparkConnectDemo \
    --policy-arn arn:aws:iam::aws:policy/service-role/AmazonEMRServicePolicy_v2
    
    aws iam attach-role-policy \
    --role-name AmazonEMR-ServiceRole-SparkConnectDemo \
    --policy-arn arn:aws:iam::aws:policy/AmazonSSMManagedInstanceCore
    
    aws iam attach-role-policy \
    --role-name AmazonEMR-ServiceRole-SparkConnectDemo \
    --policy-arn arn:aws:iam::aws:policy/service-role/AmazonElasticMapReduceRole
    

  3. Create an Amazon EMR instance role to grant permissions to EC2 instances to interact with Amazon S3 or other AWS services:
    aws iam create-role \
    --role-name EMR_EC2_SparkClusterNodesRole \
    --assume-role-policy-document '{
    "Version": "2012-10-17",
    "Statement": [{
       "Effect": "Allow",
       "Principal": {"Service": "ec2.amazonaws.com"},
       "Action": "sts:AssumeRole"
       }]
    }'
    

  4. To allow the primary instance to read from Amazon S3, attach the AmazonS3ReadOnlyAccess policy to the Amazon EMR instance role. For production environments, this access policy should be reviewed and replaced with a custom policy following the principle of least privilege, granting only the specific permissions needed for your use case:
    aws iam attach-role-policy \
    --role-name EMR_EC2_SparkClusterNodesRole \
    --policy-arn arn:aws:iam::aws:policy/AmazonS3ReadOnlyAccess
    

  5. Attaching AmazonSSMManagedInstanceCore policy enables the instances to use core Systems Manager features, such as Session Manager, and Amazon CloudWatch:
    aws iam attach-role-policy \
    --role-name EMR_EC2_SparkClusterNodesRole \
    --policy-arn arn:aws:iam::aws:policy/AmazonSSMManagedInstanceCore
    

  6. To pass the EMR_EC2_SparkClusterInstanceProfile IAM role information to the EC2 instances when they start, create the Amazon EMR EC2 instance profile:
    aws iam create-instance-profile \
    --instance-profile-name EMR_EC2_SparkClusterInstanceProfile
    

  7. Attach the role EMR_EC2_SparkClusterNodesRole created in step 3 to the newly instance profile:
    aws iam add-role-to-instance-profile \
    --instance-profile-name EMR_EC2_SparkClusterInstanceProfile \
    --role-name EMR_EC2_SparkClusterNodesRole
    

Create the Amazon EMR cluster

To create the Amazon EMR cluster, follow these steps:

  1. Set the environment variables, where your EMR cluster and load-balancer must be deployed:
    VPC_ID=<vpc-emr-and-alb>
    EMR_PRI_SB_ID_1=<emr-private-subnet-id-az1>
    ALB_PUB_SB_ID_1=<alb-public-subnet-id-az1>
    ALB_PUB_SB_ID_2=<alb-public-subnet-id-az2>
    

  2. Create the EMR cluster with the latest Amazon EMR release. Replace the placeholder value with your actual S3 bucket name where the bootstrap action script is stored:
    CLUSTER_ID=$(aws emr create-cluster \
    --name "Spark Connect cluster demo" \
    --applications Name=Spark \
    --release-label emr-7.9.0 \
    --service-role AmazonEMR-ServiceRole-SparkConnectDemo \
    --ec2-attributes InstanceProfile=EMR_EC2_SparkClusterInstanceProfile,SubnetId=$EMR_PRI_SB_ID_1 \
    --instance-groups InstanceGroupType=MASTER,InstanceCount=1,InstanceType=m5.xlarge InstanceGroupType=CORE,InstanceCount=1,InstanceType=m5.xlarge \
    --bootstrap-actions Path="s3://$BUCKET_NAME/start-spark-connect.sh" \
    --query 'ClusterId' --output text)
    echo CLUSTER_ID="$CLUSTER_ID"
    

    To modify primary node’s security group to allow Systems Manager to start a session.

  3. Get the primary node’s security group identifier. Record the identifier because you’ll need it for subsequent configuration steps in which primary-node-security-group-id is mentioned:
    PRIMARY_NODE_SG=$(aws emr describe-cluster \
    --cluster-id $CLUSTER_ID \
    --query 'Cluster.Ec2InstanceAttributes.EmrManagedMasterSecurityGroup' \
    --output text)
    echo PRIMARY_NODE_SG=$PRIMARY_NODE_SG
    

  4. Find the EC2 instance connect prefix list ID for your Region. You can use the EC2_INSTANCE_CONNECT filter with the describe-managed-prefix-lists command. Using a managed prefix list provides a dynamic security configuration to authorize Systems Manager EC2 instances to connect the primary and core nodes by SSH:
    IC_PREFIX_LIST=$(aws ec2 describe-managed-prefix-lists \
    --filters Name=prefix-list-name,Values=com.amazonaws.$REGION.ec2-instance-connect \
    --query 'PrefixLists[0].PrefixListId' \
    --output text)
    echo IC_PREFIX_LIST=$IC_PREFIX_LIST
    

  5. Modify the primary node security group inbound rules to allow SSH access (port 22) to the EMR cluster’s primary node from resources that are part of the specified Instance Connect service contained in the prefix list:
    aws ec2 authorize-security-group-ingress \
    --region $REGION \
    --group-id $PRIMARY_NODE_SG \
    --ip-permissions "[{\"IpProtocol\":\"tcp\",\"FromPort\":22,\"ToPort\":22,\"PrefixListIds\":[{\"PrefixListId\":\"$IC_PREFIX_LIST\"}]}]"
    

Optionally, you can repeat the preceding steps 1–3 for the core (and tasks) cluster’s nodes to allow Amazon EC2 Instance Connect to access the EC2 instance through SSH.

Deploy the Application Load Balancer and certificate

To deploy the Application Load Balancer and certificate, follow these steps:

  1. Create a load balancer’s security group:
    ALB_SG_ID=$(aws ec2 create-security-group \
    --group-name spark-connect-alb-sg \
    --description "Security group for Spark Connect ALB" \
    --region $REGION \
    --vpc-id $VPC_ID \
    --query 'GroupId' \
    --output text)
    

  2. Add rule to accept TCP traffic from a trusted IP on port 443. We recommend that you use the local development machine’s IP address. You can check your current public IP address here: https://checkip.amazonaws.com:
    aws ec2 authorize-security-group-ingress \
    --group-id $ALB_SG_ID \
    --protocol tcp \
    --port 443 \
    --cidr <replace-with-trusted-IP>/32
    

  3. Create a new target group with gRPC protocol, which targets the Spark Connect server instance and the port the server is listening to:
    ALB_TG_ARN=$(aws elbv2 create-target-group \
    --name spark-connect-tg \
    --protocol HTTP \
    --protocol-version GRPC \
    --port 15002 \
    --target-type instance \
    --health-check-enabled \
    --health-check-protocol HTTP \
    --health-check-path / \
    --vpc-id $VPC_ID \
    --query 'TargetGroups[0].TargetGroupArn' \
    --output text)
    echo "ALB TG created (ARN)=$ALB_TG_ARN"
    

  4. Create the Application Load Balancer:
    ALB_ARN=$(aws elbv2 create-load-balancer \
    --name spark-connect-alb \
    --type application \
    --scheme internet-facing \
    --subnets $ALB_PUB_SB_ID_1 $ALB_PUB_SB_ID_2 \
    --security-groups $ALB_SG_ID \
    --query 'LoadBalancers[0].LoadBalancerArn' \
    --output text)
    echo "ALB created (ARN)=$ALB_ARN"
    

  5. Get the load balancer DNS name:
    ALB_DNS=$(aws elbv2 describe-load-balancers \
    --load-balancer-arns $ALB_ARN \
    --query 'LoadBalancers[0].DNSName' \
    --output text)
    echo "ALB DNS=$ALB_DNS"
    

  6. Retrieve the Amazon EMR primary node ID:
    PRIMARY_NODE_ID=$(aws emr list-instances --cluster-id $CLUSTER_ID --instance-group-types MASTER --query 'Instances[0].Ec2InstanceId' --output text)
    echo PRIMARY_NODE_ID=$PRIMARY_NODE_ID
    

  7. (Optional) To encrypt and decrypt the traffic, the load balancer needs a certificate. You can skip this step if you already have a trusted certificate in ACM. Otherwise, create a self-signed certificate:
    PRIVATE_KEY_PATH=./sc-private-key.key
    CERTIFICATE_PATH=./sc-certificate.cert
    sudo openssl req -x509 -nodes -days 365 -newkey rsa:2048 -keyout $PRIVATE_KEY_PATH -out $CERTIFICATE_PATH -subj "/CN=$ALB_DNS"
    

  8. Upload to ACM:
    ACM_CERT_ARN=$(aws acm import-certificate \
    --certificate fileb://$CERTIFICATE_PATH \
    --private-key fileb://$PRIVATE_KEY_PATH \
    --region $REGION \
    --query CertificateArn \
    --output text)
    echo "Certificate created (ARN)=$ACM_CERT_ARN"
    

  9. Create the load balancer listener:
    ALB_LISTENER_ARN=$(aws elbv2 create-listener \
    --load-balancer-arn $ALB_ARN \
    --protocol HTTPS \
    --port 443 \
    --certificates CertificateArn=$ACM_CERT_ARN \
    --ssl-policy ELBSecurityPolicy-TLS13-1-2-2021-06 \
    --default-actions Type=forward,TargetGroupArn=$ALB_TG_ARN \
    --region $REGION \
    --query 'Listeners[0].ListenerArn' \
    --output text)
    echo "ALB listener created (ARN)=$ALB_LISTENER_ARN"
    

  10. After the listener has been provisioned, register the primary node to the target group:
    aws elbv2 register-targets \
    --target-group-arn $ALB_TG_ARN \
    --targets Id=$PRIMARY_NODE_ID
    

Modify the primary node’s security group to allow Spark Connect clients to connect

To connect to Spark Connect, amend only the primary security group. Add an inbound rule to the primary’s node security group to accept Spark Connect TCP connection on port 15002 from your chosen trusted IP address:

aws ec2 authorize-security-group-ingress \
--group-id $PRIMARY_NODE_SG \
--protocol tcp \
--port 15002 \
--source-group $ALB_SG_ID

Connect with a test application

This example demonstrates that a client running a newer Spark version (4.0.1) can successfully connect to an older Spark version on the Amazon EMR cluster (3.5.5), showcasing Spark Connect’s version compatibility feature. This version combination is for demonstration only. Running older versions might pose security risks in production environments.

To test the client-to-server connection, we provide the following test Python application. We recommend that you create and activate a Python virtual environment (venv) before installing the packages. This helps isolate the dependencies for this specific project and prevents conflicts with other Python projects. To install packages, run the following command:

pip install pyspark-client==4.0.1

In your integrated development environment (IDE), copy and paste the following code, replace the placeholder, and invoke it. The code creates a Spark DataFrame containing two rows and it shows its data:

from pyspark.sql import SparkSession
import os
os.environ['GRPC_DEFAULT_SSL_ROOTS_FILE_PATH'] = os.path.expanduser('sc-certificate.cert')
spark = SparkSession.builder \
    .remote("sc://:443/;use_ssl=true") \
    .config('spark.sql.execution.pandas.inferPandasDictAsMap', True) \
    .config('spark.sql.pyspark.legacy.inferMapTypeFromFirstPair.enabled', True) \
    .getOrCreate()
spark.createDataFrame([("sue", 32),("li", 3)],["first_name", "age"]).show()

The following shows the application output:

+----------+---+
|first_name|age|
+----------+---+
|       sue| 32|
|        li|  3|
+----------+---+

Clean up

When you no longer need the cluster, release the following resources to stop incurring charges:

  1. Delete the Application Load Balancer listener, target group, and the load balancer.
  2. Delete the ACM certificate.
  3. Delete the load balancer and Amazon EMR node security groups.
  4. Terminate the EMR cluster.
  5. Empty the Amazon S3 bucket and delete it.
  6. Remove AmazonEMR-ServiceRole-SparkConnectDemo and EMR_EC2_SparkClusterNodesRole roles and EMR_EC2_SparkClusterInstanceProfile instance profile.

Considerations

Security considerations with Spark Connect:

  • Private subnet deployment – Keep EMR clusters in private subnets with no direct internet access, using NAT gateways for outbound connectivity only.
  • Access logging and monitoring – Enable VPC Flow Logs, AWS CloudTrail, and bastion host access logs for audit trails and security monitoring.
  • Security group restrictions – Configure security groups to allow Spark Connect port (15002) access only from bastion host or specific IP ranges.

Conclusion

In this post, we showed how you can adopt modern development workflows and debug Spark applications from local IDEs or notebooks, so you can step through code execution. With Spark Connect’s client-server architecture, the Spark cluster can run on a different version than the client applications, so operations teams can perform infrastructure upgrades and patches independently.

As the cluster operators gain experience, they can customize the bootstrap actions and add steps to process data. Consider exploring Amazon Managed Workflows for Apache Airflow (MWAA) for orchestrating your data pipeline.


About the authors

Philippe Wanner

Philippe Wanner

Philippe is EMEA Tech Lead at AWS. His role is to accelerate the digital transformation for large organizations. His current focus is in a multidisciplinary area involving business transformation, technical strategy, and distributed systems.

Ege Oguzman

Ege Oguzman

Ege is a Software Development Engineer at AWS, and previously he was a Solutions Architect in the public sector. As a builder and cloud enthusiast, he specializes in distributed systems and dedicates his time to infrastructure development and helping organizations build solutions on AWS.

Create and update Apache Iceberg tables with partitions in the AWS Glue Data Catalog using the AWS SDK and AWS CloudFormation

Post Syndicated from Aarthi Srinivasan original https://aws.amazon.com/blogs/big-data/create-and-update-apache-iceberg-tables-with-partitions-in-the-aws-glue-data-catalog-using-the-aws-sdk-and-aws-cloudformation/

In recent years, we’ve witnessed a significant shift in how enterprises manage and analyze their ever-growing data lakes. At the forefront of this transformation is Apache Iceberg, an open table format that’s rapidly gaining traction among large-scale data consumers.

However, as enterprises scale their data lake implementations, managing these Iceberg tables at scale becomes challenging. Data teams often need to manage table schema evolution, its partitioning, and snapshots versions. Automation streamlines these operations, provides consistency, reduces human error, and helps data teams focus on higher-value tasks.

The AWS Glue Data Catalog now supports Iceberg table management using the AWS Glue API, AWS SDKs, and AWS CloudFormation. Previously, users had to create Iceberg tables in the Data Catalog without partitions using CloudFormation or SDKs and later add partitions from Amazon Athena or other analytics engines. This prevents the table lineage from being tracked in one place and adds steps outside automation in the continuous integration and delivery (CI/CD) pipeline for table maintenance operations. With the launch, AWS Glue customers can now use their preferred automation or infrastructure as code (IaC) tools to automate Iceberg table creation with partitions and use the same tools to manage schema updates and sort order.

In this post, we show how to create and update Iceberg tables with partitions in the Data Catalog using the AWS SDK and CloudFormation.

Solution overview

In the following sections, we illustrate the AWS SDK for Python (Boto3) and AWS Command Line Interface (AWS CLI) usage of Data Catalog APIs—CreateTable() and UpdateTable()—for Amazon Simple Storage Service (Amazon S3) based Iceberg tables with partitions. We also provide the CloudFormation templates to create and update an Iceberg table with partitions.

Prerequisites

The Data Catalog API changes are made available in the following versions of the AWS CLI and SDK for Python:

  • AWS CLI version of 2.27.58 or above
  • SDK for Python version of 1.39.12 or above

AWS CLI usage

Let’s create an Iceberg table with one partition, using CreateTable() in the AWS CLI:

aws glue create-table --cli-input-json file://createicebergtable.json

The createicebergtable.json is as follows:

{
    "CatalogId": "123456789012",
    "DatabaseName": "bankdata_icebergdb",
    "Name": "transactiontable1",
    "OpenTableFormatInput": { 
      "IcebergInput": { 
         "MetadataOperation": "CREATE",
         "Version": "2",
         "CreateIcebergTableInput": { 
            "Location": "s3://sampledatabucket/bankdataiceberg/transactiontable1/",
            "Schema": {
                "SchemaId": 0,
                "Type": "struct",
                "Fields": [ 
                    { 
                        "Id": 1,
                        "Name": "transaction_id",
                        "Required": true,
                        "Type": "string"
                    },
                    { 
                        "Id": 2,
                        "Name": "transaction_date",
                        "Required": true,
                        "Type": "date"
                    },
                    { 
                        "Id": 3,
                        "Name": "monthly_balance",
                        "Required": true,
                        "Type": "float"
                    }
                ]
            },
            "PartitionSpec": { 
                "Fields": [ 
                    { 
                        "Name": "by_year",
                        "SourceId": 2,
                        "Transform": "year"
                    }
                ],
                "SpecId": 0
            },
            "WriteOrder": { 
                "Fields": [ 
                    { 
                        "Direction": "asc",
                        "NullOrder": "nulls-last",
                        "SourceId": 1,
                        "Transform": "none"
                    }
                ],
                "OrderId": 1
            }  
        }
      }
   }
}

The preceding AWS CLI command creates the metadata folder for the Iceberg table in Amazon S3, as shown in the following screenshot.

Amazon S3 bucket interface showing metadata folder containing single JSON file dated November 6, 2025

You can populate the table with values as follows and verify the table schema using the Athena console:

SELECT * FROM "bankdata_icebergdb"."transactiontable1" limit 10;
insert into bankdata_icebergdb.transactiontable1 values
    ('AFTERCREATE1234', DATE '2024-08-23', 6789.99),
    ('AFTERCREATE5678', DATE '2023-10-23', 1234.99);
SELECT * FROM "bankdata_icebergdb"."transactiontable1";

The following screenshot shows the results.

Amazon Athena query editor showing SQL queries and results for bankdata_icebergdb database with transaction data

After populating the table with data, you can inspect the S3 prefix of the table, which will now have the data folder.

Amazon S3 bucket interface displaying data folder with two subfolders organized by year: 2023 and 2024

The data folders partitioned according to our table definition and Parquet data files created from our INSERT command are available under each partitioned prefix.

Amazon S3 bucket interface showing by_year=2023 folder containing single Parquet file of 575 bytes

Next, we update the Iceberg table by adding a new partition, using UpdateTable():

aws glue update-table --cli-input-json file://updateicebergtable.json

The updateicebergtable.json is as follows.

{
  "CatalogId": "123456789012",
  "DatabaseName": "bankdata_icebergdb",
  "Name": "transactiontable1",
  "UpdateOpenTableFormatInput": {
    "UpdateIcebergInput": {
      "UpdateIcebergTableInput": {
        "Updates": [
          {
            "Location": "s3://sampledatabucket/bankdataiceberg/transactiontable1/",
            "Schema": {
              "SchemaId": 1,
              "Type": "struct",
              "Fields": [
                {
                  "Id": 1,
                  "Name": "transaction_id",
                  "Required": true,
                  "Type": "string"
                },
                {
                  "Id": 2,
                  "Name": "transaction_date",
                  "Required": true,
                  "Type": "date"
                },
                {
                  "Id": 3,
                  "Name": "monthly_balance",
                  "Required": true,
                  "Type": "float"
                }
              ]
            },
            "PartitionSpec": {
              "Fields": [
                {
                  "Name": "by_year",
                  "SourceId": 2,
                  "Transform": "year"
                },
                {
                  "Name": "by_transactionid",
                  "SourceId": 1,
                  "Transform": "identity"
                }
              ],
              "SpecId": 1
            },
            "SortOrder": {
              "Fields": [
                {
                  "Direction": "asc",
                  "NullOrder": "nulls-last",
                  "SourceId": 1,
                  "Transform": "none"
                }
              ],
              "OrderId": 2
            }
          }
        ]
      }
    }
  }
}

UpdateTable() modifies the table schema by adding a metadata JSON file to the underlying metadata folder of the table in Amazon S3.

Amazon S3 bucket interface showing 5 metadata objects including JSON and Avro files with timestamps

We insert values into the table using Athena as follows:

insert into bankdata_icebergdb.transactiontable1 values
    ('AFTERUPDATE1234', DATE '2025-08-23', 4536.00),
    ('AFTERUPDATE5678', DATE '2022-10-23', 23489.00);
SELECT * FROM "bankdata_icebergdb"."transactiontable1";

The following screenshot shows the results.

Amazon Athena query editor with SQL statements and results after iceberg partition update and insert data

Inspect the corresponding changes to the data folder in the Amazon S3 location of the table.

Amazon S3 prefix showing new partitions for the Iceberg table

This example has illustrated how to create and update Iceberg tables with partitions using AWS CLI commands.

SDK for Python usage

The following Python scripts illustrate using CreateTable() and UpdateTable() for an Iceberg table with partitions:

CloudFormation usage

Use the following CloudFormation templates for CreateTable() and UpdateTable(). After the CreateTable template is complete, update the same stack with the UpdateTable template by creating a new changeset for your stack and executing it.

Clean up

To avoid incurring costs on the Iceberg tables created using the AWS CLI, delete the tables from the Data Catalog.

Conclusion

In this post, we illustrated how to use the AWS CLI to create and update Iceberg tables with partitions in the Data Catalog. We also provided the SDK for Python and CloudFormation sample code and templates. We hope this helps you automate the creation and management of your Iceberg tables with partitions in your CI/CD pipelines and production environments. Try it out for your own use case and share your feedback in the comments section.


About the authors

Acknowledgements: A special thanks to everyone who contributed to the development and launch of this feature – Purvaja Narayanaswamy, Sachet Saurabh, Akhil Yendluri and Mohit Chandak.

Aarthi Srinivasan

Aarthi Srinivasan

Aarthi is a Senior Big Data Architect with AWS. She works with AWS customers and partners to architect data lake house solutions, enhance product features, and establish best practices for data governance.

Pratik Das

Pratik Das

Pratik is a Senior Product Manager with AWS. He is passionate about all things data and works with customers to understand their requirements and build delightful experiences. He has a background in building data-driven solutions and machine learning systems in production.

IPv6 addressing with Amazon Redshift

Post Syndicated from Srini Ponnada original https://aws.amazon.com/blogs/big-data/ipv6-addressing-with-amazon-redshift/

As we witness the gradual transition from IPv4 to IPv6, Amazon Web Services (AWS) continues to expand its support for dual-stack networking across its service portfolio. In this post, we show how you can migrate your Amazon Redshift Serverless workgroup from IPv4-only to dual-stack mode, so you can make your data warehouse future ready.

An IP address serves as a digital identity for devices connected to the internet. This unique numerical identifier enables devices to communicate across IP-based networks, facilitating the exchange of data packets between source and destination.

Today’s internet operates on two IP versions:

  • IPv4 – The traditional 32-bit addressing system (such as 192.168.0.22) that has powered internet communications for over three decades. With approximately 4 billion possible addresses (2³²), IPv4’s limitations have become increasingly apparent as our digital environment expands.
  • IPv6 – The next-generation 128-bit addressing system (such as 2606:4700::6810:787f) offers an astronomical number of unique addresses (340 undecillion or 2¹²⁸). This virtually unlimited address space is designed to accommodate the explosive growth of internet-connected devices.

In the case of Amazon Redshift, dual-stack networking allows Redshift workgroups to communicate over both IPv4 and IPv6 protocols simultaneously. This networking architecture allows Redshift workgroups to be accessible using both IPv4 and IPv6 addresses, providing greater flexibility and future-proofing for network communications. Dual-stack networking provides the following advantages:

  1. Future-proofing – Facilitates compatibility with both IPv4 systems and modern IPv6 networks
  2. Enhanced connectivity – Provides more flexible networking options for diverse client applications

Enable dual-stack networking for Amazon Redshift

An Amazon Redshift workgroup operating in dual-stack mode has both IPv4 and IPv6 addresses associated with the database endpoints. We’ve introduced a new API field called ipAddressTypein the Amazon Redshift API that gives you direct control over your workgroup’s network configuration. You can now specifically choose whether your Amazon Redshift instance operates in IPv4-only mode or dual-stack mode. For complete implementation details, refer to the ipAddressType parameter in the Amazon Redshift API Reference.

Best practice

When implementing dual-stack networking in Amazon Redshift, deploy your workgroups in private subnets with virtual private cloud (VPC) endpoints for optimal security and compatibility. This approach aligns with the current Amazon Redshift support model, which requires dual-stack databases to operate in private mode only. Amazon Redshift doesn’t currently support databases with IPv6-only endpoints or publicly accessible dual-stack instances.

Prerequisites

To implement dual-stack networking in Amazon Redshift, you need to have the following prerequisites:

  • An existing Amazon Redshift serverless workgroups running in IPv4-only mode that you want to convert to dual-stack mode
  • Administrative permissions to modify Amazon Redshift workgroup network configurations
  • VPC with both IPv4 and IPv6 CIDR blocks assigned

Enable IPv6 support in your VPC subnets

Before migrating your Amazon Redshift Serverless workgroup to dual-stack mode, you must first make sure your VPC subnets support IPv6 addressing. In this section, we walk through the process of enabling IPv6 CIDR blocks for your VPC.

Existing VPC Subnets

To enable dual-stack mode in your existing VPC follow these five high-level steps:

  1. Access the VPC dashboard
  2. Navigate to the subnet settings
  3. Add the IPv6 CIDR block
  4. Repeat for the required subnets
  5. Verify IPv6 CIDR association

To access the VPC dashboard:

  1. Sign in to your account on the AWS Management Console
  2. In the search bar at the top, type VPC
  3. Choose VPC from the dropdown list to navigate to the Amazon Virtual Private Cloud (Amazon VPC) dashboard

To navigate to subnet settings:

  1. In the left navigation panel under Virtual private cloud, choose Subnets
  2. From the subnet list, identify and select the subnet(s) that your Amazon Redshift Serverless workgroup uses or will use

To add the IPv6 CIDR block:

  1. With your subnet selected, choose Actions in the dropdown list
  2. Choose Edit IPv6 CIDRs from the available options
  3. In the configuration panel that appears, choose Add IPv6 CIDR
  4. The system will automatically suggest an appropriate IPv6 CIDR block allocation
  5. Choose Save to apply the changes

Repeat for the required subnets. You must modify the subnets within the VPC that will be used by your Amazon Redshift resources. Repeat high-level steps 2–3 for each subnet in your Amazon Redshift subnet group.

Verify the IPv6 CIDR association. After completing the configuration, verify that each subnet displays both IPv4 and IPv6 CIDR blocks in the subnet details. Your subnet details should show something like the following snippet:IPv4 CIDR: 10.0.0.0/24IPv6 CIDR: 2600:1f16:c72:9d00::/64

After you’ve successfully configured IPv6 CIDR blocks for the relevant subnets, you’re ready to proceed with enabling dual-stack mode on your Amazon Redshift Serverless workgroup.

New VPC Subnets

To enable dual-stack mode in a new VPC, follow these steps:

  1. To create a dual-stack VPC, add the --amazon-provided-ipv6-cidr-block option to add an Amazon provided IPv6 CIDR block, as shown in the following example:
    aws ec2 create-vpc --cidr-block 10.0.0.0/24 \
    --amazon-provided-ipv6-cidr-block \
    --query Vpc.VpcId \
    --output text

  2. [Dual stack VPC] Get the IPv6 CIDR block that’s associated with your VPC by using the following describe-vpcs command:
    aws ec2 describe-vpcs --vpc-id vpcxxxxx \
    --query Vpcs[].Ipv6CidrBlockAssociationSet[].Ipv6CidrBlock \
    --output text

  3. If you created a dual-stack VPC, you can use the --ipv6-cidr-block option to create a dual stack subnet, as shown in the following command:
    aws ec2 create-subnet --vpc-id vpc-xxx \
    --cidr-block 10.0.1.0/20 \
    --ipv6-cidr-block 2600:1f13:cfe:3600::/64 \
    --availability-zone us-east-2a \
    --query Subnet.SubnetId \
    --output text

Migrate an existing Amazon Redshift Serverless workgroup from IPv4 to IPv6

To enable dual-stack mode for your Amazon Redshift Serverless workgroup, follow these five high-level steps:

  1. Access Amazon Redshift Serverless
  2. Select your workgroup
  3. Access network and security settings
  4. Enable dual-stack mode
  5. Verify the configuration

To access Amazon Redshift Serverless:

  1. Sign in to AWS Management Console using your credentials.
  2. In the search bar at the top of the console, enter Redshift.
  3. Choose Amazon Redshift in the dropdown list. This will take you to the Amazon Redshift dashboard. Confirm Make sure you’re on the Redshift Serverless dashboard view

To select your workgroup:

  1. In the Redshift Serverless dashboard, locate the Workgroups section
  2. Select the name of the specific workgroup you want to modify

To access network and security settings:

  1. On the workgroup details page, locate the horizontal navigation tabs
  2. Next to the Query and database monitoring section, choose the Data access tab
  3. Choose Edit to access the Edit network and security page

To enable dual-stack mode:

  1. In the network settings section, locate the IP address type options
  2. Choose Dual-stack mode. This enables connectivity for both IPv4 and IPv6
  3. Choose Save changes at the bottom of the page

To verify the configuration:

After the changes are applied, you’ll be returned to the workgroup details page. Confirm that your workgroup now displays Dual-stack mode in its network settings, as shown in the following screenshot.

Your Amazon Redshift Serverless workgroup is now configured to support both IPv4 and IPv6 traffic. This configuration allows your Redshift Serverless workgroup to communicate over both IPv4 and IPv6 protocols, providing greater flexibility for your network connectivity options.

Access Redshift dual-stack serverless workgroups

Redshift dual-stack workgroups maintain the same access methods regardless of whether you’re connecting using IPv4 or IPv6. Your existing connection endpoints remain unchanged.

To access a Redshift dual-stack workgroups from an Amazon Elastic Compute Cloud (Amazon EC2) instance, follow these steps:

  1. Create an IPv4 EC2 instance.
  2. Add the associated EC2 security group to your Redshift workgroup’s security group inbound rules.
  3. Connect to your EC2 instance. Log in to your EC2 instance on the AWS console to install the psql client to test the database connectivity. Enter the following commands from the terminal window:
    # Update your system packages
    sudo dnf update -y
    # Install the PostgreSQL 15 repository
    sudo dnf install -y postgresql15
    # Verify the installation
    psql –version

    Connect to your Redshift workgroup using application user

    psql -h your-redshift-endpoint -U your-username -d your-database -p 5439

    Enter password when prompted

  4. Enter sample queries, as shown in the following:
    SELECT * FROM your_table LIMIT 10;

  5. Validate the IPv4 connection using the following SQL by replacing your associated IPv4 EC2 instance IP address:
    SELECT * FROM sys_connection_log where user_name = 'admin' 
    and remote_host like '%172.31.83.132%' order by record_time desc;

  6. Execute identical validation steps on your IPv6-enabled EC2 instance to verify that all functionality operates correctly with the IPv6 protocol stack using the preceding commands.

Create dual-stack mode in Amazon Redshift Serverless using AWS CLI

You can create a new dual-stack mode in Amazon Redshift Serverless using AWS Command Line Interface (AWS CLI). Follow these high-level steps:

  1. Create a namespace
  2. Create a workgroup
  3. Verify the workgroup is set up in dual-stack mode

To create namespace, enter the following code:

export AWS_USE_DUALSTACK_ENDPOINT=true
aws redshift-serverless create-namespace \
--region us-east-1 \
--namespace-name ipv6-demo \
--admin-username xxx  \
--admin-user-password "yyyyyyy"

To create workgroup, enter the following code:

aws redshift-serverless create-workgroup \
--workgroup-name ipv6-demo-wg \
--namespace-name ipv6-demo \
--region us-east-1 \
--subnet-ids subnet-ppppppp subnet-qqqqqq \
--ip-address-type dualstack

To verify the workgroup is set up in dual-stack mode, refer to the steps in the previous section.

Clean up

To clean up your resources, complete the following steps:

  1. On the Amazon Redshift Serverless console, delete the Amazon Redshift workgroups and namespaces
  2. On the Amazon EC2 console, terminate the EC2 instances

Conclusion

In this post, we’ve explored the capability of Amazon Redshift Serverless to support IPv6 addressing through dual-stack mode, marking a significant advancement in the AWS data warehouse networking flexibility.

We’ve walked through the complete migration journey, from preparing your VPC subnets with IPv6 CIDR blocks to configuring your Amazon Redshift Serverless workgroup for dual-stack operation. The process is straightforward. Although IPv6-only configurations aren’t yet supported for Amazon Redshift, the dual-stack approach provides an ideal transition path, maintaining compatibility with existing IPv4 systems while introducing IPv6 capabilities. Remember that dual-stack configurations are currently limited to private access mode, with public accessibility not yet supported for dual-stack instances.

By migrating to dual-stack mode now, you can make sure your Amazon Redshift environment remains optimally connected, addressable, and ready to support your organization’s data analytics needs well into the future—regardless of how internet addressing protocols continue to evolve.

If you have questions or suggestions on the content covered in this post, leave them in the comments section.


About the authors

Srini Ponnada

Srini Ponnada

Srini is a Sr. Data Architect at AWS. He has helped customers build scalable data warehousing and big data solutions for over 20 years. He loves to design and build efficient end-to-end solutions on AWS.

Ji Yanzhu

Yanzhu Ji

Ji Yanzhu is a Senior Product Manager on the Amazon Redshift team. She has extensive experience in database security and developing product vision and strategy for industry-leading data products and platforms. She excels at building robust software products using web development, system design, database, and distributed programming techniques.

Hua Zirui

Hua Zirui

Zirui Hua is a Software Development Engineer for Amazon Redshift, where he works on developing next generation features for Amazon Redshift. His main focuses are on networking and proxy of database. Outside of work, he likes to play tennis and basketball.

Sandeep Adwankar

Sandeep Adwankar

Sandeep is a Senior Product Manager at AWS. Based in the California Bay Area, he works with customers around the globe to translate business and technical requirements into products that enable customers to improve how they manage, secure, and access data.

Sumanth Punyamurthula

Sumanth Punyamurthula

Sumanth is a Senior Data and Analytics Architect at AWS with more than 20 years of experience in leading large analytical initiatives, including analytics, data warehouse, data lakes, data governance, security, and cloud infrastructure across travel, hospitality, financial, and healthcare industries.

Niranjan Kulkarni

Niranjan Kulkarni

Niranjan is a Software Development Engineer for Amazon Redshift. He focuses on Amazon Redshift Serverless adoption and Amazon Redshift security-related features. Outside of work, he spends time with his family and enjoys watching high-quality TV series.