Tag Archives: AWS Glue

How United Airlines uses Amazon Redshift and AWS Glue Data Catalog federation to query Databricks-managed data

Post Syndicated from Vaibhav Agrawal original https://aws.amazon.com/blogs/big-data/how-united-airlines-uses-amazon-redshift-and-aws-glue-data-catalog-federation-to-query-databricks-managed-data/

This post was co-written with Ankit Aggarwal and Raja Kalluri from United Airlines.

United Airlines processes billions of events daily across its data platform, which spans Amazon Redshift and Databricks with Unity Catalog. To bridge these platforms without duplicating data, the team turned to AWS Glue Data Catalog federation.

In this post, we walk through how to configure AWS Glue Data Catalog federation to connect with Databricks Unity Catalog, so you can run live SQL queries from Amazon Redshift without moving or duplicating data.

Why United Airlines needed catalog federation

United Airlines curates petabytes of data through a medallion architecture (bronze to silver to gold) on Amazon Simple Storage Service (Amazon S3). The airline user interaction data layer alone is several double-digit terabytes of near real-time streamed data. Teams use it to measure customer engagement patterns, feature adoption, and conversion behavior across web and mobile touchpoints. Analysts need to query this curated data through Amazon Redshift Serverless. As part of the existing data platform architecture these data tables are cataloged in Databricks Unity Catalog, not in the AWS Glue Data Catalog. As a result, Amazon Redshift has no native visibility into them. Without catalog federation, the only way to make this data queryable from Amazon Redshift would have been to duplicate it into Amazon Redshift Managed Storage (RMS) and build pipelines to keep it in sync.

AWS Glue Data Catalog federation removed this need. Amazon Redshift users now query the gold layer stored in Amazon S3 directly, with Iceberg metadata resolved from Unity Catalog at query time and no data movement. AWS Glue Data Catalog federation connects Amazon Redshift to external catalogs like Unity Catalog, so analysts query cross-platform data without building sync pipelines or duplicating storage.

Amazon Redshift Serverless is powered by the same Graviton-based query engine used in the new RG instance family, which delivers up to 2x faster data lake query performance compared to prior generations. This engine is purpose-built for reading Apache Iceberg tables directly from Amazon S3, making it well-suited for such federated query workloads.

United Airlines is taking a phased approach to adopting AWS Glue Data Catalog federation across its data platform. The initial focus is the most heavily used user interaction data tables, with 30 tables currently federated in production and 70 more in active rollout. Several hundred additional tables across different business domains are planned for production in the coming months.

Solution overview

AWS Glue Data Catalog federation bridges these platforms at the metadata layer. Here’s how the architecture works.

The architecture follows a four-layer federation chain:

  • Databricks Unity Catalog exposes tables through its Iceberg REST API endpoint. For Delta tables, you can turn on UniForm format to make them Iceberg compatible.
  • AWS Glue Data Catalog creates a federated catalog that connects to Databricks Unity Catalog, making metadata visible within AWS without data movement.
  • A resource link database in the default AWS Glue catalog acts as a bridge, pointing to the federated catalog database. This is required for Amazon Redshift compute.
  • Amazon Redshift Serverless references the resource link database through an external schema. When a query runs, Amazon Redshift traverses the link, calls AWS Glue Federation, and reads the Iceberg data through the Databricks Unity Catalog REST API. AWS Lake Formation governs permissions throughout this chain.

Key services or service features used in this solution:

Figure 1: Federation chain from Databricks Unity Catalog to Amazon Redshift Serverless through AWS Glue and Lake Formation

The architecture follows a six-step flow:

  1. A SQL analyst submits a query to Amazon Redshift Serverless.
  2. Amazon Redshift resolves the external schema through the AWS Glue Data Catalog (resource link to federated catalog).
  3. The AWS Glue federated catalog calls the Databricks Unity Catalog Iceberg REST API to retrieve current table metadata.
  4. The namespace IAM role calls AWS Lake Formation GetDataAccess to obtain scoped, temporary S3 credentials.
  5. Lake Formation evaluates fine-grained access policies and vends credentials for the authorized data files.
  6. Amazon Redshift Serverless reads the Iceberg data files directly from S3 and returns results to the analyst.

Prerequisites

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

  • A Databricks workspace with Unity Catalog enabled and at least one catalog, schema, and table. Databricks uses UniForm to generate Iceberg metadata on Delta Lake tables on Amazon S3.
  • An AWS account with permissions to manage AWS Glue, AWS Lake Formation, Amazon Redshift Serverless, and IAM.
  • An Amazon Redshift Serverless workgroup and namespace already provisioned.
  • AWS Lake Formation set up with a data lake administrator.
  • AWS Command Line Interface (AWS CLI) configured with appropriate credentials.
  • Familiarity with Amazon Redshift Query Editor v2 or a SQL client.

Note: For setting up the Databricks Unity Catalog side (Phase 1), follow the steps in the AWS blog post Access Databricks Unity Catalog data using catalog federation in the AWS Glue Data Catalog. This walkthrough picks up after the federated catalog has been created in AWS Glue.

Solution walkthrough

The walkthrough is organized into six steps covering Lake Formation configuration, the resource link pattern, IAM role setup, and querying Databricks tables from Amazon Redshift.

Step 1: Configure AWS Lake Formation

1a. Add a data lake administrator

  • In Lake Formation, choose Administration, then choose Administrators and add your admin IAM user or role.

1b. Confirm the federated catalog is registered

  • Choose Data Catalog, then Catalogs and verify that databricks-federated-catalog is visible and registered.

This step is the key architectural detail in the walkthrough. Amazon Redshift resolves CREATE EXTERNAL SCHEMA only against the default AWS Glue Data Catalog. The federated catalog (databricks-federated-catalog) is a separate, non-default catalog object. To give Amazon Redshift a path to the federated data, you create a resource link database in the default catalog that points to the federated catalog’s database.

A resource link does not copy data or metadata. It’s a pointer that Lake Formation resolves at query time.

To create the resource link in the Lake Formation console:

  • Choose Data Catalog, Databases, Create database. Then select Resource link.
  • For Resource link name, enter databricks_federated_db_link.
  • For Target catalog, enter databricks-federated-catalog.
  • For Target database, enter the database name that was discovered by the AWS Glue crawler (for example, databricks_federated_db).

Alternatively, use the AWS CLI:

aws glue create-database \
  --database-input '{
    "Name": "databricks_federated_db_link",
    "TargetDatabase": {
      "CatalogId": "<account-id>:databricks-federated-catalog",
      "DatabaseName": "databricks_federated_db"
    }
  }'

Step 3: Configure the Amazon Redshift Serverless namespace IAM role

When Amazon Redshift queries through the resource link, it uses the IAM role attached to the Amazon Redshift Serverless namespace to call the Lake Formation GetDataAccess API. Lake Formation permissions must be granted to this namespace role.

Choose one of these two approaches:

  • Option A – Update your existing namespace role by adding the following policy inline.
  • Option B – Create a new dedicated role (named RedshiftServerlessNamespaceRole) and attach it to the namespace alongside existing roles.

Attach the following IAM policy to the role:

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Effect": "Allow",
      "Action": [
        "glue:GetDatabase",
        "glue:GetDatabases",
        "glue:GetTable",
        "glue:GetTables",
        "glue:GetPartitions",
        "glue:GetCatalog",
        "glue:GetCatalogs"
      ],
      "Resource": "*"
    },
    {
      "Effect": "Allow",
      "Action": "lakeformation:GetDataAccess",
      "Resource": "*"
    }
  ]
}

Note: The Resource: “*” in this policy is shown for simplicity. In production, scope resources to specific AWS Glue catalog ARNs, database ARNs, and table ARNs based on your use case.*

After creating or updating the role, associate it with your Amazon Redshift Serverless namespace:

  • In the Amazon Redshift Serverless console, choose Namespaces, select [your namespace], then choose Security and encryption, then Manage IAM roles.
  • If you use Option A, the existing role already has the new permissions, so no change is needed.
  • If you use Option B, add the new role alongside the existing roles.

Step 4: Grant Lake Formation permissions to the Amazon Redshift namespace role

4a. Grant DESCRIBE on the resource link database (default catalog)

  • In Lake Formation, choose Permissions, Data lake permissions, then Grant.
  • Principal: RedshiftServerlessNamespaceRole.
  • Resources: Named Data Catalog resources, Default catalog, databricks_federated_db_link (resouce link).
  • Database permissions: DESCRIBE.

4b. Grant SELECT and DESCRIBE on the target tables (Grant on Target)

Resource links permit only DESCRIBE and DROP permissions on the link itself. To allow Amazon Redshift to actually read data, you must separately grant SELECT on the target tables in the federated catalog. This is the Lake Formation Grant on Target pattern.

  • Principal: RedshiftServerlessNamespaceRole.
  • Resources: Named Data Catalog resources, databricks-federated-catalog, databricks_federated_db, then Tables.
  • Table permissions: SELECT, DESCRIBE.
  • Catalog permission: DESCRIBE.

Important: SELECT must be granted on the TARGET tables in the federated catalog, not on the resource link. Granting SELECT only on the resource link won’t work. This is a common configuration error.

Step 5: Create an external schema in Amazon Redshift

With the resource link in place and permissions granted, you can now create an external schema in Amazon Redshift that points to the resource link database. The external schema is the query interface. When a user runs SQL against it, Amazon Redshift traverses the link to the federated catalog and retrieves metadata and data from Databricks Unity Catalog.

The DATABASE parameter must reference the resource link database name in the default AWS Glue catalog (databricks_federated_db_link), not the federated catalog name directly. The CATALOG_ARN parameter isn’t required here because the resource link lives in the default catalog and Amazon Redshift resolves it automatically.

Connect to your Amazon Redshift cluster as a superuser (for example, using Amazon Redshift Query Editor v2) and run:

CREATE EXTERNAL SCHEMA databricks_schema
FROM DATA CATALOG
DATABASE 'databricks_federated_db_link'
IAM_ROLE '<iam-role-arn>'
REGION '<region>';

A key design principle in this architecture is the clear separation between data physically stored in Amazon Redshift and data accessed externally through federation. External schemas provide a transparent abstraction layer, so Amazon Redshift users can query data stored in S3 without ingestion. For consistency and clarity, United Airlines follows a standard naming convention for all federated schemas in Amazon Redshift: {domain}_iceberg. This convention makes it immediately clear that the data isn’t natively stored within Amazon Redshift but is accessed by using federation through AWS Glue and Lake Formation. This distinction is critical for analysts and engineers, because it improves discoverability, avoids ambiguity between storage layers, and reinforces architectural discipline when working across hybrid data environments.

The User Interactions domain exposes curated datasets representing customer interaction activity, engagement behavior, and channel usage patterns. Operational datasets follow the same pattern, providing governed access to supporting business events and reference information through a common federation framework.

You create a view layer over each external schema using WITH NO SCHEMA BINDING, so that analysts always resolve the freshest schema on each query execution. For example:

CREATE VIEW analytics.clickstream_events AS
SELECT * FROM {domain}_iceberg.interaction_events
WITH NO SCHEMA BINDING;

Step 6: Verify and query Databricks tables from Amazon Redshift

After creating the external schema, verify that the Databricks tables are visible and run a test query.

Verify table visibility

-- Confirm federated tables are visible in Redshift
SELECT * FROM SVV_EXTERNAL_TABLES
WHERE schemaname = 'databricks_schema';

Query a Databricks Unity Catalog table

-- Query a Databricks Unity Catalog table via the federated catalog
SELECT *
FROM databricks_schema.<table_name>
LIMIT 10;

When a query runs, Amazon Redshift calls Lake Formation GetDataAccess using the namespace IAM role to obtain temporary credentials. It then contacts the AWS Glue federated catalog, which in turn calls the Databricks Unity Catalog Iceberg REST API to retrieve metadata and read table data. The result is returned to the Amazon Redshift user transparently.

For SAML-authenticated users, connect using your IdP JDBC plugin:

jdbc:redshift:iam://<workgroup-name>.<account-id>.<region>.redshift-serverless.amazonaws.com:5439/<database>
?plugin_name=com.amazon.redshift.plugin.<YourIdPPlugin>
&idp_host=<your-idp-host>
&preferred_role=arn:aws:iam::<account-id>:role/RedshiftSAMLUserRole
&ssl=true

The Amazon Redshift JDBC driver handles authentication automatically. It authenticates with your IdP, receives a SAML assertion, and calls sts:AssumeRoleWithSAML for temporary IAM credentials. It then calls redshift-serverless:GetCredentials to connect as the mapped database user.

Business impact

AWS Glue Data Catalog federation delivered measurable architectural and operational improvements for United Airlines:

Area Before After Impact
Data access Delta Lake and Amazon Redshift data were completely siloed, so Amazon Redshift users had no access to curated datasets on Databricks-managed S3 data Amazon Redshift users get real-time access to Databricks-managed data through AWS Glue Data Catalog federation ~100 analysts gained access to user interaction data tables in the first phase without adding new pipelines.
Disaster recovery Cross-Region DR relied on Amazon Redshift snapshots every 3 hours (recovery point objective, or RPO, of 3 hours or more) Amazon S3 cross-Region replication on the Delta Lake provides a near-continuous RPO. A new Amazon Redshift Serverless workgroup in the DR Region can federate to the same S3 data More resilient architecture. Reduces cost for Amazon Redshift snapshot and copy maintenance across Regions
Architecture simplification Data processing happened in both Databricks and Amazon Redshift, requiring manual catalog synchronization between the two platforms which was operationally expensive and prone to drift With the federated architecture, data processing is consolidated in Databricks, and Amazon Redshift acts solely as a query engine powering user queries and dashboards through catalog federation Single processing platform, zero sync pipelines, single source of truth
Infrastructure cost Running dedicated Amazon Redshift ETL cluster with RMS storage, snapshots, and compute for data processing For this use case with federation, Amazon Redshift is not needed for ETL but only as a query engine. No RMS storage duplication, no snapshot replication required ~$30K/month in redundant ETL infrastructure cost reduced

Security considerations

At United Airlines, identity governance is unified through Azure Active Directory groups. On the AWS consumption side, users authenticate to Amazon Redshift Serverless through SAML federation. AD group membership determines database-level access to federated schemas. On the Databricks side, the same AD groups govern access to Unity Catalog schemas. This single-identity model provides consistent access control across both platforms without requiring separate user provisioning. Lake Formation handles credential vending for S3 data access during federated queries, while schema-level access decisions are managed through the AD group mappings on each platform.

The architecture also provides multiple layers of security controls built into the federation chain:

  • AWS Lake Formation governs fine-grained access control throughout the federation chain, so that principals can only access authorized databases, tables, and columns.
  • IAM roles follow least-privilege principles. The Amazon Redshift namespace role is scoped only to AWS Glue metadata operations and Lake Formation GetDataAccess.
  • SAML-based authentication integrates enterprise identity providers, so that users authenticate through existing SSO infrastructure before accessing federated data.
  • All Amazon Redshift connections enforce TLS encryption (ssl=true), protecting data in transit between clients and the Amazon Redshift endpoint.
  • Lake Formation permission vending issues short-lived, scoped credentials for each query execution rather than long-lived static credentials.

Other considerations

Review the catalog federation service limitations before deploying. Key requirements:

  • Delta Lake tables must have UniForm enabled to expose Iceberg-compatible metadata.
  • We recommend that source tables be well-partitioned and regularly compacted, because the federated query performance reflects how efficiently the data is organized at write time.

Clean up

To avoid ongoing charges for resources created in this walkthrough, remove them in the following order. This teardown doesn’t affect Databricks metadata or your underlying data stored in Amazon S3.

  • Drop the external schema in Amazon Redshift: DROP SCHEMA databricks_schema;.
  • Delete the resource link database in the default AWS Glue catalog (databricks_federated_db_link).
  • Revoke Lake Formation permissions granted to the Amazon Redshift namespace role on both the resource link database and the target tables in the federated catalog.
  • Delete the federated catalog in AWS Glue (databricks-federated-catalog).
  • Deregister the AWS Glue connection for the Databricks Unity Catalog if no longer needed.
  • Optionally, remove the IAM role (RedshiftServerlessNamespaceRole) if it was created solely for this walkthrough.

Conclusion

In this post, we showed how United Airlines uses AWS Glue Data Catalog federation to give Amazon Redshift Serverless analysts real-time access to double-digit terabytes of curated user interaction data on Amazon S3, without duplicating a single byte or building sync pipelines.

The architecture uses the Iceberg REST API, resource link databases, and Lake Formation credential vending to create a governed query path between Amazon Redshift and Unity Catalog. For United Airlines, this eliminated redundant ETL infrastructure costs, removed the need for catalog synchronization, and turned Amazon Redshift Serverless into a dedicated high-performance query engine for analysts and dashboards.

For questions or feedback, leave a comment on this post.


About the authors

Vaibhav Agrawal

Vaibhav Agrawal

Vaibhav Agrawal is a Senior Analytics Specialist Solutions Architect at AWS, focused on helping enterprise customers design and implement modern data architectures using AWS Analytics services.

Ankit Aggarwal

Ankit Aggarwal

Ankit Aggarwal is a Principal Enterprise Architect at United Airlines, where he leads the United Data Hub (UDH) platform architecture—a petabyte-scale data platform built on AWS and Databricks. He brings over 15 years of experience in data engineering and enterprise architecture.

Raja Kalluri

Raja Kalluri is a Principal Architect at United Airlines, where he leads enterprise-scale data architecture and modernization initiatives. He specializes in building cloud-native data platforms, enabling real-time analytics and AI, and transforming legacy ecosystems.

Accelerating Spark queries with Iceberg materialized views

Post Syndicated from Yuzhou Sun original https://aws.amazon.com/blogs/big-data/accelerating-spark-queries-with-iceberg-materialized-views/

In this post, you learn how to reduce Apache Spark query execution time with Apache Iceberg materialized views without changing a single SQL query.

Organizations running analytical workloads on their data lakes often hit a common wall: queries that are slow and costly, yet difficult to rewrite by hand. Multi-table joins, heavy aggregations, and window functions over large fact tables all drive up execution times, but the SQL behind them often can’t be changed. It might come from business intelligence (BI) dashboards, packaged independent software vendor (ISV) applications, or legacy reports, where editing the source introduces regression risk that outweighs the performance gain.

Starting with Amazon EMR 7.12.0 and AWS Glue 5.1, you can accelerate these queries without rewriting them. Automatic query rewrite analyzes the logical plan of each incoming query and compares it against a metadata cache of available MVs. When the optimizer finds a materialized view (MV) that satisfies all or part of a query, it rewrites the plan to read from that MV instead of the base tables. Matches can be structural (aggregations and joins) or exact (more complex patterns like window functions). If no MV matches, the original query runs unchanged with no impact on correctness.

If you have previously tried to speed up slow analytical queries, you might have considered one of the following alternatives. Here is how automatic query rewrite compares:

Query modification approach Stored results Refreshes Modification to existing queries
Standard views in AWS Glue No (re-runs each time) n/a Required
Custom ETL pipeline Yes Manual Required
Hand-rolled rewrite Yes Manual Required
Materialized views with automatic rewrite enabled Yes Automatically through AWS Glue Data Catalog on a schedule when configured Not required when supported

In this post, we:

  • Give a high-level overview of how automatic query rewrite works in Apache Spark.
  • Walk through a concrete example, showing how the same query can benefit from MVs at different levels of coverage.
  • Discuss the trade-offs so you can choose the right MV shape for your workload.

Prerequisites

To use automatic query rewrite with Iceberg materialized views, you need:

  • Amazon EMR release 7.12.0 or later, or AWS Glue 5.1 or later.
  • Source tables in Apache Iceberg or Parquet format, registered in the AWS Glue Data Catalog, in the same AWS Region and account as the materialized view. Parquet source tables are supported for automatic query rewrite starting with Amazon EMR 7.14.0 and AWS Glue 8.1.
  • An Amazon Simple Storage Service (Amazon S3) Tables (a capability of Amazon S3) bucket, or an S3 general purpose bucket, for the materialized view data.
  • Permissions for the definer role. You can use AWS Identity and Access Management (IAM) policies or AWS Lake Formation.
  • Automatic query rewrite turned on in your Spark session: --conf spark.sql.optimizer.answerQueriesWithMVs.enabled=true.
  • For Parquet source tables, set spark.sql.materializedView.v1SourceTables.enabled=true and spark.sql.materializedView.v1ETagVersioning.enabled=true.

For more Spark configurations, see Introducing Apache Iceberg materialized views in AWS Glue Data Catalog.

How it works

Here is how MVs and automatic query rewrite work together:

  • You define a SQL query with aggregations, joins, or filters across your supported source tables.
  • AWS Glue Data Catalog stores the precomputed results as an Apache Iceberg table in your Amazon S3 bucket. You can store it in a general purpose S3 bucket or in Amazon S3 Tables. Any Apache Iceberg-compatible query engine can read the materialized view, including Amazon Athena, Amazon EMR, AWS Glue, Amazon Redshift, and Iceberg-compatible third-party query engines. Automatic query rewrite is available on the AWS optimized Spark runtime in Amazon Athena, Amazon EMR, and AWS Glue. Other engines can query the materialized view directly, but they don’t rewrite queries to use it automatically.
  • Automatic refresh keeps the MV current on a schedule that you define, for example SCHEDULE REFRESH EVERY 1 DAY. You set it at creation time or later with ALTER MATERIALIZED VIEW ... ADD SCHEDULE REFRESH. At that scheduled time, the refresh process checks the current Apache Iceberg snapshot ID or Parquet file ETags and refreshes the MV when it detects source-table changes.
  • Automatic query rewrite redirects matching queries to the MV at query optimization time. Automatic query rewrite in Apache Spark uses two matching strategies:
    • Structural rewrite (adapted from Amazon Redshift) handles an MV defined as a single SELECT-FROM-WHERE-GROUP-BY block over INNER joins. The optimizer can roll up an MV’s aggregates to a coarser grain and pull extra query predicates up onto the MV scan.
    • Exact-match rewrite handles MVs defined as other shapes, such as window functions and outer joins, by matching a canonicalized form of the MV body against subtrees of the query plan.

When the optimizer evaluates a query, it consults a metadata cache of MVs from the configured catalogs and chooses the best match. It also checks MV staleness during optimization. It skips stale MVs, so rewrite won’t return stale results. If no MV matches, the original query runs unchanged.

Note that automatic query rewrite is opt-in: set spark.sql.optimizer.answerQueriesWithMVs.enabled=true when creating the Apache Spark session.

Example: One query with three potential MVs

An MV doesn’t need to cover an entire query to help it. Automatic query rewrite in Apache Spark operates on subtrees: when an MV matches a portion of your query plan, the rewriter substitutes that subtree and lets the rest of the query run on the rewrite output unchanged. The same query can therefore be served by many possible MV designs, each making a different trade-off between per-query speedup, storage cost, and reuse across other queries.

To make this concrete, consider a typical analytics query: “Top 100 preferred US customers by total store spending.” It joins fact and dimension tables, applies two selective filters on the customer dimension, aggregates per customer, ranks the result with a window function, and keeps only the top 100:

SELECT c_customer_id, total_revenue, num_transactions, avg_purchase, revenue_rank
FROM (
    SELECT cust.c_customer_id,
        SUM(sales.ss_quantity * sales.ss_sales_price) AS total_revenue,
        COUNT(*) AS num_transactions,
        AVG(sales.ss_quantity * sales.ss_sales_price) AS avg_purchase,
        RANK() OVER (ORDER BY SUM(sales.ss_quantity * sales.ss_sales_price) DESC) AS revenue_rank
    FROM base_catalog.base_db.store_sales sales
    INNER JOIN base_catalog.base_db.customer cust
        ON sales.ss_customer_sk = cust.c_customer_sk
    WHERE cust.c_birth_country = 'UNITED STATES'
        AND cust.c_preferred_cust_flag = 'Y'
    GROUP BY cust.c_customer_id
) ranked
WHERE revenue_rank <= 100
ORDER BY revenue_rank;

Query 1: The original query. Top 100 preferred US customers by total store spending, before any materialized view.

Three MV designs cover progressively more of this query, from a single-table pre-aggregate to the full query body itself:

Tier 1: Pre-aggregate store_sales only, no join, no filter. This tier is a single-table aggregate of store_sales at customer-surrogate-key grain. The query still must join the customer table, apply both filters, re-aggregate at c_customer_id grain, and run the window function.

CREATE MATERIALIZED VIEW mv_catalog.mv_db.customer_tier_1 AS
SELECT ss_customer_sk,
    SUM(ss_quantity * ss_sales_price) AS sum_revenue,
    COUNT(ss_quantity * ss_sales_price) AS count_revenue,
    COUNT(*) AS num
FROM base_catalog.base_db.store_sales
GROUP BY ss_customer_sk;

Tier 1 MV: Single-table pre-aggregate of store_sales by customer surrogate key (no join, no filter).

The following plans compare the original query plan to the rewritten plan:

Window, filter, Sort
+- Aggregate by c_customer_id
:  total_revenue = SUM(ss_quantity * ss_sales_price)
:  num_transactions = COUNT(*)
:  avg_purchase = AVG(ss_quantity * ss_sales_price)
+- Project
   +- Join Inner ON ss_customer_sk = c_customer_sk
      :- BatchScan store_sales <- reads the large store_sales table
      +- Filter c_birth_country='UNITED STATES' AND c_preferred_cust_flag='Y'
         +- BatchScan customer

Plan 1: Original plan. Scans the large store_sales table.

Window, filter, Sort
+- Aggregate by c_customer_id <- rolls up pre-aggregated sums
:  total_revenue = SUM(sum_revenue) <- sum of sum_revenue
:  num_transactions = SUM(num) <- sum of num
:  avg_purchase = SUM(sum_revenue) / SUM(count_revenue) <- sum of sum_revenue / sum of count_revenue
+- Project
   +- Join Inner ON ss_customer_sk = c_customer_sk
      :- BatchScan customer_tier_1 <- reads pre-aggregated MV
      +- Filter c_birth_country='UNITED STATES' AND c_preferred_cust_flag='Y'
         +- BatchScan customer

Plan 2: Rewritten plan (Tier 1). Reads the pre-aggregated customer_tier_1 MV.

Tier 2: Pre-join store_sales x customer, pre-apply one filter (c_preferred_cust_flag = ‘Y’). The middle tier pre-joins both tables and bakes in the preferred-customer filter. The query still must apply the country filter as a residual on the MV scan and run the RANK() window.

CREATE MATERIALIZED VIEW mv_catalog.mv_db.customer_tier_2 AS
SELECT cust.c_customer_id, cust.c_birth_country,
    SUM(sales.ss_quantity * sales.ss_sales_price) AS sum_revenue,
    COUNT(sales.ss_quantity * sales.ss_sales_price) AS count_revenue,
    COUNT(*) AS num
FROM base_catalog.base_db.store_sales sales
INNER JOIN base_catalog.base_db.customer cust
    ON sales.ss_customer_sk = cust.c_customer_sk
WHERE cust.c_preferred_cust_flag = 'Y'
GROUP BY cust.c_customer_id, cust.c_birth_country;

Tier 2 MV: Pre-joins store_sales and customer, with the preferred-customer filter applied.

Rewritten query plan:

Window, filter, Sort
+- Aggregate by c_customer_id <- rolls up pre-aggregated sums
:  total_revenue = SUM(sum_revenue) <- sum of sum_revenue
:  num_transactions = SUM(num) <- sum of num
:  avg_purchase = SUM(sum_revenue) / SUM(count_revenue) <- reads pre-aggregated MV
+- Filter c_birth_country='UNITED STATES' [residual filter on MV scan]
   +- BatchScan customer_tier_2 <- reads pre-aggregated MV

Plan 3: Rewritten plan (Tier 2). Country filter applied as a residual on the MV scan.

Tier 3: Match the entire query, including the window function and top N filter. This is the most specific tier. The MV body is the target query verbatim (minus the top-level ORDER BY, which is meaningless for a stored set). The MV stores the top-ranked rows the query asks for (rank ≤ 100).

CREATE MATERIALIZED VIEW mv_catalog.mv_db.customer_tier_3 AS
SELECT c_customer_id, total_revenue, num_transactions, avg_purchase, revenue_rank
FROM (
    SELECT cust.c_customer_id,
        SUM(sales.ss_quantity * sales.ss_sales_price) AS total_revenue,
        COUNT(*) AS num_transactions,
        AVG(sales.ss_quantity * sales.ss_sales_price) AS avg_purchase,
        RANK() OVER (ORDER BY SUM(sales.ss_quantity * sales.ss_sales_price) DESC) AS revenue_rank
    FROM base_catalog.base_db.store_sales sales
    INNER JOIN base_catalog.base_db.customer cust
        ON sales.ss_customer_sk = cust.c_customer_sk
    WHERE cust.c_birth_country = 'UNITED STATES'
        AND cust.c_preferred_cust_flag = 'Y'
    GROUP BY cust.c_customer_id
) ranked
WHERE revenue_rank <= 100;

Tier 3 MV: Stores the exact ranked output of the query (exact-match path).

This tier exercises the exact-match rewrite path: the rewriter canonicalizes the MV body and matches it against the query’s logical plan.

Rewritten plan:

Sort revenue_rank ASC
+- BatchScan customer_tier_3 <- reads around 100 stored rows

Plan 4: Rewritten plan (Tier 3). Reads around 100 stored rows.

The trade-off

The three tiers trade per-query speedup against reuse and storage. In our testing on TPC-DS 3 TB, we observed the following:

MV design Pre-computed Reuse Per-query speedup MV size
Baseline (no MV) nothing n/a 1x n/a
Tier 1: store_sales agg by customer surrogate key aggregate of all sales per customer broadest: any per-customer aggregation ~5x faster 0.07% of store_sales for TPC-DS 3 TB
Tier 2: store_sales x customer agg, one filter pre-applied join + aggregate, preferred customers only medium: any country filter, preferred customers ~10x faster 0.04% of store_sales for TPC-DS 3 TB
Tier 3: entire query body verbatim (exact-match) exact ranked output of this query narrowest: only this exact query shape 20x+ faster negligible (only 100 rows)

Performance measured on TPC-DS 3 TB. Speedup is the ratio of baseline execution time to MV-accelerated execution time. Results might vary based on data characteristics, cluster size, and query complexity.

In addition, MVs incur additional cost. Each one runs a query against your source tables once and stores the result. The more pre-computation it does (joining more tables, applying more filters), the more time it takes.

The following chart plots per-query speedup and creation time for the three tiers in our testing on TPC-DS 3 TB. Per-query speedup rises steadily, from about 5x at Tier 1 to over 20x at Tier 3. Creation time doesn’t follow the same pattern: it peaks at Tier 2. Tier 2 pre-joins and aggregates all preferred customers across every country, so it materializes the most data work. Tier 3 applies both filters, so it processes far fewer rows and costs less to create.

Chart comparing three materialized view designs. In our testing with TPC-DS 3 TB, we observed per-query speedup rises from about 5x (Tier 1) to over 20x (Tier 3), while creation time peaks at Tier 2, which materializes the most data work. Stacked bars show creation time split into catalog setup, data work, and commit.

Figure 1: Per-query speedup and creation time across the three materialized view tiers, measured on TPC-DS 3 TB

Start by identifying one expensive query that runs repeatedly with stable filters. It is likely a good candidate for an exact-match MV.

Validating automatic query rewrite

To confirm that your query benefited from automatic rewrite:

  1. Query plan inspection: Check the query’s optimized logical plan or physical plan for a leaf scan node referencing the MV (for example, BatchScan mv_catalog.mv_db.your_mv_name). If the MV appears as a scan source, rewrite succeeded.
  2. Log confirmation (Amazon EMR 7.14.0+): Look for INFO-level log entries such as AQMV outcome: rewritten=true, mvs=[mv_name], duration=12ms.
  3. No-rewrite diagnostics (Amazon EMR 7.14.0+): If rewrite didn’t occur, check the MVRewriteMetricsEvent in the Apache Spark Event Log for the specific reason the optimizer skipped the MV.

If you have set spark.sql.optimizer.answerQueriesWithMVs.enabled=true but your query still runs against the base tables, check the following common causes:

  1. Write commands block rewrite by default. INSERT and MERGE statements don’t trigger rewrite. Set spark.sql.optimizer.answerQueriesWithMVs.commandBlockingEnabled=false to turn on rewrite within write command subqueries.
  2. The MV is stale. Rewrite skips the MV when one or more source tables have changed since its last refresh. Wait for the next scheduled refresh, or force an immediate refresh with REFRESH MATERIALIZED VIEW <mv_name>.
  3. Heuristic candidate filtering. The optimizer uses heuristic checks to narrow the set of MV candidates before attempting a full match. In some cases, an MV that could benefit the query might be filtered out early by these heuristics.
  4. Spark version mismatch (Amazon EMR 7.13.0+). Automatic query rewrite skips MVs whose stored IMV_sparkVersion does not match the cluster’s current Apache Spark version. To bypass this check, set spark.sql.materializedView.sparkVersionCompatibilityCheck.enabled=false.
  5. MV metadata cache not loaded. The metadata cache loads lazily during optimization of the first rewritable query in a Spark session. If your critical query fires before the cache is warm, the MV will not be available. Run a small warm-up query (for example, SELECT 1 FROM <some_iceberg_table>) at session start to pay this cost off the critical path.
  6. MV metadata cache memory limit reached. If the cache was disabled or stopped loading MVs because of reaching its memory limit, increase spark.driver.memory.
  7. Too many tables in configured catalogs. If there are many tables or MVs in the configured catalogs, the cache might not finish loading before your query starts. Place MVs in a dedicated catalog, add it to spark.sql.materializedViews.additionalCatalogs, and set spark.sql.materializedViews.scanCurrentCatalog=false to skip scanning the current catalog.
  8. Parquet base tables have additional limitations and configuration requirements. For automatic query rewrite with Parquet base tables, set spark.sql.materializedView.v1SourceTables.enabled=true and spark.sql.materializedView.v1ETagVersioning.enabled=true. Without ETag versioning, Spark can’t determine a usable source-table version and skips the MV. Partitioned Parquet base tables are also subject to additional validation limits.

Performance considerations

Turning on automatic query rewrite has overhead: it introduces trade-offs that might affect some queries negatively:

  1. Optimization overhead. Enabling rewrite adds processing time during query optimization as the optimizer evaluates MV candidates against the query plan. This overhead applies to every query in the session, including those that ultimately don’t match any MV.
  2. Reduced task parallelism. Reading from an MV instead of the original base table might produce fewer tasks or introduce data skew, depending on the MV’s data layout. This reduces parallelism compared to a direct scan of the larger, more evenly distributed source table.

Conclusion

In this post, we showed how automatic query rewrite can accelerate your existing Apache Spark workloads. It uses Apache Iceberg materialized views in the AWS Glue Data Catalog, without changing a single line of SQL. By storing precomputed results as managed Apache Iceberg tables, the AWS Glue Data Catalog lets the Apache Spark optimizer transparently substitute matching query plans. You get the performance benefit of pre-aggregation without the application-level rewiring. BI dashboards, ISV-generated reports, and legacy pipelines all benefit the moment a matching MV exists.

We walked through three MV designs for the same analytical query, each striking a different balance between per-query speedup, storage footprint, and reuse across your workload. As the trade-off table shows, our testing found that a narrow, exact-match MV delivered 20x+ acceleration for a single query shape. A broader pre-aggregate served an entire family of queries at a more modest ~5x gain. The right choice depends on how many queries share the same join-and-aggregate pattern and how frequently your source data changes.

To get started:

  1. Launch an Amazon EMR 7.12.0+ cluster or an AWS Glue 5.1+ job.
  2. Create an MV over your most expensive repeating query using CREATE MATERIALIZED VIEW in the AWS Glue Data Catalog.
  3. Turn on automatic query rewrite by setting spark.sql.optimizer.answerQueriesWithMVs.enabled=true in your Spark session configuration.
  4. Verify the rewrite by inspecting the optimized query plan for an MV scan node, or by checking INFO-level logs on Amazon EMR 7.14.0+.

Queries with multi-table joins, heavy aggregations, or window functions over large fact tables are strong initial candidates. Start with one high-cost, frequently executed query. Validate the speedup, then expand to broader MVs as you identify shared patterns across your workload.

Special thanks to everyone who contributed to the automatic query rewrite feature and this blog: Andre Hernich, Leon Lin, Yiyang Chen, Geeta Krishna Panda, Ashok Chintalapati, Muhammad Malik, Rishabh Bhatia, and Giovanni Fumarola.

References

For more detail, see the following resources:


About the authors

Yuzhou Sun

Yuzhou Sun

Yuzhou is a software development engineer for Open Data Analytics Engines at Amazon Web Services.

Srishti Mittal

Srishti Mittal

Srishti is a product manager for Open Data Analytics Engines at Amazon Web Services.

Kinshuk Pahare

Kinshuk Pahare

Kinshuk serves as Head of Product for Analytics Engines at AWS, where he leads the product teams responsible for Amazon Redshift, AWS Glue, Amazon EMR, and Amazon Athena. With over six years at AWS, he brings deep expertise in building and scaling cloud-native analytics platforms that help organizations unlock the value of their data at any scale.

Henry Laih

Henry Laih

Henry is a software development engineer for Open Data Analytics Engines at Amazon Web Services.

Srikanth Kandula

Srikanth Kandula

Srikanth is an engineer who works in analytics and distributed systems at Amazon Web Services.

Shahryar Baki

Shahryar Baki

Shahryar is a software development engineer for Open Data Analytics Engines at Amazon Web Services.

Build declarative ETL pipelines with AWS Glue 6.0

Post Syndicated from Syed Humair original https://aws.amazon.com/blogs/big-data/build-declarative-etl-pipelines-with-aws-glue-6-0/

Data teams commonly build the extract, transform, and load (ETL) pipelines that turn raw order events into analyst-ready aggregates as a bronze, silver, and gold sequence, the medallion architecture. Bronze holds raw ingested records, silver holds cleaned and validated data, and gold holds the business-level aggregates that analysts query. Today you build this on AWS Glue with an orchestrator such as Amazon Managed Workflows for Apache Airflow (Amazon MWAA) or AWS Step Functions coordinating the stages. Many teams run production pipelines exactly this way. As a pipeline grows, the coordination work grows with it: you wire job dependencies, manage intermediate checkpoints, and add retry logic stage by stage.

AWS Glue 6.0, powered by Apache Spark 4.1, introduces Spark Declarative Pipelines (SDP), which simplifies this further. Instead of orchestrating jobs by hand, you declare what each dataset should contain and let the declarative framework resolve dependencies, manage checkpoints, and orchestrate execution order automatically. The result runs as a single declarative job, with no manual directed acyclic graph (DAG) wiring or imperative orchestration code.

In this post, you build a single AWS Glue 6.0 job that turns raw order records into validated, aggregated, analytics-ready tables through the bronze, silver, and gold sequence. You do this without writing any orchestration logic. This walkthrough uses the AWS Command Line Interface (AWS CLI), and the same operations are available through the AWS SDKs.

Solution overview

You build a single AWS Glue 6.0 job that reads raw order records from a CSV file in Amazon Simple Storage Service (Amazon S3). The job flows them through three declared datasets. These are a bronze materialized view (ingest as-is), a silver materialized view (type, validate, and classify), and a gold SQL materialized view (aggregate by region). With AWS Glue Data Catalog integration turned on, all three land as Data Catalog tables, queryable with standard SQL tooling such as Amazon Athena. SDP resolves the dependency order from the dataset references in your code, so you never orchestrate the steps yourself.

Two ways to build the pipeline

Before you build the pipeline, let’s understand this new way of writing ETL pipelines with a quick comparison of the imperative and declarative approaches.

With the imperative approach, you need three AWS Glue jobs, plus an orchestrator to handle sequencing and error handling. A typical pipeline therefore has two layers: an orchestration layer and the ETL processing layer. The following diagram shows this two-layer imperative pipeline.

Two-layer imperative pipeline: three AWS Glue jobs coordinated by an orchestrator.

Figure 1: The two-layer imperative pipeline, with three AWS Glue jobs coordinated by an orchestrator.

Compared to that, the declarative approach runs as a single ETL job with SDP. The following diagram mirrors the previous one, but here it is a single AWS Glue ETL job instead of three jobs plus an orchestrator.

Declarative pipeline: a single AWS Glue job running the bronze, silver, and gold layers with Spark Declarative Pipelines.

Figure 2: The declarative pipeline, a single AWS Glue job running the bronze, silver, and gold layers with SDP.

The declarative approach reduces more than the number of jobs. It removes the boilerplate that surrounds them. An orchestrator such as Amazon MWAA or AWS Step Functions already handles retries and parallelism, but only at the granularity of a whole job. To get finer control, teams often split a pipeline into several jobs and then hand-wire the dependencies between them. With SDP, you no longer hand-wire a DAG, manage per-stage checkpoints, or split the pipeline into separate jobs for retries and parallelism. SDP derives the dependency graph from your table references and coordinates execution at the level of individual tables. You can still invoke an SDP job from an orchestrator when a broader workflow calls for it, but the pipeline’s internal coordination is no longer code you write and maintain.

SDP separates the what from the how: you declare datasets (the outputs you want), and SDP builds the flows that produce them and runs them as one pipeline, resolving dependencies and execution order automatically.

You declare these abstractions through Python decorators. This post covers three of them, @dp.table, @dp.materialized_view, and @dp.temporary_view, each with its own purpose:

  • @dp.table defines a streaming table, which processes new data incrementally on each run. Typical use cases are raw event ingestion and change data capture (CDC) feeds.
  • @dp.materialized_view defines a materialized view for batch use cases. Today, this dataset type fully recomputes on each run. Common uses include parsing, aggregations, and machine learning (ML) feature engineering.
  • @dp.temporary_view is for temporary computations and aggregations. It’s pipeline-scoped and isn’t persisted outside the pipeline. Use it for enrichment lookups and subqueries.

Streaming tables append only new arrivals. Materialized views fully recompute. This post uses @dp.materialized_view for all three layers to keep the walkthrough focused. In production, you would typically use @dp.table for the bronze layer to process only new files as they arrive rather than re-reading the full source each run.

Running and refreshing the pipeline

When you rerun a pipeline, you don’t always want the same work to happen. Sometimes you only want to confirm the pipeline is well-formed before spending compute. Other times you want to run it but recompute only the datasets that changed rather than the entire graph. SDP handles both cases through two independent controls, and it helps to keep them separate:

  • Execution mode (the spark.glue.sdp.jobMode key) answers run or only validate?
  • Refresh scope (the spark.glue.sdp.runMode key) answers given that I’m running, what do I recompute?

Execution mode. VALIDATE runs the pipeline in dry-run mode: SDP checks the YAML syntax, dependency resolution, and SQL and Python compilation without writing any data. Use it to verify your pipeline is well-formed before committing compute. RUN (the default) executes the pipeline normally, resolving the dependency graph and materializing datasets.

# Dry run: validate the graph, write nothing
aws glue start-job-run \
  --job-name "${JOB_NAME}" \
  --arguments '{"--conf":"spark.glue.sdp.jobMode=VALIDATE"}' \
  --region "${AWS_REGION}"

# Normal execution
aws glue start-job-run \
  --job-name "${JOB_NAME}" \
  --arguments '{"--conf":"spark.glue.sdp.jobMode=RUN"}' \
  --region "${AWS_REGION}"

Refresh scope. By default, a RUN recomputes every materialized view. You can narrow or widen that with spark.glue.sdp.runMode:

  • --refresh <datasets> updates only the named datasets (comma-separated, no spaces).
  • --full-refresh <datasets> resets and recomputes only the named datasets (for streaming tables, this also clears their checkpoints).
  • --full-refresh-all resets and recomputes every dataset in the pipeline.
# Selective refresh of named datasets
aws glue start-job-run \
  --job-name "${JOB_NAME}" \
  --arguments '{"--conf":"spark.glue.sdp.jobMode=RUN --conf spark.glue.sdp.runMode=--refresh silver_orders,gold_sales_summary"}' \
  --region "${AWS_REGION}"

# Full reset and recompute of the entire pipeline
aws glue start-job-run \
  --job-name "${JOB_NAME}" \
  --arguments '{"--conf":"spark.glue.sdp.jobMode=RUN --conf spark.glue.sdp.runMode=--full-refresh-all"}' \
  --region "${AWS_REGION}"

Selective refresh is useful during development, so you can iterate on a single layer without reprocessing the entire graph. Note that --refresh and --full-refresh each take an explicit list of datasets. To reset the whole pipeline, use --full-refresh-all. Because materialized views hold no incremental state, resetting a materialized view and refreshing it both fully recompute it. The reset-versus-refresh distinction matters for streaming tables, where a refresh processes only new data and a reset clears the checkpoint and reprocesses from scratch.

The multiple values are passed as a single --conf argument string ("spark.glue.sdp.jobMode=RUN --conf spark.glue.sdp.runMode=..."). This is the serialization the AWS Glue SDP mode expects for the run.

Materialized views: Batch transforms with automatic dependency resolution

Materialized views recompute their full result set on each run. SDP infers dependencies from table references: in this pipeline, silver_orders references bronze_orders, so SDP runs bronze first, as shown in the following diagram.

Dependency graph showing Spark Declarative Pipelines running the bronze layer before the silver layer.

Figure 3: SDP infers the dependency order from table references and runs bronze before silver.

The core pattern is a decorated function that returns a DataFrame:

@dp.materialized_view(comment="Raw orders loaded from CSV")
def bronze_orders() -> DataFrame:
    return spark.read.schema(ORDERS_SCHEMA).option("header", "true").csv(ORDERS_PATH)

The silver layer references bronze_orders through spark.table("bronze_orders"), with no explicit dependency declaration. SDP builds the DAG by analyzing table references in your code and runs bronze first automatically.

Bronze reads every column as a string by design: the bronze layer preserves raw source data without coercion. Type casting, validation, and filtering happen in the silver layer.

SQL and Python coexistence

SDP supports both Python and SQL definitions in the same pipeline project. A SQL materialized view can reference a Python-defined table directly, for example the gold layer aggregating the silver table:

CREATE MATERIALIZED VIEW gold_sales_summary
COMMENT 'Completed-order metrics by region'
AS
SELECT
  region,
  COUNT(*) AS order_count,
  CAST(ROUND(SUM(amount), 2) AS DECIMAL(10, 2)) AS total_sales,
  CAST(ROUND(AVG(amount), 2) AS DECIMAL(10, 2)) AS average_order_value
FROM silver_orders
GROUP BY region;

In this post, Python files define ingestion and validation logic, and SQL files define reporting views and aggregations. SDP discovers both through the libraries glob pattern in the pipeline specification and resolves the cross-language dependencies automatically. The complete source for all three layers follows in the step-by-step walkthrough.

Build the pipeline: Step by step

The rest of this post is a hands-on walkthrough. You build a single AWS Glue 6.0 job that reads orders.csv and processes it through the bronze, silver, and gold layers. The steps are:

  1. Prerequisites: AWS account, AWS Identity and Access Management (IAM) role, and S3 bucket.
  2. Set up sample data: create orders.csv and upload it to Amazon S3.
  3. Build the pipeline files (the spark-pipeline.yml specification plus the three transformation files).
  4. Package the pipeline into a zip and upload it to Amazon S3.
  5. Create the database: a Data Catalog database with an S3 location.
  6. Configure the job: create the AWS Glue 6.0 job with the SDP flag.
  7. Validate: run in dry-run mode to verify the graph.
  8. Run the pipeline to materialize all datasets.
  9. Query results: inspect the tables with Amazon Athena.
  10. Clean up: delete the resources you created.

Step 1 – Prerequisites

To follow along, you need:

  • An AWS account with access to AWS Glue 6.0.
  • A dedicated IAM role trusted by glue.amazonaws.com (set up in the following section).
  • A private, encrypted Amazon S3 bucket with Block Public Access enabled.
  • The AWS CLI configured with credentials for a non-production account.

IAM role for the pipeline

Create a role that AWS Glue can assume, with the following trust policy:

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

Attach the AWS managed policy AWSGlueServiceRole, which grants the AWS Glue Data Catalog and Amazon CloudWatch Logs access the job needs. Then add an inline policy that scopes Amazon S3 access to your bucket, covering the input data, the pipeline zip, the pipeline storage (state) path, and the warehouse location:

{
  "Version": "2012-10-17",
  "Statement": [{
    "Effect": "Allow",
    "Action": ["s3:GetObject", "s3:PutObject", "s3:DeleteObject", "s3:ListBucket"],
    "Resource": [
      "arn:aws:s3:::amzn-s3-demo-bucket",
      "arn:aws:s3:::amzn-s3-demo-bucket/*"
    ]
  }]
}

For a full breakdown of the baseline permissions, see Setting up IAM permissions for AWS Glue.

Set the walkthrough variables

Set the following variables, replacing the example values (us-east-1, amzn-s3-demo-bucket, the account ID 111122223333, and the role name) with your own:

export AWS_REGION="us-east-1"
export BUCKET="amzn-s3-demo-bucket"
export PREFIX="simple-sdp-demo"
export DATABASE="simple_sdp_demo_db"
export ROLE_ARN="arn:aws:iam::111122223333:role/AWSGlueServiceRole-sdp-demo"
export JOB_NAME="simple-sdp-demo"

Step 2 – Set up sample data

The pipeline reads a CSV of order records. Save the following as orders.csv:

order_id,customer_id,region,amount,status,order_ts
O-1001,C-101,EMEA,120.50,COMPLETE,2026-07-23T08:00:00Z
O-1002,C-102,AMER,750.00,COMPLETE,2026-07-23T08:15:00Z
O-1003,C-103,EMEA,-10.00,INVALID,2026-07-23T08:30:00Z
O-1004,C-104,APAC,320.25,COMPLETE,2026-07-23T09:00:00Z
O-1005,C-105,AMER,250.00,COMPLETE,2026-07-23T09:15:00Z
O-1006,C-106,EMEA,90.00,COMPLETE,2026-07-23T09:30:00Z

Upload the file to the input/ location under your project prefix, which is where the bronze layer reads it (the ORDERS_PATH in 01_bronze.py, shown in Step 3). Use the variables you exported in Step 1:

aws s3 cp orders.csv \
  "s3://${BUCKET}/${PREFIX}/input/orders.csv" \
  --region "${AWS_REGION}"

The file includes one invalid order (O-1003, a negative amount), which the silver layer filters out to demonstrate the validation step. The AMER and EMEA regions each have two completed orders, so the gold layer’s order_count and average_order_value are meaningful aggregations rather than single-row passthroughs.

Step 3 – Build the pipeline files

The pipeline project uses the structure introduced earlier: a transformations/ folder holding the three layer definitions (01_bronze.py, 02_silver.py, 03_gold.sql), plus the spark-pipeline.yml specification. The following screenshot shows this layout in a code editor.

Pipeline project layout in a code editor, showing the transformations folder and the spark-pipeline.yml file.

Figure 4: The pipeline project layout in a code editor.

The complete contents of each file follow.

3a. spark-pipeline.yml

The specification names the pipeline, points to the Data Catalog database, configures state storage, and discovers transformation files. As with the transformation files, it uses the __DATABASE__, __BUCKET__, and __PREFIX__ tokens, which you substitute at packaging time in Step 4:

name: simple_sdp_demo
catalog: spark_catalog
database: __DATABASE__
storage: s3://__BUCKET__/__PREFIX__/state/
libraries:
  - glob:
      include: transformations/**
configuration:
  spark.sql.shuffle.partitions: "4"

3b. transformations/01_bronze.py

Bronze preserves the raw source as strings. No coercion, no filtering:

"""Bronze layer: preserve source order records as strings."""
from pyspark import pipelines as dp
from pyspark.sql import DataFrame, SparkSession
from pyspark.sql.types import StringType, StructField, StructType

spark = SparkSession.active()

ORDERS_PATH = "s3://__BUCKET__/__PREFIX__/input/orders.csv"

ORDERS_SCHEMA = StructType([
    StructField("order_id", StringType(), True),
    StructField("customer_id", StringType(), True),
    StructField("region", StringType(), True),
    StructField("amount", StringType(), True),
    StructField("status", StringType(), True),
    StructField("order_ts", StringType(), True),
])


@dp.materialized_view(comment="Raw orders loaded from CSV")
def bronze_orders() -> DataFrame:
    return (
        spark.read
        .schema(ORDERS_SCHEMA)
        .option("header", "true")
        .csv(ORDERS_PATH)
    )

The path uses the tokens __BUCKET__ and __PREFIX__ rather than hardcoded values. AWS Glue reads these files from the packaged zip at runtime, so shell variables like ${BUCKET} are not expanded inside them. You substitute the tokens with your real values when you package the project in Step 4, which keeps every file consistent with the variables you exported in Step 1.

3c. transformations/02_silver.py

Silver casts types, filters to complete orders with positive amounts, and derives an amount_band classification:

"""Silver layer: type, validate, and classify complete orders."""
from pyspark import pipelines as dp
from pyspark.sql import DataFrame, SparkSession
from pyspark.sql.functions import col, to_timestamp, trim, when

spark = SparkSession.active()


@dp.materialized_view(comment="Validated complete orders with typed values")
def silver_orders() -> DataFrame:
    typed = (
        spark.table("bronze_orders")
        .select(
            trim(col("order_id")).alias("order_id"),
            trim(col("customer_id")).alias("customer_id"),
            trim(col("region")).alias("region"),
            col("amount").cast("double").alias("amount"),
            trim(col("status")).alias("status"),
            to_timestamp("order_ts", "yyyy-MM-dd'T'HH:mm:ss'Z'").alias("order_ts"),
        )
        .filter(
            col("order_id").isNotNull()
            & col("region").isNotNull()
            & col("order_ts").isNotNull()
            & (col("status") == "COMPLETE")
            & (col("amount") > 0)
        )
    )
    return typed.select(
        "*",
        when(col("amount") >= 500, "large")
        .when(col("amount") >= 100, "medium")
        .otherwise("small")
        .alias("amount_band"),
    )

Silver reads bronze with spark.table("bronze_orders"), so SDP infers the dependency and runs bronze first. Two details matter here:

  • The to_timestamp call passes an explicit format, "yyyy-MM-dd'T'HH:mm:ss'Z'". The source timestamps are ISO 8601 with a Z suffix. Giving the format treats Z as a literal and produces the same wall-clock value regardless of the job’s session time zone, which keeps the result deterministic.
  • The transformation runs in two projections: the first casts and filters, and the second derives amount_band from the already-typed amount column. Deriving columns with .select(...) rather than a separate .withColumn(...) step keeps SDP’s reference to bronze_orders resolvable as a pipeline dependency. This way, SDP consistently orders the bronze layer before the silver layer. The order matters here too. Spark 4.1 enables ANSI mode by default, so comparing the raw string amount against a number would fail. amount_band therefore reads the already-cast amount.

3d. transformations/03_gold.sql

The gold layer aggregates order metrics by region using SQL:

CREATE MATERIALIZED VIEW gold_sales_summary
COMMENT 'Completed-order metrics by region'
AS
SELECT
  region,
  COUNT(*) AS order_count,
  CAST(ROUND(SUM(amount), 2) AS DECIMAL(10, 2)) AS total_sales,
  CAST(ROUND(AVG(amount), 2) AS DECIMAL(10, 2)) AS average_order_value
FROM silver_orders
GROUP BY region;

Step 4 – Package the project

Substitute the __BUCKET__, __PREFIX__, and __DATABASE__ tokens with the values you exported in Step 1. Then package spark-pipeline.yml and the transformations/ folder into a zip with both at the zip root. Because AWS Glue reads these files from the zip at runtime, the substitution has to happen now, at packaging time, not through shell variables at run time:

# Render the tokens into a build/ copy, leaving your source files untouched
rm -rf build/package && mkdir -p build/package/transformations

sed -e "s|__BUCKET__|${BUCKET}|g" \
    -e "s|__PREFIX__|${PREFIX}|g" \
    -e "s|__DATABASE__|${DATABASE}|g" \
    spark-pipeline.yml > build/package/spark-pipeline.yml

sed -e "s|__BUCKET__|${BUCKET}|g" \
    -e "s|__PREFIX__|${PREFIX}|g" \
    transformations/01_bronze.py > build/package/transformations/01_bronze.py
cp transformations/02_silver.py transformations/03_gold.sql build/package/transformations/

# Zip with the spec and transformations at the zip root
(cd build/package && zip -r -q ../simple-sdp-demo.zip spark-pipeline.yml transformations)

# Upload
aws s3 cp build/simple-sdp-demo.zip "s3://${BUCKET}/${PREFIX}/pipeline/simple-sdp-demo.zip" --region "${AWS_REGION}"

Only spark-pipeline.yml and 01_bronze.py carry tokens, so the other files are copied as-is. The uploaded object is named simple-sdp-demo.zip, which is the same name the job references in Step 6.

Step 5 – Create the database

The database named in spark-pipeline.yml must already exist in the AWS Glue Data Catalog, with an S3 location URI, before the pipeline runs. SDP does not create it automatically:

aws glue get-database --name "${DATABASE}" --region "${AWS_REGION}" >/dev/null 2>&1 \
|| aws glue create-database \
--database-input "{\"Name\":\"${DATABASE}\",\"LocationUri\":\"s3://${BUCKET}/${PREFIX}/warehouse/\"}" \
--region "${AWS_REGION}"

Step 6 – Configure the job

Create an AWS Glue 6.0 job with the zip as ScriptLocation and the SDP flag enabled:

aws glue create-job \
--name "${JOB_NAME}" \
--role "${ROLE_ARN}" \
--command "{\"Name\":\"glueetl\",\"ScriptLocation\":\"s3://${BUCKET}/${PREFIX}/pipeline/simple-sdp-demo.zip\",\"PythonVersion\":\"3\"}" \
--glue-version "6.0" \
--worker-type "G.1X" \
--number-of-workers 2 \
--default-arguments "{\"--enable-spark-declarative-pipeline\":\"true\",\"--enable-glue-datacatalog\":\"true\"}" \
--region "${AWS_REGION}"

Key arguments:

Argument Purpose
--enable-spark-declarative-pipeline Activates the SDP executor (required)
--enable-glue-datacatalog Uses the AWS Glue Data Catalog as the Spark Hive metastore, so the pipeline’s output tables register in the catalog
ScriptLocation Points to the pipeline zip, not a .py file

Table 2: Key arguments for the create-job command.

The create-job command sets ScriptLocation to the pipeline zip. You can also point it to an Amazon S3 prefix: upload the unzipped spark-pipeline.yml and transformations/ to a prefix and set ScriptLocation to that prefix (with a trailing /). No other change is needed, and the --enable-spark-declarative-pipeline flag stays the same. The zip keeps the upload to a single object.

Step 7 – Validate (dry run)

Run the job in validation mode first to verify the dependency graph without materializing data:

aws glue start-job-run \
  --job-name "${JOB_NAME}" \
  --arguments '{"--conf":"spark.glue.sdp.jobMode=VALIDATE"}' \
  --region "${AWS_REGION}"

Validation analyzes the project structure, dependency graph, and SQL and Python compilation without creating tables, executing transforms, or writing data. Confirm that the database has no tables after validation completes.

On AWS Glue, validation runs as a job (jobMode=VALIDATE), so you create the job in Step 6 and then validate it here. If you develop locally with the open source spark-pipelines CLI, you can run its dry-run against the project before packaging and uploading.

Step 8 – Run the pipeline

Start the pipeline in normal execution mode:

aws glue start-job-run \
  --job-name "${JOB_NAME}" \
  --arguments '{"--conf":"spark.glue.sdp.jobMode=RUN"}' \
  --region "${AWS_REGION}"

After the run completes, list the materialized tables:

aws glue get-tables \
  --database-name "${DATABASE}" \
  --region "${AWS_REGION}" \
  --query 'TableList[].Name' \
  --output table

Expected tables: bronze_orders, silver_orders, gold_sales_summary.

After the run, the AWS Glue console shows the three output tables in the simple_sdp_demo_db database. The database’s Location is the warehouse path you configured, s3://amzn-s3-demo-bucket/simple-sdp-demo/warehouse/, and each table stores its data under that prefix. The following screenshot shows the database properties and the three tables (bronze_orders, silver_orders, and gold_sales_summary), each registered in the AWS Glue Data Catalog.

The bronze_orders, silver_orders, and gold_sales_summary tables in the AWS Glue Data Catalog.

Figure 5: The three output tables in the AWS Glue Data Catalog.

Step 9 – Query results

Query the tables with Amazon Athena. If this is your first time using Athena in this Region, set an Amazon S3 query-results location for your workgroup first (Athena console, Settings). Also make sure your identity can read the simple_sdp_demo_db tables in the Data Catalog and the underlying S3 data.

-- Bronze preserves all 6 source rows
SELECT * FROM simple_sdp_demo_db.bronze_orders ORDER BY order_id;

-- Silver retains the 5 complete orders with positive amounts
SELECT * FROM simple_sdp_demo_db.silver_orders ORDER BY order_id;

-- Gold aggregates by region
SELECT * FROM simple_sdp_demo_db.gold_sales_summary ORDER BY region;

Expected gold result:

region order_count total_sales average_order_value
AMER 2 1000.00 500.00
APAC 1 320.25 320.25
EMEA 2 210.50 105.25

Table 3: Gold layer aggregation results by region.

Running the query in the Amazon Athena console returns the aggregated result. The following screenshot shows the gold query and its three result rows (AMER, APAC, and EMEA), matching the values in the preceding table.

Amazon Athena console showing the gold query and its AMER, APAC, and EMEA result rows.

Figure 6: The gold table results in the Amazon Athena console.

Cost considerations

AWS Glue 6.0 bills ETL jobs by the data processing unit (DPU)-hour, per second, with a 1-minute minimum per run. AWS Glue 6.0 is also priced 30 percent lower per DPU-hour than AWS Glue 5.1, with no change to your workload, so the same job costs less to run on 6.0. This walkthrough runs on 2 G.1X workers (2 DPUs), reads a 6-row CSV, and completes each run in about 2 minutes. It produces three tables in one AWS Glue Data Catalog database.

To estimate the cost of a run, multiply the 2 DPUs by the run time in hours by your Region’s AWS Glue 6.0 DPU-hour rate. You can find that rate on the AWS Glue pricing page, and rates differ by AWS Region. The Amazon S3 objects created are the 6-row CSV, the pipeline zip, and the three tables’ data. To stop further charges, delete the resources when you finish, as shown in the next step.

Step 10 – Clean up

To avoid ongoing charges, delete the resources you created:

# Delete the AWS Glue job
aws glue delete-job --job-name "${JOB_NAME}" --region "${AWS_REGION}"

# Delete the Data Catalog database and its table metadata
aws glue delete-database --name "${DATABASE}" --region "${AWS_REGION}"

# Remove the S3 objects
aws s3 rm "s3://${BUCKET}/${PREFIX}/" --recursive --region "${AWS_REGION}"

What’s next

You now have a single pipeline that turns raw order records into validated, aggregated analytics tables, without writing orchestration logic. From here you can:

  • Extend: Add transformation stages (additional @dp.materialized_view functions) and connect them by referencing upstream tables. The pipeline picks up the new dependency automatically.
  • Scale: This walkthrough uses materialized views throughout, so every layer fully recomputes on each run (materialized views don’t support incremental refresh). To process only new data as it arrives, convert the bronze layer to a streaming table, which maintains state across runs with checkpoints. For that cross-run state to persist, a streaming table’s data and checkpoint state must not be stored locally. Hive or AWS Glue managed tables require the database’s LocationUri to point to an Amazon S3 path, while Apache Iceberg tables manage their table metadata themselves.
  • Govern: Protect the Data Catalog tables SDP produces with AWS Lake Formation fine-grained access control. It enforces table-, row-, column-, and cell-level permissions on read queries in AWS Glue Spark jobs (Glue 5.0 and later, for Hive and Iceberg tables). Because this enforcement covers batch reads, it applies to SDP’s materialized views but not to streaming tables, which read through Spark Structured Streaming.
  • Automate: Store the pipeline project in source control. Have your continuous integration and continuous delivery (CI/CD) pipeline package and upload it to Amazon S3 so each job run maps to a known build. Version the zip by object key, or upload the unzipped project to an S3 prefix and turn on Amazon S3 bucket versioning.
  • Monitor: Use Amazon CloudWatch metrics and AWS Glue job run insights for pipeline observability, latency tracking, and failure alerting.

Conclusion

In this post, you used Spark Declarative Pipelines, the declarative alternative to explicitly orchestrated ETL, now available in AWS Glue 6.0. Two decorated Python functions and one SQL file define the bronze, silver, and gold datasets, and SDP resolves the dependencies and manages execution order for you.

With SDP, you declare what each dataset should contain and the declarative framework handles ordering and execution. A three-layer pipeline that would otherwise need separate transform and orchestration logic runs as one job that you can ship and maintain.

To get started, open the AWS Glue console and build the walkthrough pipeline, or adapt the pattern to your own bronze, silver, and gold datasets. For the full set of features, see the AWS Glue 6.0 launch announcement. To move existing jobs to the Spark 4.1 runtime, see Upgrade AWS Glue jobs to AWS Glue 6.0 with AI-powered Spark upgrades. For job configuration details, see the AWS Glue Developer Guide.


About the authors

Syed Humair

Syed Humair

Syed is a Senior Analytics Specialist Solutions Architect at Amazon Web Services, based in Dubai. He has nearly 20 years of experience in data strategy, data engineering, AI, and enterprise architecture across industries including financial services, retail, telecom, and healthcare. At AWS, he works with enterprise customers to build AI-ready data foundations, from lakehouse architectures and open data formats to real-time analytics and data governance. He is the co-author of the AWS Certified Data Engineer Study Guide (Wiley, 2025).

Shrey Malpani

Shrey Malpani

Shrey is a Senior Product Manager Technical at Amazon Web Services (AWS), where he works at the intersection of distributed data processing and data integration. He helps customers build AI-ready data platforms for analytics and machine learning. His focus is scaling data integration and data management across services like AWS Glue, Amazon EMR, and Amazon Redshift.

Bo Li

Bo Li

Bo is a Senior Software Development Engineer on the AWS Glue team. He is devoted to designing and building end-to-end solutions to address customers’ data analytic and processing needs with cloud-based, data-intensive and generative AI technologies.

Kartik Panjabi

Kartik Panjabi

Kartik is a Software Development Manager on the AWS Glue team. His team builds generative AI features for data integration and distributed systems for data integration.

Build a real-time event pipeline with Spark Real-Time Mode on AWS Glue 6.0

Post Syndicated from Shoukat Ghouse original https://aws.amazon.com/blogs/big-data/build-a-real-time-event-pipeline-with-spark-real-time-mode-on-aws-glue-6-0/

Real-time event pipelines rarely get to work with a uniform schema. Whether it’s IoT metrics, ecommerce clickstreams, or financial pricing vectors, each event type brings its own schema. An equity trade and a rates trade, for instance, carry almost entirely different fields. Ingesting these multi-schema streams has traditionally forced suboptimal architectural choices. You build separate tables for each event type or maintain a wide STRUCT where every possible field across all event types must be declared upfront (fast reads, but sparse and rigid). The other option is to flatten everything into an unwieldy schema with hundreds of columns. To sidestep that maintenance burden, many teams dump events into a plain JSON string column that introduces significant performance penalty. Querying a single nested field requires your engine to deserialize the entire JSON blob for every row. At scale, you burn compute and budget scanning terabytes of raw text to extract a few bytes of data.

Adding to the challenge, these pipelines typically demand mixed processing speeds. You need a real-time path (not near-real-time) to flag anomalies or high-risk events with sub-second latency, while simultaneously pushing those same events into analytical storage for deep historical analysis in batch.

With AWS Glue 6.0, you can tackle all of these challenges (schema heterogeneity, JSON scanning overhead, and mixed-latency requirements) from a single pipeline. Built on Apache Spark 4.1 with Apache Iceberg v3 support, AWS Glue 6.0 brings Variant columns, Variant shredding, Spark Real-Time Mode (RTM), and Arrow-native user-defined functions (UDFs) to a fully managed, serverless environment.

In this post, we walk you through how to build this multi-layer architecture using a financial services use case: a market risk pipeline processing trade pricing vectors. While the example is finance, the patterns apply wherever you deal with heterogeneous schemas, expensive JSON parsing, and mixed real-time/batch requirements such as IoT device fleets, multi-tenant SaaS platforms, logistics tracking, and beyond. We will show you how to flag high-risk trades with sub-second latency, stream everything into an Iceberg v3 data lake as Variants, and run batch Value at Risk (VaR) computations efficiently using Arrow-native UDFs.

Solution overview

A bank’s Market Risk team receives a continuous stream of trade pricing vectors from front-office systems. Each trade event carries:

  1. A trade ID and book/desk IDs.
  2. A pricing vector as semi-structured data. The schema varies by asset class (for example, equities carry risk sensitivities known as Greeks such as delta/gamma, foreign exchange (FX) carries volatility surfaces, rates carry curve sensitivities).
  3. A region ID for jurisdictional reporting (uses a column DEFAULT value, so rows that omit it get auto-populated).

The team needs three things from this stream, each at a different speed. We build the pipeline in three layers, each addressing a distinct requirement with a purpose-built AWS Glue 6.0 capability.

Layer 1: Real-time trade position breach detection (sub-second latency)

Positions must be updated in sub-second time, not seconds of traditional micro-batch streaming. For a team monitoring position limits, those seconds mean trades can breach limits before the system reacts. Spark Real-Time Mode (RTM) eliminates the micro-batch boundary entirely, letting records flow continuously through the pipeline so that high-risk trades trigger alerts within sub-second latency of arrival.

Layer 2: Near-real-time analytical lakehouse (seconds latency)

Every trade must land in a queryable data lake within seconds, with heterogeneous pricing vectors stored without declaring a fixed schema upfront. The Iceberg v3 Variant type handles this natively. The raw semi-structured payload goes into a single column regardless of asset class schema. At write time, Variant shredding automatically extracts fields observed in the data into typed Parquet columns, so downstream analytical queries read only the columns they need without deserializing the full blob. Trade amendments and cancellations are handled at a lower cost with deletion vectors (merge-on-read), and column DEFAULT values reduce boilerplate in ingestion code.

Layer 3: Batch risk computation (minutes to hours latency)

Risk metrics like Value at Risk (VaR) must be computed in Python across millions of trades. Traditional row-by-row pickle serialization between the Java Virtual Machine (JVM) and Python is the bottleneck. Arrow-native UDFs process data as vectorized columnar batches, eliminating serialization overhead and accelerating Python-based risk calculations.

The solution uses three separate AWS Glue 6.0 jobs, each independently scalable:

  • Real-time path (Scala, gluestreaming): Reads trades from Amazon Managed Streaming for Apache Kafka (Amazon MSK), enriches them with risk scores and breach flags, and writes alerts to a downstream Kafka topic. It runs with a fixed set of workers that are always on. Downstream fraud detection and position limit systems consume the alerts topic for real-time blocking decisions.
  • Near-real-time path (PySpark, gluestreaming): Reads from the same MSK topic and lands the full trade history into an Iceberg v3 table. It uses Glue auto scaling and can scale down between batches, keeping costs lower.
  • Batch analytics (PySpark, glueetl): Reads from the Iceberg v3 table, extracts fields using variant_get, computes VaR across the portfolio, and writes aggregated risk reports to a downstream summary table.

The following diagram illustrates the solution architecture.

Architecture diagram of a real-time market risk pipeline on AWS Glue 6.0. All compute runs inside a VPC within an AWS Account. A Sample Trades Producer (AWS Glue job, simulating front office trading systems) publishes to an Amazon MSK topic named trade-risk-vectors. From MSK, three processing paths branch out. The Real-Time Path uses an AWS Glue 6.0 Spark Real-Time Mode job in Scala for continuous processing (JSON extraction, risk scoring, breach detection), writing alerts to an Amazon MSK trade-alerts topic that feeds CloudWatch Alarms, SNS notifications, and position limit systems. The Near-Real-Time Path uses an AWS Glue 6.0 micro-batch PySpark job that applies PARSE_JSON to Variant, TIMESTAMP_NTZ with nanosecond precision, and shredding, writing to an Amazon S3 Apache Iceberg v3 table named trade_risk_vectors. The Batch Consumption path reads that table with an AWS Glue 6.0 PySpark job using variant_get extraction and an Arrow UDF for Value at Risk computation and jurisdiction classification, writing to an Amazon S3 Iceberg v3 table named daily_risk_summary that feeds downstream analytics. Amazon S3, AWS Glue Data Catalog, and CloudWatch are regional services shown outside the VPC but inside the AWS Account, accessed privately through VPC endpoints.

Figure 1: Real-time market risk pipeline on AWS Glue 6.0

Prerequisites

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

  1. An AWS account in a Region where AWS Glue 6.0 is available.
  2. An AWS Identity and Access Management (IAM) role with permissions to deploy AWS CloudFormation stacks and create resources including AWS Glue, Amazon MSK, AWS Lambda, Amazon Simple Storage Service (Amazon S3), and the AWS Glue Data Catalog.

Deploy the CloudFormation stack

We provide an AWS CloudFormation template that provisions all the resources needed for this walkthrough.

The stack provisions the following resources:

  • An Amazon MSK cluster with two topics: trade-risk-vectors (input) and trade-alerts (real-time alerts output).
  • An Amazon S3 bucket for Iceberg table storage and streaming checkpoints.
  • An AWS Glue database (risk_analytics_<account-id>_glue6b1).
  • An IAM role (GlueRole-<account-id>-glue6b1) with permissions for Glue, MSK, S3, and CloudWatch.
  • Virtual private cloud (VPC) networking: A Glue network connection (connection-<account-id>-glue6b1), S3 gateway endpoint, and Glue interface endpoint.
  • AWS Glue job rtm-alerts-<account-id>-glue6b1 (Scala): This job reads trades from MSK, scores risk in real time using Spark RTM, writes alerts to the trade-alerts topic.
  • AWS Glue job nrt-ingestion-<account-id>-glue6b1 (PySpark): This job reads trades from MSK, writes to Iceberg v3 table with Variant + shredding enabled.
  • AWS Glue job batch-var-<account-id>-glue6b1 (PySpark): This job reads from Iceberg v3 table, computes VaR with Arrow UDF, demonstrates deletion vectors.
  • AWS Glue job producer-<account-id>-glue6b1-helper (PySpark): This job generates sample trade events (equities, FX, rates) to the trade-risk-vectors topic.

Deploy the CloudFormation stack:

  1. Download the CloudFormation template from the GitHub repository.
  2. Sign in to the AWS CloudFormation console
  3. Choose Create stack > With new resources > Upload a template file, and upload the downloaded template.
  4. Enter the following parameters:
    • VpcId: Your VPC ID.
    • SubnetIds: At least two subnets in different Availability Zones.
    • SecurityGroupId: A dedicated security group that allows all inbound TCP traffic from itself (self-referencing rule).
    • RouteTableId: The main route table for your VPC.
  5. Acknowledge the IAM capabilities and choose Create stack.

Stack creation takes approximately 20 minutes.

After the stack completes, open the AWS Glue console and start the jobs in this order:

  1. Start rtm-alerts-<account-id>-glue6b1 and nrt-ingestion-<account-id>-glue6b1.
  2. Once both show RUNNING, start producer-<account-id>-glue6b1-helper.
  3. After the producer finishes (~3.5 minutes), run batch-var-<account-id>-glue6b1 for risk aggregation.

The consumers must be running before the producer starts so that trades are scored in real time and landed in the Iceberg table as they arrive. The batch job runs last because it reads from the Iceberg table that the near-real-time path populates.

Understand the Iceberg v3 table design

The CloudFormation stack provisions Glue jobs that create two Iceberg v3 tables, trade_risk_vectors (primary trade store) and daily_risk_summary (batch VaR output), using new data types and features:

  1. VARIANT: Stores semi-structured pricing vectors without requiring a fixed schema.
  2. DEFAULT values: Automatically applies provided defaults when fields aren’t provided.
  3. Deletion vectors (merge-on-read): Enables fast row-level updates and deletes.

Open the AWS Glue console under Data Catalog > Tables > trade_risk_vectors.

AWS Glue Data Catalog console showing the trade_risk_vectors table with its Variant and default-valued columns

Figure 2: The trade_risk_vectors table in the AWS Glue Data Catalog

The following is the Create Table command:

CREATE TABLE {TABLE} (
    trade_id STRING, book_id STRING, desk STRING,
    asset_class STRING DEFAULT 'UNKNOWN',
    execution_time STRING,
    pricing_vector VARIANT,
    var_contribution DOUBLE DEFAULT 0.0,
    risk_weight DOUBLE DEFAULT 1.0,
    trade_date DATE, region STRING DEFAULT 'EMEA'
) USING iceberg
TBLPROPERTIES ('format-version'='3', 'write.delete.mode'='merge-on-read',
    'write.update.mode'='merge-on-read',
    'write.parquet.shred-variants'='true')
PARTITIONED BY (trade_date, asset_class)

Note the use of DEFAULT values for asset_class, var_contribution, risk_weight, and region. This is an Iceberg v3 feature that applies defaults automatically when values aren’t provided during writes, reducing boilerplate in ingestion code. The pricing_vector column is defined as a Variant type, and write.parquet.shred-variants='true' automatically extracts Variant fields into separate typed Parquet columns at write time for faster downstream queries.

Sample trade event generator

The CloudFormation stack includes a Glue job (producer-<accountid>-glue6b1-helper) that produces realistic trade events to the trade-risk-vectors MSK topic. Each event carries a pricing_vector with a completely different schema per asset class. This is exactly the problem Variant solves.

Equity trade (greeks, scenarios with sector/region breakdowns):

JSON pricing vector for an equity trade showing greeks and per-sector and per-region scenario breakdowns

Figure 3: Sample equity trade pricing vector

Rates trade (curve sensitivities per tenor, calibration params):

JSON pricing vector for a rates trade showing curve sensitivities per tenor and calibration parameters

Figure 4: Sample rates trade pricing vector

Completely different structures: greeks vs curve sensitivities, BlackScholes vs HullWhite. Both land in the same pricing_vector VARIANT column with no schema changes required.

Ingest trades with Spark Real-Time Mode

Traditional Spark Structured Streaming uses micro-batches: collect records, schedule a job, process, commit, wait. Even with small batches, the fixed overhead of planning and scheduling adds noticeable latency per batch. For a risk team monitoring position limits, the delay can let a trade breach a limit before the system reacts.

The following Scala job reads trade events from Amazon MSK, applies lightweight risk rules based on data directly available in the event, and writes alerts to a Kafka topic, all with sub-second latency. The real-time path intentionally avoids external lookups (market data, volatility surfaces) to stay fast. The full VaR computation happens later in the batch layer where latency is less critical.

You can view the complete job code in the AWS Glue console under the rtm-alerts-<accountid>-glue6b1 job. Additionally, all the scripts are available in the GitHub repository.

Scala real-time job code that reads from Kafka, scores risk, and writes alerts to a Kafka topic

Figure 5: Scala real-time job that scores trades and writes alerts

The Trigger.RealTime("1 minute") is what distinguishes this from a traditional micro-batch. Records flow through the pipeline continuously. Records are processed the instant they arrive. The 1-minute parameter controls how often Spark checkpoints its progress for recovery. It does not control how often records are processed. RTM on AWS Glue 6.0 currently supports Kafka-source, stateless, Scala workloads with fixed workers (no auto scaling) and update output mode only. This makes it ideal for stateless transformations that require sub-second latency, such as the filter, enrich, score, and route pattern shown here. The heavier computation (VaR, aggregations) runs in the micro-batch/batch layer where sub-second latency is less critical.

The real-time path acts as a circuit breaker: trades over $50M notional are flagged CRITICAL, over $25M flagged HIGH. Downstream systems consume the trade-alerts topic and can block or escalate before the next trade executes. The detailed VaR computation (which requires market data, volatility surfaces, and the full pricing vector) runs in the batch consumption layer where latency is less sensitive.

After the streaming phase completes, the job reads back from the trade-alerts topic and measures end-to-end latency. It compares two MSK timestamps: when the trade was received by MSK from the producer, and when the alert was received by MSK from RTM.

To verify the alerts and latency, open the Amazon CloudWatch console > Log groups > /aws-glue/jobs/output and select the RTM job’s log stream. You will see the alert summary showing each flagged trade with its end-to-end latency. The following is a sample.

CloudWatch log output listing flagged trades with CRITICAL and HIGH labels and their end-to-end latency

Figure 6: CloudWatch output showing flagged trades and end-to-end latency

Store trades in Iceberg v3 with Variant shredding enabled

The near-real-time path reads from the same MSK topic but writes to an Iceberg v3 table using standard micro-batch streaming. This job runs separately with auto scaling enabled, scaling between batches, keeping costs lower than the always-on real-time path.

You can view the complete job code in the AWS Glue console under the nrt-ingestion-<accountId>-glue6b1 job. The critical aspects are the Variant conversion and the Iceberg write:

PySpark code applying PARSE_JSON to build a Variant column and writing to the Iceberg v3 table

Figure 7: Near-real-time PySpark job writing trades to Iceberg v3 as a Variant

The PARSE_JSON() function converts the raw pricing vector into a native Variant, regardless of the asset class schema. Whether the incoming trade is an equity with greeks, an FX option with a volatility surface, or a rates swap with curve sensitivities, it all goes into the same column. Since the table has write.parquet.shred-variants enabled, fields observed in the initial sample are automatically extracted into typed Parquet columns for fast downstream queries.

How shredding works

During the write process, Spark automatically extracts the Variant fields it observes into separate typed Parquet columns at write time, a feature called shredding. At the start of each write, the engine buffers a sample of rows (controlled by write.parquet.variant-inference-buffer-size), infers which fields exist and their types, then uses that schema to shred all subsequent rows in the file. Every field observed in that sample gets its own typed column, including nested objects. Rows that lack a particular field simply store NULL in that shredded column. For our risk table, fields like $.greeks.delta, $.dv01, and $.model all live in their own typed Parquet columns, even if only one asset class carries a specific field. The result: faster read performance because queries access only the typed columns they need, skipping the rest of the document entirely. Shredding is transparent to queries. variant_get() calls work the same way whether the field is shredded or not. The query engine automatically routes to the shredded column when available, falling back to the binary Variant blob for fields that aren’t part of the inferred schema.

Note: Shredding adds write latency because the engine must infer the schema and write additional typed columns. In this pipeline, we enable shredding on the near-real-time path and absorb that cost, since the downstream read benefits (batch VaR, ad-hoc queries, audit) far outweigh the write penalty. For latency-sensitive pipelines where every millisecond on the write path matters, you can disable shredding on the streaming table and instead write shredded data in a separate batch job that reads from the unshredded table and inserts into a shredded copy. This approach trades architectural simplicity for lower ingestion latency.

Build the batch consumption layer

The third AWS Glue 6.0 job reads from the Iceberg v3 table, extracts risk metrics from the Variant column, computes VaR using an Arrow-native UDF, and writes aggregated results to a summary table.

You can view the complete job code in the AWS Glue console under the batch-var-<accountid>-glue6b1 job. The key aspects are the variant_get extraction from deeply nested structures and the Arrow-native UDF:

PySpark code using variant_get to extract deeply nested fields from the Variant column

Figure 8: Extracting nested Variant fields with variant_get

Notice how variant_get reaches into arbitrarily nested structures: $.greeks.delta (2 levels), $.scenarios[0].breakdown.by_sector.financials (5 levels), $.model_params.calibration.fit_error (4 levels). All with the same function call. No pre-flattening, no schema-per-asset-class tables, no ETL to restructure the data before querying.

Once the risk metrics are extracted, we need to run a Monte Carlo-style Historical VaR that simulates 1,000 daily profit and loss (P&L) scenarios per trade and returns the 99th percentile loss. This is where the @arrow_udf decorator comes in.

Python Arrow UDF code running a Monte Carlo Historical VaR simulation for each trade

Figure 9: Arrow-native UDF computing Historical VaR

The @arrow_udf decorator is new in Spark 4.1. Your function receives and returns pyarrow.Array directly, operating on the entire batch of rows at once. There is no pickle serialization, no row-by-row invocation, and no Pandas conversion. Data flows as native Arrow columnar arrays between the JVM and Python. For compute-heavy operations like VaR across hundreds of thousands of rows, this can be significantly faster than traditional scalar UDFs. Additionally, you can use the built-in UDF profilers to identify performance and memory bottlenecks in compute-heavy UDFs such as VaR calculations.

Handle late trade corrections with deletion vectors

In financial markets, trade amendments and cancellations are common. The batch VaR job demonstrates this after completing the risk computation. It amends one trade and cancels another.

PySpark code amending one trade and cancelling another in the Iceberg table

Figure 10: Amending and cancelling trades with merge-on-read

Prior to Iceberg v3, row-level deletes required either rewriting entire data files (copy-on-write) or maintaining separate positional delete files that store (file_path, row_position) pairs as Parquet rows (merge-on-read). Both approaches are expensive at scale. Copy-on-write rewrites gigabytes for a single amendment, and positional deletes degrade read performance as delete files accumulate (each read must parse and hash-join all delete records against the data file).

Because we configured the table with write.delete.mode='merge-on-read', UPDATEs and DELETEs write deletion vectors instead of positional delete files used in Iceberg v2. A deletion vector is a Roaring Bitmap stored in a Puffin file (.puffin), one per affected data file, marking which row positions are deleted. At read time, the engine loads a single bitmap and skips flagged positions with a bit check. No file joins, no linear scan through multiple delete files. The bitmap is compact regardless of how many rows are deleted, and read performance remains predictable as amendments accumulate.

The batch job also verifies the deletion vectors were created. You can see the results in the job’s output logs.

Job output confirming deletion vector Puffin files were created for the affected data files

Figure 11: Output verifying deletion vectors were created

Clean up

To avoid incurring further charges, delete the CloudFormation stack. This removes all resources provisioned as part of this post, including the S3 bucket, Glue jobs, Iceberg tables, MSK cluster, and IAM roles.

Conclusion

In this post, we built a multi-layer market risk pipeline using AWS Glue 6.0 (real-time alerting, near-real-time ingestion, and batch analytics):

  • Spark Real-Time Mode (RTM) on the real-time path delivers sub-second trade scoring and breach alerting, eliminating the micro-batch boundary so position limits are enforced before the next trade executes.
  • Iceberg v3 Variant on the near-real-time path stores heterogeneous pricing vectors without schema flattening. One table handles equities, FX, and rates with different schemas per row.
  • Variant shredding delivers faster reads by automatically extracting fields into separate typed Parquet columns at write time with no manual tuning required.
  • Arrow-native UDFs eliminate pickle serialization overhead for Python-based risk calculations, processing data as vectorized columnar batches on the batch layer.
  • Deletion vectors handle trade amendments and cancellations without costly data file rewrites, using compact Roaring Bitmaps instead of accumulating positional delete files.
  • Default values reduce boilerplate in ingestion code.

To get started with AWS Glue 6.0, see the AWS Glue documentation. For more information about Apache Iceberg v3, see the Iceberg specification.


About the authors

Shoukat Ghouse

Shoukat Ghouse

Shoukat is a Senior Specialist Solutions Architect for Big Data, Analytics, and Data Governance at Amazon Web Services (AWS). He partners with enterprise and financial services customers worldwide to design and scale production-grade data lakehouse platforms on Apache Spark, Apache Iceberg, AWS Glue, Amazon Athena, Amazon EMR and Amazon SageMaker Unified Studio. His focus spans distributed data processing, fine-grained data governance, and helping organizations build AI-ready data foundations that power analytics and machine learning at scale.

Shrey Malpani

Shrey Malpani

Shrey is a Senior Product Manager Technical at Amazon Web Services (AWS), where he works at the intersection of distributed data processing and data integration. He is focused on building and scaling data integration and data management capabilities across services like AWS Glue, Amazon EMR, and Amazon Redshift that help customers build AI-ready data platforms for their analytics and machine learning workflows.

Danylo Prozorov

Danylo Prozorov

Danylo is a Software Development Engineer at Amazon Web Services (AWS), where he works at the intersection of distributed data processing, AI-powered Spark troubleshooting, and AI-driven engineering automation. He focuses on the AWS Glue data integration libraries and AI-powered Spark troubleshooting capabilities across AWS Glue and Amazon EMR, delivering scalable and reliable data integration for customers’ ETL and analytics workloads.

Kartik

Kartik

Kartik is a Software Development Manager on the AWS Glue team. His team builds generative AI features for the Data Integration and distributed system for data integration.

Build with geospatial and variant types in Iceberg v3 on AWS Glue 6.0

Post Syndicated from Shoukat Ghouse original https://aws.amazon.com/blogs/big-data/build-with-geospatial-and-variant-types-in-iceberg-v3-on-aws-glue-6-0/

As organizations build data lakes that combine geospatial data, high-frequency event streams, and heterogeneous payloads, the limitations of older table formats become acute. Without a native geospatial type, coordinates require separate float columns (latitude/longitude) with no spatial predicates. Without nanosecond-precision timestamps, sub-microsecond event ordering is lost. Without a variant type, semi-structured data forces a choice between rigid flattening and untyped JSON strings. Each workaround adds complexity, slows queries, and increases maintenance burden.

AWS Glue 6.0, powered by Apache Spark 4.1, removes these workarounds by adding support for Apache Iceberg v3, bringing new column-level capabilities to your data lake tables. These include new data types: native geospatial types (GEOMETRY with spatial predicates, and GEOGRAPHY), nanosecond-precision timestamps, and the VARIANT type for semi-structured data with automatic shredding. Iceberg v3 also adds support for DEFAULT column values. These are table format features. After they’re written, they’re readable by any Iceberg v3-compatible engine that supports these features.

In this post, we build a connected vehicle fleet monitoring pipeline that uses these capabilities in a single Iceberg v3 table. Vehicles emit telemetry events with GPS coordinates (geospatial), sub-microsecond event times (nanosecond), and sensor payloads that vary by vehicle type (variant). We ingest these events, run spatial queries to detect geofence violations, sequence events at nanosecond precision, and extract typed metrics from heterogeneous payloads, all without workarounds, flattening, or external libraries.

Solution overview

A logistics company operates a mixed fleet of delivery vehicles: vans, electric bikes, and delivery robots. Each vehicle type produces telemetry events with a different sensor payload schema. The operations team needs to:

  1. Detect geofence violations: flag vehicles that enter restricted zones (airports, pedestrian areas, private property).
  2. Sequence events precisely: at fleet scale, many events land in the same microsecond window. Nanosecond timestamps give a deterministic order and prevent ties when sequencing or deduplicating events during processing.
  3. Extract metrics from heterogeneous payloads: query battery level from delivery robots, fuel level from vans, and pedal cadence from bikes, all stored in the same column.

We address all three requirements with a single Iceberg v3 table on AWS Glue 6.0. The following data definition language (DDL) shows the table structure. The AWS Glue job we provision in subsequent steps executes this statement.

CREATE TABLE fleet_monitoring_db.vehicle_telemetry (
event_id STRING,
vehicle_id STRING,
vehicle_type STRING DEFAULT 'UNKNOWN',
event_time TIMESTAMP_NTZ(9),
location GEOMETRY(4326),
service_area GEOGRAPHY(4326),
sensor_payload VARIANT,
speed_kmh DOUBLE DEFAULT 0.0,
region STRING DEFAULT 'EMEA'
) USING ICEBERG
TBLPROPERTIES (
'format-version' = '3',
'write.delete.mode' = 'merge-on-read'
)
PARTITIONED BY (days(event_time), vehicle_type)

In the preceding statement, the database is shown as fleet_monitoring_db for readability. The deployed stack creates it as fleet_monitoring_<account-id>.

The following list describes the key columns:

  • event_time TIMESTAMP_NTZ(9): Stores the event timestamp at nanosecond precision.
  • location GEOMETRY(4326): Stores GPS coordinates as native spatial objects using (SRID 4326). You can use predicates like ST_Intersects directly in SQL, replacing hand-coded spatial math on raw latitude/longitude doubles (WGS 84).
  • service_area GEOGRAPHY(4326): Stores geographic coordinates using a spherical (geodesic) model, distinct from GEOMETRY’s planar model. AWS Glue 6.0 writes and reads GEOGRAPHY in Iceberg v3, and the type is portable to any Iceberg v3-compatible engine. Geodesic spatial predicates over GEOGRAPHY are engine-dependent today. In this post we run spatial queries on the GEOMETRY location column, which Glue 6.0 supports natively.
  • sensor_payload VARIANT: Each vehicle type produces a different JSON schema. Vans report fuel and engine metrics, robots report battery and camera status, bikes report cadence and heart rate. All land in this single column without schema unions or separate tables using variant data type.
  • vehicle_type STRING DEFAULT ‘UNKNOWN’ and speed_kmh DOUBLE DEFAULT 0.0: When an ingestion writer omits these fields, Iceberg applies the declared defaults automatically. Useful when multiple producers write to the same table and not all of them populate every column.

The table uses PARTITIONED BY (days(event_time), vehicle_type) so that analytical queries can prune by date range and vehicle type without scanning the full table. 'write.delete.mode' = 'merge-on-read' supports fast row-level corrections (for example, correcting a misreported GPS coordinate) through compact deletion vectors (Roaring Bitmaps) instead of accumulating positional delete files.

In this post, we insert sample data directly to focus on the new Iceberg data types and how to use them together. In production, these events would stream from Amazon Managed Streaming for Apache Kafka (Amazon MSK) into an AWS Glue 6.0 streaming job.

The following diagram illustrates the production architecture for reference:

Architecture diagram showing a vehicle fleet of vans, delivery robots, and electric bikes sending telemetry through Amazon MSK into an AWS account. Within a VPC, a hot path uses AWS Glue 6.0 Spark Real-Time Mode to detect geofence violations and send alerts to a Kafka topic, while a cold path uses a Glue 6.0 micro-batch job to write events into an Apache Iceberg v3 table with GEOMETRY, TIMESTAMP_NTZ(9), VARIANT, and DEFAULT columns. Amazon S3 stores the Iceberg data and the AWS Glue Data Catalog holds metadata. A batch analytics Glue job reads the Iceberg table for geofence detection, nanosecond event sequencing, and per-vehicle-type metric extraction using variant_get

Figure 1: Reference architecture for a fleet telemetry pipeline on AWS Glue 6.0

The architecture processes vehicle telemetry through two paths, with a downstream batch analytics layer:

Hot path (real-time, milliseconds): A Spark Real-Time Mode (RTM) job reads telemetry from Amazon MSK and evaluates geofence violations using spatial predicates like ST_Intersects, routing alerts to a downstream Kafka topic within milliseconds.

Cold path (near-real-time, seconds): A micro-batch job reads the same MSK topic and writes events into an Iceberg v3 table, converting payloads to GEOMETRY, TIMESTAMP_NTZ(9), and VARIANT columns with DEFAULT values applied.

Batch analytics: An AWS Glue job reads the Iceberg v3 table to run batch analytics on geofence detection, nanosecond event sequencing, and per-vehicle-type metric extraction.

Prerequisites

To follow along, you need:

  • An AWS account and an AWS Region where AWS Glue 6.0 is available.
  • An AWS Identity and Access Management (IAM) role with permissions to deploy AWS CloudFormation stacks and create resources including AWS Glue, Amazon Simple Storage Service (Amazon S3), and Amazon CloudWatch Logs.

Deploy the CloudFormation stack

We provide an AWS CloudFormation template that provisions all the resources needed for this walkthrough.

The stack provisions the following resources:

  • An Amazon S3 bucket for Iceberg table storage.
  • An IAM role with permissions for AWS Glue, Amazon S3, and Amazon CloudWatch Logs.
  • An AWS Glue database (fleet_monitoring_<account-id>).
  • An AWS Glue job fleet-telemetry-ingest-<account-id> (PySpark): creates the Iceberg v3 table vehicle_telemetry described earlier and inserts sample telemetry from three vehicle types.
  • An AWS Glue job fleet-telemetry-queries-<account-id> (PySpark): demonstrates geofence detection, nanosecond sequencing, variant extraction, and default values.

Deploy the CloudFormation stack:

  1. Download the CloudFormation template from the GitHub repository.
  2. Sign in to the AWS CloudFormation console.
  3. Choose Create stack, With new resources, Upload a template file, and upload the downloaded template.
  4. Acknowledge the IAM capabilities and choose Create stack.

Stack creation takes approximately 2–5 minutes. No parameters are required.

After the stack completes, navigate to the AWS Glue console and run the jobs in this order:

  1. Run fleet-telemetry-ingest-<account-id>. This job creates the Iceberg v3 table and inserts sample data (approximately 2 minutes).
  2. After it succeeds, run fleet-telemetry-queries-<account-id>. This job executes all demonstration queries (approximately 2 minutes).

The following sections describe each job in detail.

Job 1: Ingest sample telemetry data

The ingestion job creates the Iceberg v3 table described earlier and inserts four sample telemetry events: one for each of the three vehicle types (van, robot, bike), plus one with omitted fields to demonstrate DEFAULT values. You can view the complete script in the GitHub repository. Note that the geospatial types require one additional Spark configuration (spark.sql.geospatial.enabled=true), which is already set in the job’s --conf argument by the CloudFormation template. All other types work with no extra configuration.

The following are the key snippets from the script:

Van telemetry: GPS coordinates with engine metrics and route information:

spark.sql(f"""
INSERT INTO {TABLE} VALUES (
'EVT-001', 'VAN-042', 'VAN',
CAST('2026-07-28 09:15:30.123456789' AS TIMESTAMP_NTZ(9)),
ST_SetSrid(ST_GeomFromWKB(X'0101000000E17A14AE47E1C0BF1F85EB51B84E4940'), 4326),
ST_SetSrid(ST_GeogFromWKB(X'0101000000E17A14AE47E1C0BF1F85EB51B84E4940'), 4326),
PARSE_JSON('{{"fuel_pct": 0.72, "cargo_kg": 450, "door_open": false,
"engine": {{"rpm": 2100, "temp_c": 88.5}},
"route": {{"stops_remaining": 4, "eta_minutes": 35}}}}'),
35.2, 'EMEA'
)
""")

Delivery robot telemetry: Same table, completely different sensor schema (battery, cameras, navigation):

spark.sql(f"""
INSERT INTO {TABLE} VALUES (
'EVT-002', 'ROB-117', 'ROBOT',
CAST('2026-07-28 09:15:30.123456790' AS TIMESTAMP_NTZ(9)),
ST_SetSrid(ST_GeomFromWKB(X'01010000000000000000001040000000000000F03F'), 4326),
ST_SetSrid(ST_GeogFromWKB(X'01010000000000000000001040000000000000F03F'), 4326),
PARSE_JSON('{{"battery_pct": 0.62, "obstacle_distance_m": 2.8,
"navigation_mode": "autonomous",
"cameras": {{"front": "active", "rear": "recording"}}}}'),
48.0, 'EMEA'
)
""")

Note: EVT-001 and EVT-002 are exactly 1 nanosecond apart (.123456789 vs .123456790). Without TIMESTAMP_NTZ(9), both would round to the same microsecond and be indistinguishable.

Default values test: Event inserted with vehicle_type, speed_kmh, and region omitted:

spark.sql(f"""
INSERT INTO {TABLE}
(event_id, vehicle_id, event_time, location, service_area, sensor_payload)
VALUES (
'EVT-004', 'UNK-999',
CAST('2026-07-28 10:00:00.000000000' AS TIMESTAMP_NTZ(9)),
ST_SetSrid(ST_GeomFromWKB(X'0101000000000000000000F03F000000000000F03F'), 4326),
ST_SetSrid(ST_GeogFromWKB(X'0101000000000000000000F03F000000000000F03F'), 4326),
PARSE_JSON('{{"status": "initializing"}}')
)
""")

The omitted columns automatically receive their DEFAULT values: vehicle_type = 'UNKNOWN', speed_kmh = 0.0, region = 'EMEA'.

Job 2: Query the data

The query job demonstrates all four data types working together. After the job succeeds, select the run in the AWS Glue console and choose Output logs to see the results.

The following sections walk through the key queries from the job and the results of each.

Geofence detection with ST_Intersects

The job defines a polygon and finds all vehicles inside it:

POLY = "010300...."
SELECT event_id, vehicle_id, vehicle_type, speed_kmh
FROM fleet_monitoring_db.vehicle_telemetry
WHERE ST_Intersects(
location,ST_SetSrid(ST_GeomFromWKB(X'{POLY}'), 4326)
)
ORDER BY event_id

The polygon covers coordinates (0,0)-(5,0)-(5,2)-(0,2). Three vehicles are inside (ROBOT at (4,1), BIKE at (3,1), UNKNOWN at (1,1)). The VAN at (-0.1278, 51.5074) is outside.

Query results listing the ROBOT, BIKE, and UNKNOWN vehicles inside the geofence polygon, with the VAN excluded

Figure 2: Geofence query results showing the three vehicles inside the polygon

Nanosecond event sequencing

Order events by their sub-microsecond timestamps:

SELECT event_id, vehicle_id, CAST(event_time AS STRING) AS precise_time
FROM fleet_monitoring_db.vehicle_telemetry
WHERE event_id IN ('EVT-001', 'EVT-002', 'EVT-003')
ORDER BY event_time ASC

EVT-001 and EVT-002 are correctly distinguished and ordered despite being only 1 nanosecond apart. With standard TIMESTAMP_NTZ (microsecond precision), both would show .123456 and their relative order would be undefined.

Query results showing EVT-001 and EVT-002 ordered by nanosecond-precision timestamps one nanosecond apart

Figure 3: Nanosecond-precision ordering distinguishing two events one nanosecond apart

Variant extraction with variant_get

Different sensor schemas per vehicle type, all extracted with variant_get:

SELECT vehicle_id, vehicle_type,
CASE vehicle_type
WHEN 'VAN' THEN variant_get(sensor_payload, '$.fuel_pct', 'DOUBLE')
WHEN 'ROBOT' THEN variant_get(sensor_payload, '$.battery_pct', 'DOUBLE')
WHEN 'BIKE' THEN variant_get(sensor_payload, '$.battery_pct', 'DOUBLE')
ELSE NULL
END AS energy_level,
variant_get(sensor_payload, '$.engine.temp_c', 'DOUBLE') AS engine_temp,
variant_get(sensor_payload, '$.cameras.front', 'STRING') AS front_cam,
variant_get(sensor_payload, '$.deliveries.completed', 'INT') AS deliveries_done
FROM fleet_monitoring_db.vehicle_telemetry
WHERE vehicle_type != 'UNKNOWN'
ORDER BY vehicle_id
Query results showing variant_get extracting energy level, engine temperature, and camera status for each vehicle type

Figure 4: Variant extraction returning typed values from heterogeneous sensor payloads

variant_get takes three arguments: the column, a dot-path expression, and the expected return type. It supports arbitrary nesting depth. $.engine.temp_c reaches two levels deep, $.deliveries.completed reaches into a different structure entirely. When a path doesn’t exist in a particular row’s payload, it returns NULL.

Default values

Confirm that omitted columns received their defaults:

SELECT event_id, vehicle_type, speed_kmh, region
FROM fleet_monitoring_db.vehicle_telemetry
WHERE event_id = 'EVT-004'
Query results showing event EVT-004 with the default values UNKNOWN, 0.0, and EMEA applied

Figure 5: Default column values applied to the event inserted with omitted fields

EVT-004 was inserted without vehicle_type, speed_kmh, or region. The declared defaults were applied automatically.

Combined query: Combining spatial, temporal, and variant operations

The following query runs a geospatial predicate, nanosecond ordering, and variant extraction in a single SELECT statement:

SELECT vehicle_id, vehicle_type,
CAST(event_time AS STRING) AS precise_time,
CASE vehicle_type
WHEN 'VAN' THEN variant_get(sensor_payload, '$.fuel_pct', 'DOUBLE')
WHEN 'ROBOT' THEN variant_get(sensor_payload, '$.battery_pct', 'DOUBLE')
WHEN 'BIKE' THEN variant_get(sensor_payload, '$.battery_pct', 'DOUBLE')
ELSE NULL
END AS energy_level,
speed_kmh
FROM fleet_monitoring_db.vehicle_telemetry
WHERE ST_Intersects(location, ST_SetSrid(ST_GeomFromWKB(X'0103000000...'), 4326))
ORDER BY event_time ASC
Query results combining spatial filtering, nanosecond ordering, and variant extraction in a single query

Figure 6: Combined query results over a single Iceberg v3 table

This single query combines a spatial predicate, nanosecond ordering, and variant extraction over one table, with no external libraries, pre-processing, or joins to separate geometry or payload tables.

Clean up

To avoid ongoing charges from the AWS Glue jobs and Amazon S3 storage, delete the CloudFormation stack when you’re done:

  1. Open the AWS CloudFormation console.
  2. Select the stack you deployed earlier and choose Delete.

Conclusion

In this post, we stored and analyzed geospatial coordinates, nanosecond timestamps, and heterogeneous sensor payloads in a single Iceberg v3 table on AWS Glue 6.0, with sensible defaults applied automatically, no external libraries, and no schema flattening.

  • GEOMETRY columns replace latitude/longitude doubles and support native spatial predicates like ST_Intersects for geofence detection. GEOGRAPHY is stored natively.
  • TIMESTAMP_NTZ(9) preserves full nanosecond precision for event sequencing where microsecond resolution is insufficient.
  • VARIANT stores heterogeneous payloads (different schema per vehicle type) in one column with typed extraction through variant_get.
  • DEFAULT values keep field population consistent across multiple ingestion writers without duplicating logic.

All capabilities require Iceberg format-version 3. Geospatial requires one additional configuration (spark.sql.geospatial.enabled=true). Nanosecond timestamps, Variant, and DEFAULT values work with no extra configuration.

These capabilities apply wherever schemas vary by source (IoT fleets, multi-tenant software as a service (SaaS), event-driven architectures), timestamps need sub-microsecond precision (trading, sensor fusion, autonomous systems), or spatial operations replace coordinate workarounds (logistics, real estate, delivery networks).

For more information, see the AWS launch announcement (launch URL to be added before publishing), the AWS Glue documentation, and the Apache Iceberg v3 specification. AWS Glue 6.0 includes additional capabilities such as Spark Real-Time Mode and Spark Declarative Pipelines, which we cover in separate posts.


About the authors

Shoukat Ghouse

Shoukat Ghouse

Shoukat is a Senior Specialist Solutions Architect for Big Data, Analytics, and Data Governance at Amazon Web Services (AWS). He partners with enterprise and financial services customers across EMEA to design and scale production-grade data lakehouse platforms on Apache Spark, Apache Iceberg, AWS Glue, Amazon EMR, and Amazon SageMaker Unified Studio. His focus spans distributed data processing, fine-grained data governance, and helping organizations build AI-ready data foundations that power analytics and machine learning at scale.

Shrey Malpani

Shrey Malpani

Shrey is a Senior Product Manager Technical at Amazon Web Services (AWS), where he works at the intersection of distributed data processing and data integration. He is focused on building and scaling data integration and data management capabilities across services like AWS Glue, Amazon EMR, and Amazon Redshift that help customers build AI-ready data platforms for their analytics and machine learning workflows.

Kartik

Kartik

Kartik is a Software Development Manager on the AWS Glue team. His team builds generative AI features for the Data Integration and distributed system for data integration.

Introducing AWS Glue 6.0 for faster and more cost-effective data integration

Post Syndicated from Aarthi Srinivasan original https://aws.amazon.com/blogs/big-data/introducing-aws-glue-6-0-for-apache-spark/

Organizations running large data processing pipelines want lower costs, faster job runtimes, and dependable support for open table formats, without adding operational overhead. AWS Glue, a serverless, scalable data integration service that you can use to discover, prepare, move, and integrate data from multiple sources, has now launched AWS Glue 6.0, the new version of AWS Glue that addresses these needs. This version upgrade lowers AWS Glue pricing by 30% and improves performance with AWS optimized Apache Spark 4.1. It also augments developer experience with new features and adds support for Apache Iceberg V3 specifications that are suitable for enterprise adoption. The newly available AWS Glue 6.0 makes data processing workloads more manageable, faster to run, and easier to operate.

In this post, we cover the key capabilities of AWS Glue 6.0 and their performance benefits. We share code examples to help you take full advantage of the release, and we show you how to get started.

AWS Glue 6.0 highlights

AWS Glue 6.0 brings together four major improvements designed to transform how you build and run data integration workloads.

First, it upgrades the underlying runtime to Apache Spark 4.1.1, Python 3.13, Scala 2.13, and AWS SDK for Java 2.x, delivering performance improvements that can help with faster job completion times and lower costs.

Second, this release reduces current AWS Glue usage rate by 30%, and when combined with the performance improvements, you may realize an even lower effective cost.

Third, AWS Glue 6.0 introduces support for more capabilities of Apache Iceberg V3. This includes the VARIANT data type with automatic shredding, deletion vectors, row lineage tracking, nanosecond timestamps, and geo types. With these capabilities, you can build modern lakehouse architectures on the latest open table format standards.

Finally, new features like Spark Declarative Pipelines, Real-Time Mode for streaming and Python virtual environments with S3 caching are designed to further improve performance and developer experience. The following sections dive deeper into each of these areas.

Runtime upgrades

AWS Glue 6.0 upgrades the core runtime stack across the board, bringing newer versions of Apache Spark, Python, Scala, and the AWS SDK to your serverless data integration workloads.

  • Apache Spark 4.1.1 – AWS Glue 6.0 runs an AWS optimized build of Apache Spark 4.1.1, a major generational leap from Spark 3.5 on AWS Glue 5.1. This release introduces improvements focused on intent-driven data engineering, real-time streaming with sub-second latencies down to single-digit milliseconds for stateless tasks, faster PySpark performance, and expanded SQL features.
  • Python 3.13 – Supports Python 3.13, a stable release that brings interpreter changes, Python data model enhancements, standard library updates, and security updates.
  • Scala 2.13 – Upgrades to Scala 2.13 which includes a collections library overhaul, language and syntax feature changes, standard library additions, and compiler performance updates.

Reduced Pricing

AWS Glue 6.0 cuts current AWS Glue pricing by 30%. This means every job you run on AWS Glue 6.0 costs 30% less per DPU-hour compared to AWS Glue 5.1, with no changes required to your workload configuration. When you combine this pricing reduction with the performance improvements delivered by runtime upgrades, your effective cost savings can compound because jobs can complete faster and consume fewer DPU-hours on a lower price point. If you run large-scale Extract, Transform, and Load (ETL) pipelines or recurring batch workloads, this compounding effect can help reduce your monthly spend.

To quantify the comparison, we ran the industry-standard TPC-DS benchmark at 3 TB scale on Parquet data stored in Amazon Simple Storage Service (Amazon S3), using 30 G.2X workers on AWS Glue. The following table compares the results we obtained in our tests for AWS Glue 6.0 and AWS Glue 5.1. Thus, based on TPC-DS benchmark at 3 TB scale, AWS Glue 6.0 delivers up to 36% better price performance than AWS Glue 5.1.

. AWS Glue 6.0 AWS Glue 5.1
Estimated Cost ($) USD 5.61 USD 8.87

Table 1: 3TB TPC-DS benchmark comparison between AWS Glue 6.0 and AWS Glue 5.1

Updated Open Table Format (OTF) support

AWS Glue 6.0 ships with updated versions of all three major open table formats – Iceberg 1.11.0, Hudi 1.1.1, and Delta Lake 4.2.0 – providing better performance, improved merge-on-read capabilities, streamlined concurrency control, and expanded SQL compatibility.

Besides supporting the latest open table format versions, AWS Glue 6.0 delivers Apache Iceberg V3 specification that is suitable for enterprise use. The highlight is Variant shredding, which AWS Glue uses to automatically decompose semi-structured data into physically optimized, columnar sub-fields, which should result in faster query read performance. Combined with deletion vectors for efficient row-level updates, UNKNOWN column types, default column values, and richer data type support, AWS Glue 6.0 is designed to make your open data lakes faster, more flexible, and more cost-efficient. AWS Glue 6.0 also adds support for geospatial data types (Geometry and Geography) and nanosecond-precision timestamps from the Apache Iceberg V3 specification, neither of which are currently supported in open-source Apache Spark 4.1. Additional features like row lineage tracking round out the Apache Iceberg V3 capabilities available on AWS Glue 6.0.

In the following sections, we illustrate select capabilities from Apache Iceberg V3 specification on AWS Glue 6.0.

  1. VARIANT column type

Apache Iceberg V3 introduces the Variant type to store semi-structured data (think JSON, XML, logs, and deeply nested event data) in a compact binary format. Variant shredding is designed to automatically decompose VARIANT columns into physically optimized, columnar sub-fields, facilitating predicate pushdowns and reducing scan overhead. It aims to provide simpler management of semi-structured data, without the need for complex flattening logic. With Variant type, you get the flexibility of embedding a JSON data type in your table columns while shredding is designed to help accelerate read queries and reduce costs.

  1. UNKNOWN column type

The UNKNOWN type in Apache Iceberg V3 acts as a flexible placeholder for columns where the data type is not yet determined at the time of table creation or data ingestion. Tables can accept all-null data initially, and the column type can be upgraded later without breaking ingestion pipelines or consuming applications. This can simplify schema evolution for rapidly changing data sources. Apache Iceberg V3’s UNKNOWN column type maps to Spark 4.1’s VOID type.

  1. DEFAULT column values

Apache Iceberg V3’s DEFAULT column values allow specifying a default value for a column in the table metadata. When you add a new column, the query engine is designed to automatically apply this default to older rows, without rewriting data or running manual backfill operations.

The following code demonstrates creating an Apache Iceberg V3 table that uses VARIANT and UNKNOWN types, and DEFAULT values for a column.

Prerequisites

To get started with this code example, make sure you have the following prerequisites.

  1. An AWS account.
  2. An AWS Identity and Access Management (IAM) role with permissions for AWS Glue, the AWS Glue Data Catalog, and Amazon S3. For more information, see Setting up IAM permissions for AWS Glue. This will be the AWS Glue job execution role.
  3. An S3 bucket to store the Iceberg table data.

Steps

To create an AWS Glue 6.0 job, use the following steps.

  1. Log in to your AWS account and open the AWS Glue console.
  2. Create a new ETL job, with Script editor option.
    1. Choose engine as Spark in the drop-down menu.
    2. Start fresh, Create script and copy-paste the following code.
    3. Replace the demo S3 bucket name with your bucket name in the code.
# Example pySpark script for testing few Iceberg v3's new data types
from pyspark.sql import SparkSession

CATALOG = "glue_catalog"
DATABASE = "sample_glue6_iceberg_db"
TABLE_NAME = "sample_glue6_table"
TABLE = f"{CATALOG}.{DATABASE}.{TABLE_NAME}"
TABLE_LOCATION = "s3://amzn-s3-demo-table-bucket/glue6blog-newdatatypes/"

# Configure Spark to use Apache Iceberg with the AWS Glue Data Catalog.
spark = (
    SparkSession.builder
    .appName("Glue6NewDataTypes")
    .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
    .config(f"spark.sql.catalog.{CATALOG}", "org.apache.iceberg.spark.SparkCatalog")
    .config(f"spark.sql.catalog.{CATALOG}.catalog-impl", "org.apache.iceberg.aws.glue.GlueCatalog")
    .config(f"spark.sql.catalog.{CATALOG}.io-impl", "org.apache.iceberg.aws.s3.S3FileIO")
    .config(f"spark.sql.catalog.{CATALOG}.warehouse", "s3://amzn-s3-demo-table-bucket/glue6blog-newdatatypes")
    .config("spark.sql.defaultColumn.enabled", "true")
    .getOrCreate()
)

spark.sql(f"CREATE DATABASE IF NOT EXISTS {CATALOG}.{DATABASE}")

# Create an Iceberg v3 table with VARIANT, unknown, and a default value.
# Spark's VOID type is stored as the Iceberg v3 unknown type.
spark.sql(
    f"""
    CREATE TABLE {TABLE} (
        record_id BIGINT,
        payload VARIANT,
        reserved_field VOID,
        status STRING DEFAULT 'active'
    )
    USING ICEBERG
    LOCATION '{TABLE_LOCATION}'
    TBLPROPERTIES ('format-version' = '3')
    """
)

# Insert two rows. The omitted columns use null and the declared default.
spark.sql(
    f"""
    INSERT INTO {TABLE} (record_id, payload)
    VALUES
    (1, parse_json('{{"event_type":"created","score":98.5}}')),
    (2, parse_json('{{"event_type":"processed","score":87.2}}'))
    """
)

# Query the row and extract values from the VARIANT column.
spark.sql(
    f"""
    SELECT
        record_id,
        variant_get(payload, '$.event_type', 'string') AS event_type,
        variant_get(payload, '$.score', 'double') AS score,
        reserved_field,
        status
    FROM {TABLE}
    """
).show(truncate=False)

spark.stop()
  1. Provide the following details in the Job details tab.
    1. A Name for the job.
    2. The IAM role you have from Prerequisites (2) for the IAM role of the job.
    3. Choose Glue 6.0 for the Glue version.
    4. Leave the rest as defaults.The following screenshot shows the Job details tab with illustrated values in the AWS Glue console.

      AWS Glue Job details tab with the Glue version set to Glue 6.0 and other settings left as defaults

      Figure 1: Job details tab with the Glue version set to Glue 6.0

    5. Scroll down. Under Advanced properties, for Job parameters, add the following additional Job parameter key-value pair:--datalake-formats=icebergThe following screenshot shows the Job parameters with the illustrated key-value pair in the AWS Glue console.

      Advanced properties section showing the Job parameters key –datalake-formats set to the value iceberg

      Figure 2: Job parameters with the datalake-formats key set to iceberg

  2. Save the job and choose Run.
  3. After the job is completed successfully, from the Runs tab – Run details, you can inspect the Output logs that take you to the logs in the Amazon CloudWatch console. The following shows the sample output for the SELECT query in the script.
+---------+----------+-----+--------------+------+
|record_id|event_type|score|reserved_field|status|
+---------+----------+-----+--------------+------+
|1        |created   |98.5 |NULL          |active|
|2        |processed |87.2 |NULL          |active|
+---------+----------+-----+--------------+------+

Notice that we inserted two rows with values only in the record_id and the variant column. Variant column inserts were done using parse_json(). The reserved_field is of VOID type, hence returns NULL values. The status column is declared with a default active value and returns active, since the column was omitted during the insert operation.

  1. Deletion Vectors
    Apache Iceberg V3 replaces the traditional positional delete files used in Apache Iceberg V2 to deletion vectors. This change can help improve Merge-on-Read (MoR) performance. This shift replaces heavy, multi-file Parquet reads with highly compressed, direct binary bitmaps that can provide lower storage overhead and faster reads on delete-heavy tables. In scenarios with heavy table updates, such as streaming change data capture (CDC) from operational databases, Apache Iceberg V3 can offer read performance advantage over Apache Iceberg V2.

    To validate the performance of deletion vectors, we created two identical AWS Glue streaming jobs and ingested the events into two different Iceberg tables, one in Apache Iceberg V2 and another in Apache Iceberg V3 format. The streaming CDC events were approximately 150,000 events per second, merge-on-read, update-heavy. Every micro-batch writes row-level deletes. We froze both tables at the same delete-heavy state and disabled compaction, leaving the tables with roughly 1.7 million rows in valid state out of the 26.4 million physical rows. The following table summarizes the read performance latency of the two Iceberg tables. We observed in this testing that reading from the delete-heavy Apache Iceberg V3 is at least 1.5 times faster than the reading from a similar Apache Iceberg V2 table. For larger enterprise scale Apache Iceberg V3 tables, the read performance could improve further.

Read latency comparison showing Apache Iceberg V3 deletion vectors reading at least 1.5 times faster than Apache Iceberg V2 delete files

Table 2 – Read latency comparison between Apache Iceberg V2 delete files and Apache Iceberg V3 deletion vectors

New ETL features

AWS Glue 6.0 introduces several additional capabilities designed to simplify how you build and manage data pipelines, some of which are discussed in the following list.

  • Spark Declarative Pipelines (SDP) where you define the outcomes you want for your entire data pipelines in a declarative fashion with SQL statements or Python decorators while AWS Glue handles execution flow, dependency resolution, parallelism, checkpointing, CDC, and recovery automatically. This helps you focus on business logic rather than orchestration plumbing.
  • Real-Time Mode (RTM) for streaming delivers continuous execution for Structured Streaming with sub-second latencies, down to single-digit milliseconds for stateless tasks. This can help support real-time use cases like fraud detection, live dashboards, and event-driven architectures without managing dedicated streaming infrastructure.
  • Arrow-Native UDFs/UDTFs execute Python functions directly on PyArrow batches without Pandas conversion overhead, which can result in faster performance for custom transformation logic at scale.
  • Recursive Common Table Expressions (CTE) adds WITH RECURSIVE queries natively, allowing graph traversals and hierarchical queries without workarounds or external libraries.
  • Python data source filter pushdown evaluates filters at the data source designed to minimize data movement, reduce the volume of data scanned, and improve job performance.
  • With Python virtual environments and S3 caching, you can provide pre-built Python dependencies, which should result in reduced startup latency for AWS Glue jobs by eliminating runtime dependency resolution. For existing jobs that use --additional-python-modules, no action is required. AWS Glue automatically handles the conversion to virtual environments when your job runs on AWS Glue 6.0.

Dependent library upgrades

The following table summarizes the key runtime and library version upgrades on AWS Glue 6.0.

Feature Glue 6.0 Glue 5.1
Spark 4.1.1-amzn-0 3.5.6-amzn-1
Python 3.13.14 3.11.15
Scala 2.13.17 2.12.18
AWS SDK for Java 2.44.6
(Version 1.x removed)
2.35.5
Boto3 1.42.84 1.40.61
Java 17.0.20 17.0.19
Iceberg 1.11.0 1.10.0
Hudi 1.1.1 .0.2
Delta Lake 4.2.0 3.3.2
AWS Glue Data Catalog client 4.11.0 4.9.0
EMR DynamoDB connector 6.1.0 5.7.0
Arrow 18.3.0 2.0.1
Hive 2.3.10-amzn-1 2.3.9-amzn-4

Table 3: Runtime and library version comparison across AWS Glue 6.0 and AWS Glue 5.1

Getting started

To get started with AWS Glue 6.0, you can use one of the following methods.  

Clean up

To avoid incurring costs, clean up the resources you created for this post.

  1. Delete the Data Catalog database and the Iceberg table.
  2. Delete the data and metadata folders of the Iceberg table from your S3 bucket.
  3. Delete the AWS Glue job and the AWS Glue job execution IAM role.

Conclusion

AWS Glue 6.0 is designed to be faster, more cost-effective, and easier to use for building your open data lakehouse architectures and orchestrating your data pipelines. In this post, we discussed the key highlights of AWS Glue 6.0 and illustrated usage of Apache Iceberg V3 features with code samples. You can create new AWS Glue jobs on AWS Glue 6.0 or migrate your existing AWS Glue jobs to benefit from these improvements.

With Apache Spark 4.1.1, Apache Iceberg V3, Python 3.13, upgraded open table format libraries, and new streaming capabilities, AWS Glue 6.0 aims to help you build new data applications or to operate your existing data pipelines more efficiently and with less maintenance overhead.

We encourage you to test AWS Glue 6.0 in your development environment today. Check out this blog that talks about upgrading your AWS Glue jobs to AWS Glue 6.0. Also, in the coming days and weeks, look out for blogs on individual topics illustrating various features of Spark 4.1.1 and Apache Iceberg V3 on AWS Glue 6.0.

Acknowledgements: We thank the numerous engineers and leaders who helped build AWS Glue 6.0 to support customers with a highly performant Spark runtime and other value-added capabilities.


About the authors

Aarthi Srinivasan

Aarthi Srinivasan

Aarthi is a Senior Big Data Architect working on data, analytics and GenAI topics with the worldwide specialist org at AWS. She works with AWS customers and partners to architect open data lake solutions, enhance product features, and establish best practices for data governance and analytics services adoption.

Shrey Malpani

Shrey Malpani

Shrey is a Senior Product Manager Technical at Amazon Web Services (AWS), where he works at the intersection of distributed data processing and data integration. He is focused on building and scaling data integration and data management capabilities across services like AWS Glue, Amazon EMR, and Amazon Redshift that help customers build AI-ready data platforms for their analytics and machine learning workflows.

Angel Conde Manjon

Angel Conde Manjon

Angel is a Senior Solutions Architect at AWS where he helps partners develop businesses centered on Data and AI. He has previously worked on research related to Data Analytics and Artificial Intelligence in diverse European research projects. Angel is also an Apache Iceberg contributor.

Peter Tsai

Peter Tsai

Peter is a Software Development Engineer at AWS, where he enjoys solving challenges in the design and performance of the AWS Glue runtime. In his leisure time, he enjoys hiking and cycling.

Danylo Prozorov

Danylo Prozorov

Danylo is a Software Development Engineer at AWS Glue, where he works on building data integration and generative AI solutions for analytics customers. Outside of work, he enjoys sports, hiking, riding motorcycles, and building his overland rig.

Bo Li

Bo Li

Bo is a Senior Software Development Engineer on the AWS Glue team. He is devoted to designing and building end-to-end solutions to address customers’ data analytic and processing needs with cloud-based, data-intensive and GenAI technologies.

Kartik Panjabi

Kartik Panjabi

Kartik is a Software Development Manager on the AWS Glue team. His team builds generative AI features for the Data Integration and distributed system for data integration.

Mohit Saxena

Mohit Saxena

Mohit leads AWS Glue and AWS Data Analytics agentic AI initiatives that help customers build and operate big data applications on Apache Spark, Amazon S3, and cloud data lakes and warehouses, spanning across AWS Glue, Amazon EMR, and Amazon Athena.

Upgrade AWS Glue jobs to Glue 6.0 with AI-powered Spark upgrades

Post Syndicated from Prasad Nadig original https://aws.amazon.com/blogs/big-data/upgrade-aws-glue-jobs-to-glue-6-0-with-ai-powered-spark-upgrades/

Upgrading PySpark jobs to a new Apache Spark major version can introduce breaking changes. Removed configuration keys, stricter type casting, and Python library incompatibilities can cause runtime failures or silent behavior differences. With AWS Glue 6.0 now running Apache Spark 4.1 and Python 3.13, you need a reliable way to migrate your existing jobs while validating correctness.

In this post, we walk through upgrading a PySpark ETL job from AWS Glue 5.1 to AWS Glue 6.0. We use the generative AI upgrades for Apache Spark in the AWS Glue console. The upgrade analysis automatically identifies incompatibilities, iteratively resolves them, validates the result with data quality checks, and presents recommended changes for your review. AWS Glue 6.0 also delivers up to 36% better price performance* along with Iceberg v3, Spark Declarative Pipelines, Real-Time Mode, and Arrow-native Python UDFs.

What changes with AWS Glue 6.0

AWS Glue 6.0 runs Apache Spark 4.1, which introduces several behavioral changes from the Spark 3.5 runtime used in AWS Glue 5.1:

Behavior Spark 3.5 (AWS Glue 5.1) Spark 4.1 (AWS Glue 6.0)
ANSI SQL mode Disabled by default Enabled by default
Legacy Parquet datetime configs Supported Removed (renamed)
Python runtime 3.11 3.13

Beyond version compatibility, AWS Glue 6.0 also introduces:

  • Apache Iceberg v3 with VARIANT Shredding for efficient semi-structured data handling.
  • Spark Declarative Pipelines — agent-authorable ETL.
  • Real-Time Mode — single-digit millisecond streaming latency.
  • Arrow-native Python UDFs (PyArrow) for improved performance.
  • Built-in observability with structured metrics and enhanced Spark UI.
  • Up to 36% better price performance compared to AWS Glue 5.1*.

These runtime changes mean your existing AWS Glue jobs might encounter removed configuration keys, stricter type casting behavior, or Python package version incompatibilities when running on AWS Glue 6.0. Fixing these manually is time-consuming and error-prone. The following sections show how the generative upgrade analysis handles this automatically.

The sample job

Our example is a daily ecommerce order analytics pipeline running on AWS Glue 5.1:

What the job does:

  • Ingests 10,000 orders from Parquet files with INT96 timestamps (including pre-1900 historical dates from a legacy system migration).
  • Computes revenue metrics by casting string prices to numeric values and calculating line totals with discounts and tax.
  • Segments customers using recency, frequency, and monetary (RFM) scoring through mapInPandas with pandas and scikit-learn.
  • Writes enriched results back to Amazon Simple Storage Service (Amazon S3).

Job configuration (AWS Glue 5.1):

Glue version: 5.1
Worker type: G.1X
Workers: 10
Python modules: pandas==2.2.2, scikit-learn==1.5.0, numpy==1.26.4
Spark configs:
  spark.sql.legacy.parquet.datetimeRebaseModeInWrite=LEGACY
  spark.sql.legacy.parquet.int96RebaseModeInWrite=LEGACY
  spark.sql.parquet.datetimeRebaseModeInRead=LEGACY
  spark.sql.parquet.int96RebaseModeInRead=LEGACY

This job runs successfully on AWS Glue 5.1. The following sections walk through how the upgrade analysis identifies and resolves incompatibilities when upgrading this job to AWS Glue 6.0. Before starting, confirm you have the prerequisites in place.

Prerequisites

  • An AWS account with access to the AWS Glue console.
  • An existing AWS Glue job on version 5.1 or earlier with at least one successful run.
  • An Amazon S3 path for storing the upgrade analysis results.

Running the upgrade analysis from the console

The following steps walk through the upgrade analysis workflow using the AWS Glue console.

Step 1: Select your job

Navigate to your job in the AWS Glue Studio console. Confirm the job has a successful run history on AWS Glue 5.1 before starting the upgrade analysis.

AWS Glue Studio job run history showing a successful run on AWS Glue 5.1

Figure 1: Job run status for the job on AWS Glue 5.1

Step 2: Start the upgrade analysis

From the job’s Actions menu, select Upgrade with generative AI. Configure the following:

  • Target AWS Glue version: 6.0.
  • Results S3 path: An S3 location where the analysis stores its artifacts and recommendations.
Actions menu in AWS Glue Studio with the Upgrade with generative AI option

Figure 2: The Upgrade with generative AI option in the Actions menu

Configure the target AWS Glue version and the S3 results path, then choose Run.

Upgrade window with the target AWS Glue version set to 6.0 and an Amazon S3 results path

Figure 3: The Upgrade with generative AI window for setting the target AWS Glue version and results path

Choose Run. The analysis begins by running your job on AWS Glue 5.1 to establish a baseline. It then iteratively tests the job on AWS Glue 6.0, identifies failures, applies recommended fixes, and validates the job. If the upgrade analysis cannot resolve an incompatibility within its attempt budget, the analysis stops and reports the unresolved issue for manual review. Your original job remains unchanged.

Note: The upgrade analysis executes your job multiple times (one baseline run plus one or more validation attempts), and each run consumes Data Processing Units (DPUs). For large or long-running jobs, consider using the run configuration option to specify fewer workers or a smaller dataset to optimize analysis cost.

Step 3: Monitor progress

The console displays the analysis progressing through multiple validation attempts. Each attempt either succeeds or fails with a specific error, and the upgrade analysis uses that error signal to determine and apply the appropriate fix for the next attempt.

Upgrade analysis progress showing multiple validation attempts with success and failure states

Figure 4: Upgrade analysis progress across multiple validation attempts

What the upgrade analysis found and fixed

The analysis completed in four validation attempts, identifying and resolving three distinct incompatibilities. The upgrade uses deterministic migration rules for known config changes, and automated diagnosis for runtime or code errors.

Iteration 1: Removed Parquet legacy configuration

The analysis first sanitizes any Spark configurations that were removed in Spark 4.1. Our job used spark.sql.legacy.parquet.datetimeRebaseModeInWrite and spark.sql.legacy.parquet.int96RebaseModeInWrite, which no longer exist.

Migration rule applied: The SQL configs with the spark.sql.legacy prefix were removed in Spark 4.1. They have been renamed to their non-legacy equivalents, preserving the original values.

Recommended change:

Before:
spark.sql.legacy.parquet.datetimeRebaseModeInWrite=LEGACY
spark.sql.legacy.parquet.int96RebaseModeInWrite=LEGACY

After:
spark.sql.parquet.datetimeRebaseModeInWrite=LEGACY
spark.sql.parquet.int96RebaseModeInWrite=LEGACY

The read-side configs (datetimeRebaseModeInRead, int96RebaseModeInRead) already used the correct non-legacy names and required no changes.

However, with this fix applied, the validation run still failed because the Python module installation encountered an error on the AWS Glue 6.0 image.

Iteration 2: Python module version incompatibility

The pinned module versions (pandas==2.2.2, scikit-learn==1.5.0, numpy==1.26.4) could not be installed in the AWS Glue 6.0 Python 3.13 environment.

Error:

LAUNCH ERROR | Installation of Additional Python Modules failed

Recommended change: The upgrade analysis updated the version specifications from exact pins to minimum version constraints, allowing pip to resolve compatible versions for Python 3.13:

Before: pandas==2.2.2, scikit-learn==1.5.0, numpy==1.26.4
After:  pandas>=2.1.0, scikit-learn>=1.3.0, numpy>=1.24.0

With modules installing successfully, the job launched on AWS Glue 6.0 but encountered a runtime error.

Iteration 3: ANSI mode strict type casting

Spark 4.1 enables ANSI SQL mode by default (spark.sql.ansi.enabled=true). Our revenue calculation casts string prices to double, but approximately 1.8% of records contain non-numeric placeholder values such as “N/A”, “pending”, or “null” from the upstream system. From a business perspective, this meant 1.8% of revenue orders were silently excluded from revenue metrics. This data quality issue was invisible to the original pipeline.

On AWS Glue 5.1 (ANSI mode off), cast("N/A" as double) silently returns null. On AWS Glue 6.0 (ANSI mode on), this throws an exception:

Error:

NumberFormatException: [CAST_INVALID_INPUT] The value 'null' of the type
"STRING" cannot be cast to "DOUBLE" because it is malformed. Correct the
value as per the syntax, or change its target type. Use try_cast to
tolerate malformed input and return NULL instead. SQLSTATE: 22018

Migration rule applied: As of Spark 4.1, spark.sql.ansi.enabled is on by default. Casting a malformed value now raises CAST_INVALID_INPUT instead of returning NULL. The upgrade analysis resolved this by updating the script to use try_cast(), which safely returns NULL for malformed input while preserving ANSI mode protections for the rest of the job.

Recommended change:

Before (AWS Glue 5.1):

F.col("unit_price").cast("double")

After (AWS Glue 6.0, fixed by the upgrade analysis):

F.expr("try_cast(unit_price as double)")

This is a targeted fix that handles the known dirty data without disabling ANSI mode globally, keeping overflow detection and type safety active throughout the job.

Final validation and data quality check

After applying all three fixes, the analysis ran the job on AWS Glue 6.0 one final time and performed a data quality comparison between the AWS Glue 5.1 baseline output and the AWS Glue 6.0 output.

Result: The job completed successfully and all data validations passed with no mismatches detected between the source and target outputs.

Completed upgrade analysis status with links to the results output path in Amazon S3

Figure 5: Final analysis status with links to the results output path in Amazon S3

Reviewing the upgrade summary

The analysis produces a detailed summary stored in your S3 results path. This summary documents each validation attempt, the errors encountered, the migration rules applied, and the recommended configuration changes:

s3://amzn-s3-demo-bucket/scripts/auto-upgrade/ja-{analysis-id}/
    summary/
        summary.md                       # Full iteration-by-iteration report
        data_validation_summary.md       # Data quality comparison results
    artifact/
        attempt_N/
            script/main.py               # Recommended script (if modified)
            job_config_modifications.json  # Recommended parameter changes
            requirements.txt             # Updated dependency versions

The following is a snippet from the upgrade summary (summary.md) showing the recommended changes and validation attempt details:

The summary documents each validation attempt, the changes applied, and the data quality results.

Upgrade summary showing validation attempts, applied changes, and data quality comparison results

Figure 6: Upgrade summary snippet showing validation attempt details, data quality, and analysis results

After reviewing the recommendations, accept the changes to upgrade your job to AWS Glue 6.0. This updates your job definition with the recommended configuration, including the renamed Spark configs, updated module versions, and any script modifications. Because the analysis has already validated the job on AWS Glue 6.0 and confirmed data quality parity with the original, your job is ready for production.

After reviewing the recommendations, you can apply the upgraded script to your job.

AWS Glue Studio prompt to apply the upgraded script to the job

Figure 7: The option to apply the upgraded script to the job

Choose Apply to confirm the upgrade.

Confirmation dialog with the Apply button to upgrade the job to AWS Glue 6.0

Figure 8: The Apply button that confirms upgrading the job to AWS Glue 6.0

After applying, the job definition reflects the new AWS Glue version.

AWS Glue job details showing version 6.0 after applying the upgrade

Figure 9: The AWS Glue version for the job after applying the upgrade

Python virtual environments in AWS Glue 6.0

AWS Glue 6.0 introduces --python-virtual-env-storage-prefix, a service-managed virtual environment with S3 caching that simplifies Python dependency management.

For existing jobs that use --additional-python-modules, no action is required. AWS Glue automatically handles the conversion to virtual environments when your job runs on AWS Glue 6.0. Your jobs continue to work without any changes.

For new jobs on AWS Glue 6.0, we recommend using the virtual environment approach:

{
    "DefaultArguments": {
        "--python-virtual-env-storage-prefix": "s3://amzn-s3-demo-bucket/glue-venv-cache/",
        "--additional-python-modules": "pandas>=2.1.0,scikit-learn>=1.3.0,numpy>=1.24.0"
    }
}

How it works:

  • On the first run, AWS Glue installs your modules into a virtual environment, packages it, and caches the result to your specified S3 path (approximately 15–30 seconds of additional startup time).
  • On subsequent runs, AWS Glue downloads and extracts the cached virtual environment instead of running pip install.
  • The cache is automatically invalidated when your module list, versions, or AWS Glue version changes.

This approach provides faster cold starts after the first run, requires no Docker image management (unlike --python-virtual-env), and is entirely service-managed with no maintenance burden.

Conclusion

The generative upgrade analysis identified and resolved three distinct compatibility issues in our AWS Glue 5.1 job, so the job now runs successfully on AWS Glue 6.0 with Apache Spark 4.1:

  • The upgrade analysis renamed legacy Parquet datetime configuration keys (removed in Spark 4.1) to their current equivalents.
  • The upgrade analysis updated Python module version specifications that were incompatible with Python 3.13 to use flexible minimum version constraints.
  • The upgrade analysis addressed the new ANSI SQL mode default (which causes runtime failures on malformed data) with a targeted fix using try_cast() to safely handle non-numeric values while preserving ANSI mode protections.

The analysis validated that the upgraded job produces output consistent with the original, and presented all changes as recommendations for review before applying them to your job.

Next steps

After you have reviewed and accepted the upgrade changes, you can delete the analysis results stored in your S3 results path.

*Based on 3TB TPC-DS benchmark comparing AWS Glue 6.0 to AWS Glue 5.1.


About the authors

Prasad Nadig

Prasad Nadig

Prasad is a Senior Analytics Specialist Solutions Architect at Amazon Web Services (AWS), specializing in large-scale data analytics and AI. Prasad partners with customers to design, migrate, and modernize their analytics platforms on AWS into scalable, cost-effective solutions, with deep expertise in data lakes, data warehousing, distributed processing, and performance tuning at petabyte scale.

Shrey Malpani

Shrey Malpani

Shrey is a Senior Product Manager Technical at Amazon Web Services (AWS), where he works at the intersection of distributed data processing and data integration. He is focused on building and scaling data integration and data management capabilities across services like AWS Glue, Amazon EMR, and Amazon Redshift that help customers build AI-ready data platforms for their analytics and machine learning workflows.

Rishabh Nair

Rishabh Nair

Rishabh is a Software Development Engineer in the AWS analytics organization, where he combines generative AI with distributed systems to build agentic workflows that modernize large-scale data processing. He is passionate about the infrastructure that makes these workflows reliable and scalable for customers.

Keerthi Chadalavada

Keerthi Chadalavada

Keerthi is a Senior Software Development Engineer in the AWS analytics organization. She focuses on combining generative AI and data integration technologies to design and build comprehensive solutions for analytics and data engineering workloads.

AWS Weekly Roundup: Student Rewards on AWS Builder Center, Local Zone in Las Vegas, and more (August 24, 2026)

Post Syndicated from Esra Kayabali original https://aws.amazon.com/blogs/aws/aws-weekly-roundup-student-rewards-on-aws-builder-center-local-zone-in-las-vegas-and-more-august-24-2026/

During my time at AWS, I have always looked for opportunities to work with students. I have delivered over 50 talks at universities across the region, and watching the potential in the room is always a strong motivator. It reminds me of why I do this work, and that the students I meet today may well become our customers and collaborators tomorrow. That is why I am happy to open this week with Student Rewards on AWS Builder Center.

Rick Suttles published Introducing Student Rewards on AWS Builder Center, a new benefit for verified higher education students. When you verify your enrollment through SheerID and complete your Builder Center profile, you unlock 12 months of premium AWS Skill Builder access (900+ courses, hands-on labs, certification exam prep, and game-based learning). From there, you earn badges through actions on Builder Center: publishing articles, commenting, and maintaining engagement. At 7 badges, you unlock $10 in AWS Credits. At 14 badges, another $20 in credits. At 21 badges, you earn an AWS Foundational Certification exam voucher ($100 value).

This represents a commitment of over $500 million in resources during this back-to-school season, providing students with the training, tools, and certification needed to start building their careers in cloud and AI. Student Rewards is available to students 18 years or older and enrolled at accredited higher education institutions worldwide, subject to verification and applicable terms.

Verify your student status and start learning, earning badges, and unlocking rewards!

Last week’s launches
Here’s what else happened this week.

  • A new AWS Local Zone in Las Vegas, Nevada – This new Local Zone supports Amazon EC2 C7i, M7i, R7i, and C8gn instances, Amazon EBS, Amazon ECS, Amazon EKS, Application Load Balancer, and AWS Direct Connect. AWS Local Zones are now available in more than 30 metropolitan areas worldwide. In addition, AWS added a fourth Availability Zone to the Europe (London) Region, delivering next-generation AI and ML capacity with Trn3 and P6 accelerated instances alongside general-purpose compute.
  • Amazon EC2 Auto Scaling now supports batch instance termination – You can now pass up to 100 instance IDs to the TerminateInstanceInAutoScalingGroup API to terminate them as a batch, reducing the number of API calls needed to scale down your Auto Scaling groups. Batch termination is designed for workloads that need to rapidly scale down, such as AI/ML training jobs, container orchestrators, or event-driven architectures that spin up large fleets temporarily.
  • AWS CloudShell now includes a built-in visual file editor – CloudShell now includes a visual file editor that you can launch directly from your shell session using a single edit command. The editor supports syntax highlighting, find-and-replace, multi-line selection, copy-paste, and undo-redo in a single browser session. Whether you are updating a deployment script, modifying an agent steering file, editing a CloudFormation template, or fixing a Lambda function, the editor provides a seamless edit-and-run experience without leaving CloudShell.
  • Amazon Bedrock now supports SpaceXAI Grok 4.6 with cross-Region inference – Grok 4.6, a frontier model built for coding, agentic tasks, and knowledge work, is now available on Amazon Bedrock. The model runs on the bedrock-runtime endpoint with support for the Responses, Chat Completions, and Converse APIs, and works with existing account-level controls including model invocation logging, Amazon CloudWatch metrics, and cost itemization in AWS Cost Explorer.
  • Amazon Bedrock expands API support and introduces cross-Region inference for OpenAI models – Amazon Bedrock now supports OpenAI GPT-5.6 models (Sol, Terra, and Luna) with the Responses, Converse, and Chat Completions APIs, and adds cross-Region inference. Geo cross-Region inference routes requests within a predefined geography (including new US Geo support with this launch), while Global cross-Region inference serves requests from any commercial AWS Region at a lower per-token cost.
  • AgentCore payments is now generally available in Amazon Bedrock AgentCore – At general availability, AgentCore payments includes Quick Create for Coinbase credential provisioning directly within the AgentCore console, a curated Coinbase Bazar MCP server of pay-per-use x402 endpoints via AgentCore gateway, support for the Machine Payment Protocol (MPP), and the “upto” scheme in the x402 protocol for pay-per-inference and dynamic pricing use cases. To learn more, visit the AI Blog post.
  • AWS Glue 6.0 delivers 30% price reduction and Iceberg v3 support – AWS Glue 6.0 is built on a fully modernized runtime, Apache Spark 4.1, Python 3.13, and Scala 2.13, delivering 30% lower pricing than previous AWS Glue versions. With Iceberg v3, Glue 6.0 adds the VARIANT data type with automatic shredding for faster reads on semi-structured data, deletion vectors for high-performance row-level updates, geometry and geography data types for spatial processing, and flexible schema evolution.

For a full list of AWS announcements, be sure to keep an eye on the What’s New with AWS page.

Other AWS news
Here are some additional posts you may find useful:

  • Updates to your AWS Sign-In experience – AWS is gradually introducing updates to the sign-in and sign-up experience. The redesigned sign-in page introduces a unified email entry point for root users and customers using the new email-based sign-in method, while IAM users continue signing in with their account ID, username, and password. The page also includes sign-in options for customers whose AWS account was created using a supported identity provider (Google, GitHub, Apple, or Amazon.com). A redesigned session selection page simplifies viewing and managing multiple active account and role sessions. If your organization relies on browser automation or scripted workflows that interact with the sign-in page, review the post to understand how these changes might affect your configuration.
  • In the works: AWS Builder Lofts in Berlin, Hyderabad, and São Paulo – My colleague Channy announced plans to open new Builder Lofts in three cities. Since the first Builder Loft opened in San Francisco in July 2025, it has welcomed more than 22,500 developers through its doors. Each new location will be a permanent community space offering free workshops, networking events, pitch nights, content creation spaces, and co-working areas. Berlin will focus on digital sovereignty and security-readiness, Hyderabad on AI and cloud-native architecture, and São Paulo on supporting Latin America’s developer ecosystem.
  • AWS and Amazon WorkSpaces recognized as a Leader in the 2026 Gartner Magic Quadrant for Desktop as a Service – AWS has been named a Leader in the 2026 Gartner Magic Quadrant for Desktop as a Service (DaaS) for the third consecutive year, evaluated on Completeness of Vision and Ability to Execute. Gartner noted strengths in operations, geographic strategy, and overall viability. This is also the first year the evaluation includes Amazon WorkSpaces for AI agents, a capability that runs AI agents within the same desktop environment, security perimeter, and audit trail as human users.

For a full list of AWS blog posts, be sure to keep an eye on the AWS Blogs page.

Upcoming AWS events
Check your calendar and sign up for upcoming AWS events:

Visit the AWS Builder Center to meet other builders, contribute solutions, and find resources that help you keep building.

Summer is slowly coming to an end, and I am already planning a few days off in the coming months to keep me motivated through the rainy autumn ahead. I hope you are doing the same. Come back next week for more!

— Esra

AWS Glue 6.0 now available with 30% lower price and full Apache Iceberg v3 support

Post Syndicated from Channy Yun (윤석찬) original https://aws.amazon.com/blogs/aws/aws-glue-6-0-now-available-with-30-lower-price-and-full-apache-iceberg-v3-support/

Today, we are announcing the general availability of AWS Glue 6.0, delivering 30% lower pricing than previous AWS Glue versions and introducing full support for Apache Iceberg v3 features. AWS Glue 6.0 is built on a fully modernized runtime, Apache Spark 4.1, Python 3.12, and Scala 2.13, delivering faster performance.

With this release, AWS Glue provides the most complete Iceberg v3 implementation on any fully serverless managed Spark service, along with new capabilities that simplify ETL authoring, improve PySpark performance, and enable real-time streaming with single-digit millisecond latency.

What is new in AWS Glue 6.0
AWS Glue 6.0 delivers the complete Apache Iceberg v3 specification, built on Iceberg 1.11.0. The headline feature is the VARIANT data type with shredding support, which achieves faster query read performance compared to traditional string data type columns for semi-structured data.

With VARIANT shredding, you can store and query JSON, logs, and event data without flattening schemas, eliminating duplicate data copies, custom parsing code, and pipeline breakage when schemas change. This capability transforms how teams handle semi-structured data at scale.

Additional Iceberg v3 capabilities include:

  • Geometry and Geography data types: Enable native spatial processing for GIS analytics, location intelligence, and geospatial data pipelines directly on managed Spark.
  • Nanosecond-precision timestamps: Support IoT sensor data, scientific computing, and high-frequency financial workloads that require precision beyond standard milliseconds.
  • Unknown type handling: Process data with unexpected or evolving schemas without pipeline failures, providing resilience against upstream schema changes.

AWS Glue 6.0 also includes most significant upgrade in Spark 4.1, the modern runtime engine:

  • Spark declarative pipelines: Spark Declarative Pipelines introduces a simplified approach to ETL authoring. Data engineers declare transformations, specifying what data should look like, while the engine automatically determines execution order and optimization. This reduces the complexity of pipeline development and eliminates manual orchestration overhead.
  • Arrow-native Python UDFs and UDTFs: AWS Glue 6.0 introduces Arrow-native execution for Python User-Defined Functions (UDFs) and User-Defined Table Functions (UDTFs). This eliminates serialization overhead between Python and the JVM, improving PySpark performance for complex transformations.
  • Real-time streaming mode: For stateless streaming use cases, AWS Glue 6.0 introduces a real-time streaming mode that achieves single-digit millisecond latency. Built on Spark 4.1’s Real-Time Mode with Glue-optimized execution, this capability supports real-time event processing, low-latency data transformation pipelines, and time-sensitive data routing.

Getting started with AWS Glue 6.0
No API changes are required to use AWS Glue 6.0. You can select the new version using the existing --glue-version parameter in the create-job or update-job APIs through AWS Command Line Interface (AWS CLI), AWS SDK, AWS Glue Studio, Amazon SageMaker Unified Studio, and your preferred IDE.

To get started with AWS Glue 6.0 jobs in the AWS Glue Studio console, open the AWS Glue job and on the Job Details tab, choose the version Glue 6.0 – Supports Spark 4.1, Scala 2, Python 3. You can create new AWS Glue jobs on AWS Glue 6.0 to get the benefit from the improvements, or migrate your existing AWS Glue jobs.

To start using AWS Glue 6.0 on an AWS Glue Studio notebook or an interactive session through a Jupyter notebook, set 6.0 in the %glue_version magic. You can also upgrade existing jobs to Glue 6.0 using the Spark upgrade agent on AWS Glue Studio or use the auto-upgrade feature in their existing Glue jobs to automatically upgrade them to Glue 6.0.

To learn more, visit the AWS Glue 6.0 version detail and Migrating AWS Glue for Spark jobs to AWS Glue version 6.0 in the AWS documentation.

Now available
AWS Glue 6.0 is generally available today in all AWS Regions where AWS Glue operates. For Regional availability and a future roadmap, visit the AWS Capabilities by Region. If you want to call APIs, search documentation, find regional availability, and check troubleshooting about this new feature, try using the AWS MCP Server and plugins with your preferred AI tool.

You pay an hourly rate, billed by the second, for crawlers (discovering data) and extract, transform, and load (ETL) jobs (processing and loading data). For the AWS Glue Data Catalog, you pay a simplified monthly fee for storing and accessing the metadata. The first million objects stored are free, and the first million accesses are free. To learn more, visit AWS Glue Pricing page.

Give it a try in the AWS Glue Studio console, and send feedback to AWS re:Post for AWS Glue or through your usual AWS support contacts.

— Channy

Automate creating AWS Glue Data Catalog views with AWS SDK for data mesh use case

Post Syndicated from Aarthi Srinivasan original https://aws.amazon.com/blogs/big-data/automate-creating-aws-glue-data-catalog-views-with-aws-sdk-for-data-mesh-use-case/

AWS Glue Data Catalog view is a multi-dialect view that supports querying from multiple SQL query engines, such as Amazon Athena, Amazon Redshift Spectrum, Apache Spark in Amazon EMR and AWS Glue. You can create a Data Catalog view in one account, using an AWS Identity and Access Management (IAM) definer role in the same or different account and use AWS Lake Formation to share the view across multiple accounts. The definer role has the required full SELECT on the base tables to create the view and share it with other users for querying. The Data Catalog assumes the definer role and manages access of the base tables when the view is queried, thus allowing to share a subset of data without sharing the underlying base tables.

AWS Glue now adds AWS SDK support for creating and updating the ATHENA dialect of Glue views. With this addition, you can now create ATHENA and SPARK dialects of Glue views simultaneously, using a cross account IAM definer role. This feature enhances the automation to create and update Glue views, like that of Data Catalog tables. In our earlier blog Create AWS Glue Data Catalog views using cross-account definer roles, we had introduced IAM definer roles in a cross-account use case to create Data Catalog views with SPARK dialects using the APIs – CreateTable() and UpdateTable() – while creating and adding ATHENA dialects using Athena query editor. As a continuation to it, this post shows you how to use the Catalog objects API CreateTable() to programmatically create ATHENA and SPARK dialects using cross-account IAM definer roles, and how to add the ATHENA dialect programmatically for the views that were created earlier with only SPARK dialect.

Cross account definer roles enable enterprise data mesh architectures where multiple accounts are interconnected in a central governance and multiple producers and consumers. The central governance account hosts the database, tables and permissions, while the producer accounts maintain CI/CD pipelines to create and manage those data assets. Having the definer role in producer accounts allows those CI/CD pipelines to be fully managed by IAM roles in the individual accounts.

Key points on creating multi-dialect views using cross-account definer roles

  • ATHENA dialects are validated and asynchronously created. Hence, a cross-account Glue connection is required for validation for every producer account-central governance account pair. This is a one-time setup.
  • SPARK dialects are not validated. Hence SPARK dialect’s create syntax requires SubObjects list of the base tables and StorageDescriptor fields for the columns of the view.
  • Though queries on cross account views can be run using database resource link names, the view definition SQL query for creating the view requires the original database and base table names from the central governance account.
  • If a view has SPARK and ATHENA dialects available, we recommend updating both the dialects of the view simultaneously using update_table() API/SDK, for any changes in the SQL definition of the view or the base table. This will keep both the dialects queryable.
  • Creating and updating both SPARK and ATHENA dialects using cross account definer role is supported using AWS CloudFormation.
  • The Data Catalog view that can be created using cross account IAM definer roles are available in SPARK and ATHENA dialects and currently not supported for Redshift Spectrum dialect.

Prerequisites

We use the same setup used in Create AWS Glue Data Catalog views using cross-account definer roles for the sample database, tables, definer role, resource link, IAM and Lake Formation permissions on those resources and principals between the two AWS accounts. Summarizing the requirements as below.

  • The setup includes a central governance account with Data Catalog database bankdata_icebergdb and two tables transaction_table1 and transaction_table2, a producer account with a Data-Analyst role used as view definer role.
  • Lake Formation permissions on the central account’s database and tables are granted to the producer account Data-Analyst role as per the earlier blog. The definer role in producer account should have database DESCRIBE and CREATE_TABLE permissions, table SELECT and DESCRIBE permission on all columns and rows of the base tables. The IAM permissions required on the definer role are detailed in Prerequisites for creating views. Similarly, follow the earlier blog to create resource link for the shared database and grant Lake Formation permissions on the resource link to the Data-Analyst
  • An Athena data source named centraladmin in the producer account, pointing to the Data Catalog of the central governance account.

Creating ATHENA and SPARK dialects at the same time

Creating both ATHENA and SPARK dialects of a Glue catalog view simultaneously is now supported by the AWS SDK. In the producer account, create a new Glue connection, required for the Athena dialect validation. This is a prerequisite for creating the ATHENA dialect of the Glue catalog view using cross account definer role. Then we create a Glue view with both dialects.

  1. Sign in to the producer account as the Lake Formation admin role, or any role with permission to create AWS Glue connections.
  2. Using an AWS Command Line Interface (AWS CLI) environment, such as AWS CloudShell, create an AWS Glue connection as follows.
    aws glue create-connection --cli-input-json file://athena-validation-connection.json

    The content of athena-validation-connection.json is as follows.

    {
        "CatalogId": "<producer-account-id>",
        "ConnectionInput": {
            "Name": "glue-view-validation-connection",
            "Description": "Glue view Athena cross-account validation connection",
            "ConnectionType": "VIEW_VALIDATION_ATHENA",
            "ConnectionProperties": {
                "WORKGROUP_NAME": "primary",
                "DATA_SOURCE": "centraladmin"
            }
        }
    }

    Note: If you are using Athena for the first time in your account or using Primary workgroup, setup the query results location bucket using Specify a query result location.

  3. Sign out as the Lake Formation admin and sign back in to the producer account as the definer IAM role, Data-Analyst.
  4. Create an AWS Glue view using the create-table CLI command and JSON file, or using the AWS SDK for Python (Boto3) script.
    aws glue create-table --cli-input-json file://create_multipledialects.json

    The content of create_multipledialects.json is as follows.

     {
       "DatabaseName": "rl_bank_iceberg",
       "TableInput": {
         "Name": "view_2dialects_2basetables_fromcli",
         "StorageDescriptor": {
           "Columns": [
             {
               "Name": "transaction_id",
               "Type": "string"
             },
             {
               "Name": "transaction_type",
               "Type": "string"
             },
             {
               "Name": "transaction_amount",
               "Type": "double"
             },
             {
               "Name": "transaction_location",
               "Type": "string"
             },
             {
               "Name": "transaction_date",
               "Type": "date"
             }
         },
         "ViewDefinition": {
           "SubObjects": [
             "arn:aws:glue:us-west-2:<central-account-id>:table/bankdata_icebergdb/transaction_table1",
             "arn:aws:glue:us-west-2:<central-account-id>:table/bankdata_icebergdb/transaction_table2"
            ],
           "IsProtected": true,
           "Representations": [
             {
               "Dialect": "SPARK",
               "DialectVersion": "1.0",
               "ViewOriginalText": "SELECT a.transaction_id, a.transaction_type, a.transaction_amount, b.transaction_location, b.transaction_date FROM bankdata_icebergdb.transaction_table1 a RIGHT JOIN bankdata_icebergdb.transaction_table2 b ON a.transaction_id = b.transaction_id",
               "ViewExpandedText": "SELECT a.transaction_id, a.transaction_type, a.transaction_amount, b.transaction_location, b.transaction_date FROM bankdata_icebergdb.transaction_table1 a RIGHT JOIN bankdata_icebergdb.transaction_table2 b ON a.transaction_id = b.transaction_id"
             },
             {
                "Dialect": "ATHENA",
                "DialectVersion": "3",
                "ViewOriginalText": "SELECT a.transaction_id, a.transaction_type, a.transaction_amount, b.transaction_location, b.transaction_date FROM bankdata_icebergdb.transaction_table1 a RIGHT JOIN bankdata_icebergdb.transaction_table2 b ON a.transaction_id = b.transaction_id",
                "ValidationConnection": "glue-view-validation-connection"
             }
           ]
         }
       }
    }

    Notes about fields in the above CLI input JSON (applies to all SDK):

    • The definer is by default the API caller, but a Definer field can be set to explicitly specify a different IAM role.
    • In the ViewDefinition, database qualifiers are required for SPARK dialect. That is, the SQL definition provided for ViewOriginalText and ViewExpandedText should be in <source_database_name>.<source_table_name> format.
  5. After the view is created, you can inspect the details on the Lake Formation console. The SQL definitions show both ATHENA and SPARK as shown in the following screenshot.

Lake Formation console showing the SQL definitions tab for the new Data Catalog view, with both ATHENA and SPARK dialects listed

If your view creation fails for any of the dialects, you can use the AWS Glue get-table CLI command with --include-status-details to see what the error is and rectify it.

aws glue get-table --database-name <rl_database_name> --name <view_name> --include-status-details

Glue PySpark script

The PySpark script for creating a view with ATHENA and SPARK dialects are provided below. Download and edit the Pyspark script with your bucket name, producer and central account ids, region and relevant Glue resource names: bdb_5773_createview_bothdialects.py

Provide the following settings to run the script in your Glue Studio. For details on running a Spark job in Glue, refer Working with Spark jobs in AWS Glue.

  • Choose Data-Analyst as the job execution IAM role.
  • Choose Glue 5.1 for Glue version.
  • For the Requested number of workers, provide >=4. This is an FGAC Spark driver requirement, which is needed for Glue catalog views. Below screenshot shows these settings.
  • Add the following 2 properties as additional job parameters. A screenshot is shown for reference.
    --datalake-formats = iceberg
    --enable-lakeformation-fine-grained-access=true

    AWS Glue ETL job configuration page showing the additional job parameters set for the multi-dialect view creation script

  • Save and run the Glue job. Check the stdout logs to review the query on the newly created view.

A sample update_table script is also provided below, to illustrate changing the view definition with additional columns. Note the REPLACE keyword:

bdb_5773_updateview_bothdialects.py

Adding ATHENA dialect using SDK to an existing AWS Glue view

You can update an existing AWS Glue view that was created with the SPARK dialect and add the ATHENA dialect using the SDK. The following example uses the update-table CLI command.

aws glue update-table --cli-input-json file://add-athena-dialect.json

The content of add-athena-dialect.json is as follows.

{
    "DatabaseName": "rl_bank_iceberg",
    "ViewUpdateAction": "ADD",
    "TableInput": {
        "Name": "view_sparkfirst_athenanext",
        "ViewDefinition": {
            "Representations": [
                {
                    "Dialect": "ATHENA",
                    "DialectVersion": "3",
                    "ViewOriginalText": "SELECT a.transaction_id, a.transaction_type, a.transaction_amount, b.transaction_location, b.transaction_date FROM bankdata_icebergdb.transaction_table1 a RIGHT JOIN bankdata_icebergdb.transaction_table2 b ON a.transaction_id = b.transaction_id",
                    "ValidationConnection": "glue-view-validation-connection"
                }
            ]
        }
    }
}

Verify the added dialect on the view by reviewing the SQL definitions of the view in Lake Formation console or using GetTable(). If you want to edit the SQL definition or change the base tables of an existing view that has both SPARK and ATHENA dialects, you can do so using the update_table API (using SDK or CLI), with "ViewUpdateAction": “REPLACE” and provide both the dialect definition under ViewDefinition.

You can run queries on the view from the producer account as Data-Analyst. The view can be shared using Lake Formation Tags or named method, just like sharing tables, to additional consumer accounts from the central governance account. The consumer accounts will create a resource link and query the views.

Cleanup

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

  1. Revoke the Lake Formation permissions granted to the Data-Analyst role and the producer account from the central governance account.
  2. Drop the Data Catalog tables, views, and the database.
  3. Delete the Athena query results from your Amazon Simple Storage Service (Amazon S3) bucket.
  4. Delete the Data-Analyst role from IAM.
  5. Delete the AWS Glue connection and the Athena data source.
  6. Delete the AWS Glue job, if you tried the Python script as an AWS Glue job.

Conclusion

In this post, I demonstrated how to use cross-account IAM definer roles with AWS Glue Data Catalog views, how to create and update ATHENA and SPARK dialects using the Data Catalog CreateTable() and UpdateTable() APIs. The multi-dialect Data Catalog views allow sharing a subset of data from different tables using Lake Formation permissions, including LF-Tags based access control. The cross-account definer roles support multi-account data mesh architectures so that the producer IAM roles can run the CI/CD pipelines in its account. We encourage you to try the feature and share your feedback in the comments.

Acknowledgements: I would like to thank all the team members who worked to add AWS SDK support for creating ATHENA and SPARK dialects together for AWS Glue views – Daniil Arushanov, Wyatt Hawes, Yuxi Wu, Santhosh Padmanabhan and Karthik Devaraj.


About the author

Aarthi Srinivasan

Aarthi Srinivasan

Aarthi is a Senior Big Data Architect working on data, analytics and GenAI topics with the worldwide Specialists Org at AWS. She works with AWS customers and partners to architect data lake solutions, enhance product features, and establish best practices for data governance and analytics services adoption.

How Mapfre USA modernized fraud claims with Amazon EMR Serverless

Post Syndicated from Lijan Kuniyil original https://aws.amazon.com/blogs/architecture/how-mapfre-usa-modernized-fraud-claims-with-amazon-emr-serverless/

Insurance fraud remains a significant challenge for the insurance industry because fraudulent claims can increase loss costs, reduce trust, and consume investigation capacity that could otherwise be focused on serving customers. Traditional fraud detection approaches typically rely on rules-based controls, manual investigation triggers, historical claim patterns, and structured-data-only analysis. These approaches are useful for known fraud patterns, but they can struggle to detect sophisticated fraud rings or hidden relationships across claimants, policies, vehicles, providers, addresses, and prior suspicious activities.

Mapfre USA is the number one auto and home insurer in Massachusetts, serving customers in 11 states nationwide. Our coverage includes auto, home, motorcycle, watercraft, business insurance, and more. As part of Mapfre Group, we’re a worldwide leader serving over 31.1 million customers in more than 100 countries with a team of 31,000 employees. In collaboration with AWS and Neo4j, Mapfre USA modernized its fraud prevention capabilities by combining graph-based features with machine learning (ML) models deployed on AWS. This initiative, focused initially on Massachusetts Auto insurance and later expanded to Home (HO), has delivered significant business impact, exceeding $5 million Net Present Value (NPV), with realized savings already outperforming projections.

In this post, we share how Mapfre USA designed and implemented this solution, highlight the technical architecture running on AWS specifically on the Mapfre Data Platform called Atenea, and explore lessons learned that can apply to other industries facing complex fraud challenges.

Business challenge

Fraudulent claims aren’t always isolated events. They often involve hidden networks of policyholders, vehicles, providers, and prior suspicious activities. Detecting these complex relationships requires going beyond traditional structured data analysis.

Mapfre set out with a clear goal:

  • Goal: Improve fraud detection accuracy and claims handling efficiency.
  • Key KPI: Identify fraudulent claims missed by traditional methods.
  • Approach: Develop several ML models that use both traditional structured data and graph-based features derived from claim relationships.
  • Deployment: Integrate seamlessly with Guidewire Claims, so front-line adjusters automatically receive fraud alerts with explanations.

Each flagged claim exposure generates a Guidewire activity showing the top three model drivers, helping investigators understand why the claim was identified and act quickly.

Technical solution on AWS (Atenea Data Platform)

The fraud detection platform is built on a modern data architecture on AWS, designed to scale efficiently and provide long-term governance.

At its core, the solution uses Apache Iceberg tables stored on Amazon Simple Storage Service (Amazon S3), with metadata managed through the AWS Glue Data Catalog and access governed through AWS Lake Formation as part of the Atenea lakehouse governance model. The platform feature store is implemented through feature-store-managed Iceberg tables that manage model features, predictions, and Guidewire activities. The implementation is structured across three logical layers:

  • Silver Layer – Iceberg tables that contain source data from each of the sources, used as the initial consumption point of the platform.
  • Gold Layer – Iceberg tables storing intermediate data, such as unified Guidewire activity logs, Auto features, and Home features.
  • Platinum Layer – Feature Store-managed Iceberg tables containing encoded features and model predictions, making them reusable across models and ensuring strong metadata governance.

Processing pipelines are executed on Amazon EMR Serverless, with orchestration managed by Apache Airflow operators running on Amazon Managed Workflows for Apache Airflow (Amazon MWAA). This provides elastic, cost-efficient compute for both batch processing and fast-time scoring, while keeping orchestration, monitoring, and recovery centralized.

For graph enrichment, the platform connects to Neo4j using a dedicated driver, enabling advanced network-based features like suspicious claim linkages, provider fraud ratios, and centrality metrics.

This architecture supports efficient, reliable, and transparent production execution through repeatable Airflow orchestration, environment-based continuous integration and delivery (CI/CD) promotion, centralized monitoring, failure notifications, retry mechanisms, dead-letter queue handling for Guidewire integration, and controlled secret management. At the same time, the layered lakehouse design keeps the platform flexible enough to evolve with new business needs and fraud detection use cases.

Architecture diagram showing the Mapfre USA fraud detection platform on AWS, including data ingestion, graph enrichment, model scoring, and Guidewire integration

The data sources are policy, claims, vehicles, and notes (from AS400 and Guidewire), which include structured data and derived features capturing entity relationships (graph data).

The following list describes the architecture overview:

  1. Data ingestion – Claim batch data uploaded to Amazon S3. Data gets standardized and materialized in Iceberg tables within the Silver layer.
  2. Graph enrichment – Data processed to update Neo4j graph database hosted on AWS.
  3. Model training and scoring – Batch scoring for several ML models.
  4. Model orchestration – Unified orchestration for ingestion, training, and inference using Apache Airflow operators. CI/CD pipelines for promotion across environments.
  5. Execution platform – Amazon EMR Serverless for cost-efficient Spark processing. Migration to Apache Iceberg plus AWS Glue Data Catalog for scalable metadata handling.
  6. Integration with claims systems – Fraud predictions automatically create Guidewire activities, enriched with a description for investigators.
  7. Secrets and security – AWS Secrets Manager securely stores credentials and tokens for Guidewire API integration, with environment-specific and Region-specific access controls.
  8. Monitoring and reliability – Amazon CloudWatch and Amazon Simple Notification Service (Amazon SNS) provide visibility into pipeline health and notify teams on failures. Data quality checks are executed at key stages of the pipeline to validate data availability, schema consistency, completeness, and business-rule expectations before outputs are consumed by models or sent to Guidewire.

Guidewire integration with MLOps on AWS

One of the most important parts of Mapfre’s solution was closing the loop between ML predictions and the claims handling system. This required a resilient integration between the Atenea data platform on AWS and Guidewire Claims.

The following describes the integration flow:

  1. When an ML use case finishes scoring, the results are written as JSON files into the S3 path: <bucket_name>/guidewire/.
  2. An S3 event notification triggers an AWS Lambda function.
  3. This Lambda function:
    • Reads the JSON file.
    • Calls the Guidewire Predictive Model API.
    • Because Guidewire doesn’t support batch requests, the Lambda function sends each JSON payload individually. This keeps the integration compatible with Guidewire and isolates failures at the individual activity level, but it increases the number of API calls and makes retry, throttling, DLQ handling, and monitoring controls important.
  4. If successful, the API responds with HTTP 201 (activity created).
    • If not, the Lambda function retries up to two times.
    • Failed requests are sent to an Amazon Simple Queue Service (Amazon SQS) dead-letter queue (DLQ), and an Amazon SNS notification is published for monitoring.
  5. Secrets are stored in AWS Secrets Manager and injected as Lambda environment variables, along with AWS Region-specific URLs for token retrieval and API endpoints.
  6. The following JSON shows the example structure for Guidewire integration:
{
  "method": "createPredictiveActivity",
  "params": [
    {
      "claimNumber": "AUXXXXXXX",
      "exposureNumber": 1,
      "subject": "Fraud alert from ML model",
      "description": "Claim flagged as potential fraud based on graph + ML features",
      "shortSubject": "ML_Fraud_Flag",
      "priority": "high",
      "availableForClosedClaim": true,
      "autoCloseOnExposureClosure": false,
      "targetDays": 4,
      "escalationDays": 6
    }
  ]
}

Diagram showing the Guidewire integration flow with AWS Lambda, Amazon SQS dead-letter queue, and AWS Secrets Manager

Key benefits of this integration:

  • Real-time actionability – Fraud predictions automatically create Guidewire activities for front-line adjusters.
  • Resilience – Built-in retries, DLQ handling, and Amazon SNS alerts make sure failed events aren’t lost.
  • Security – Secrets and tokens are managed using AWS Secrets Manager, with strict environment separation (dev, pre, pro).
  • Scalability – Any new MLOps use case writes results into the S3 output path, automatically flowing into Guidewire.

This integration shows that fraud models don’t exist in isolation but actively augment daily claim workflows in production. It connects Atenea’s MLOps pipelines on AWS directly with business decisioning systems, which is critical to realizing the fraud savings impact.

Data quality and resilience

For robustness, data quality checks are applied on ingestion pipelines and graph features. Automated validation detects anomalies early, monitoring dashboards track KPIs and model performance, and standardized recovery and promotion processes operate across environments.

Visualization and investigative tools

Neo4j Bloom supports SIU workflows by visually exploring entity relationships, such as a provider linked across multiple suspicious claims, accelerating fraud ring identification.

Conclusion

The fraud detection model in auto claims has enhanced Mapfre USA’s ability to identify fraudulent activity, driving significant savings and improving overall claims efficiency.

During the pilot phase alone, savings exceeded projections, and in production the initiative has proven a Net Present Value (NPV) of more than $5M. These results confirm the business case and highlight the strength of combining structured data with graph-based features to uncover fraud networks that traditional approaches miss.

The results have been compelling:

  • Accuracy gains – Detection improved by 50–135 percent compared to baseline methods.
  • Substantial realized value – Both during the pilot and in production.
  • Cross-functional success – The initiative brought together Claims, IT Data, Advanced Analytics, and Neo4j teams in an agile, collaborative model.

Beyond the financial outcomes, several lessons have emerged. First, cross-functional collaboration between groups like Claims, Data Engineering, Advanced Analytics, and technology partners like AWS and Neo4j was critical to success. Second, explainability proved essential. By presenting adjusters with the top model drivers directly in Guidewire, trust and adoption of the system increased substantially. Finally, building resilience into the architecture through monitoring, retries, and data quality processes helped the models operate reliably in production.

Looking ahead, the platform is well-positioned to expand beyond fraud detection. New use cases such as underwriting anomaly detection and customer entity resolution are already on the roadmap. With robust architecture built on AWS using Amazon EMR Serverless, Apache Iceberg on Amazon S3 supported by AWS Glue Data Catalog and Lake Formation, a custom-built Feature Store, and Neo4j, Mapfre now has a scalable foundation to continue driving innovation and business impact.

To learn more about Amazon EMR Serverless, see the Amazon EMR Serverless documentation.

Multi-cloud lakehouse architecture on AWS for Agentic AI, Part 1: Architecture and best practices

Post Syndicated from Sakti Mishra original https://aws.amazon.com/blogs/big-data/multi-cloud-lakehouse-architecture-on-aws-for-agentic-ai-part-1-architecture-and-best-practices/

Enterprise data architectures have become fundamentally distributed. Over the past decade, organizations have made deliberate investments across multiple platforms such as relational databases for transactional workloads, cloud data warehouses for analytics, object stores for unstructured data, and SaaS applications for domain-specific functions. Each was chosen to solve a specific problem, serve a specific team, or meet a specific performance requirement. The result is not accidental sprawl. It is a deeply heterogeneous data landscape shaped by intentional, workload-driven decisions. The challenge now is not consolidation, but interoperability: enabling these systems to function as a unified foundation for the next generation of AI-driven applications.

Agentic AI systems that autonomously reason, plan, and take action on behalf of users are moving rapidly from experimentation to enterprise production. These systems do not just retrieve information. They synthesize it, act on it, and learn from it. And unlike traditional analytics tools that can work with a well-scoped dataset, AI agents require something more demanding: unified, governed, and real-time access to all relevant enterprise data, regardless of where it lives.

This is the gap that matters most right now. Enterprises that have invested in building strong data capabilities across multiple providers are well-positioned, but only if those platforms can be accessed together, consistently, and with the governance controls that enterprise AI requires. Without a unified data foundation, AI agents operate with incomplete context, governance becomes inconsistent, and the promise of autonomous AI remains out of reach.

Solution approach

The following high-level architecture explains how you can onboard metadata catalogs and MCP servers to your context layer, which becomes the primary input for your AI agents.

Assuming your data products have a well-defined metadata catalog, you can take a unified-catalog-first approach, then build the context layer on top of it to let your AI agents discover all the context from one place. This helps bring in centralized governance and audit control, because every request gets routed through the centralized metadata catalog and context layer to simplify implementation of unified governance. In addition, this brings simplicity to enable business semantics, define attribute priorities, and define authoritative sources for the consumer use cases.

Architecture showing metadata catalogs and MCP servers onboarded to a context layer that feeds AI agents

If any of the data sources does not have a well-defined metadata catalog, you can define Model Context Protocol (MCP) servers on them, and then directly onboard them to the context layer. For example, if you have semi-structured or unstructured datasets for which you do not have a well-defined metadata catalog, or you want to onboard third-party data sources through REST APIs, then you can add their respective MCP server to the context layer directly. The following architecture explains the extended flow for it.

Extended architecture where data sources without a metadata catalog expose MCP servers directly to the context layer

In this series of posts, we demonstrate how you can unify the metadata catalog access across multiple providers, how you can enable AI agents to query the unified catalog, and how the context layer can be integrated to unify metadata from catalogs and MCP servers. We have divided the series into the following parts.

  • Part 1: Architecture approach with tradeoffs to unify a multi-cloud lakehouse architecture that can power Agentic AI (this post).
  • Part 2: Implementing an example solution to unify catalogs from multiple providers and deploy AI agents to query the unified data access layer.
  • Part 3: Integrate a context layer on top of the unified catalog for AI agents.
  • Part 4: Onboard additional data sources to the context layer through MCP servers and demonstrate the full solution.

This post focuses on explaining the architecture approach to build the open lakehouse architecture on AWS, unifying the metadata catalog across providers for the AI agents to access. In addition, it highlights the architecture trade-offs and best practices.

Use case

Every AI initiative launched on a fragmented data foundation is an initiative that will need to be rebuilt. Organizations that establish unified data access today are the ones that will scale Agentic AI with confidence tomorrow. Consider a large enterprise managing petabytes of data across a diverse set of environments:

  • On-premises: Network device telemetry, customer records, and operational databases.
  • Multiple cloud platforms: Marketing analytics, HR systems, and enterprise applications distributed across cloud providers.
  • Data platforms: Data science workloads, feature engineering pipelines, and finance and supply chain analytics running on specialized platforms.
  • SaaS applications: Salesforce, SAP, Zendesk, ITSM, and other business tools that each hold a critical piece of the enterprise data picture.

The business objective is to build a unified analytics and AI platform that can:

  • Query and analyze data across all environments without requiring full data migration.
  • Enforce consistent data governance and access control regardless of data location.
  • Power AI agents that can autonomously discover, query, and act on enterprise data.
  • Reduce total cost of ownership by eliminating redundant pipelines and storage.

This architecture directly addresses these needs by combining flexible data integration patterns, an open-table-format-based lakehouse architecture (with an example of Apache Iceberg), AI agent deployment to access unified metadata, and centralized governance.

Reference architecture

Before going deeper into a specific architecture, let’s revisit at a high level how the AWS open lakehouse architecture enables data ingestion and query or catalog federation to power analytics, machine learning development, and generative AI application development.

The following architecture diagram represents an end-to-end flow that includes:

  • Data ingestion to the data lake or data warehouse through Zero-ETL and batch or stream processing using AWS native services, or accessing data from Google Cloud Platform using AWS Interconnect – multicloud.
  • A centralized metadata catalog layer that includes data on AWS and metadata representation of non-AWS data sources using query or catalog federation.
  • A context layer that you can integrate to create a knowledge graph with ontology and business semantics that can enrich context for AI agents.
  • The consumption layer, which can include analytics, machine learning model development with Amazon SageMaker AI, and generative AI application development with Amazon Bedrock AgentCore, Amazon Quick, or other AWS and non-AWS AI applications.

End-to-end AWS open lakehouse architecture spanning ingestion, catalog, context, and consumption layers

Let’s look at an expanded version of this architecture that details the data ingestion and data consumption patterns to build a unified data access layer on AWS that spans multiple cloud and ISV providers.

Expanded technical architecture walkthrough

The following architecture demonstrates the comprehensive AWS approach for metadata catalog consolidation through flexible integration patterns, and it also highlights patterns for building a lakehouse on AWS. Built on the open standards of Apache Iceberg for storage and governance through AWS Lake Formation, it creates a unified data foundation that connects existing investments without requiring wholesale migration, and it makes enterprise data AI-ready from day one. This architecture delivers value at every layer: business teams query across platforms without data movement, IT teams manage governance through a single federated layer with the flexibility to federate or ingest per use case, and compliance teams enforce policies once across all sources with full lineage and audit coverage.

Expanded lakehouse architecture on AWS showing federation and ingestion patterns across multiple cloud and ISV providers

The following are the key components of the architecture.

Data access methods

This section provides options to access data that is not available in AWS Glue Data Catalog and not available on AWS.

1. Iceberg catalog federation (Reference points 2, 6.1, 6.2)

  • AWS Glue Data Catalog implements the Iceberg REST Catalog API specification, which enables seamless federation with Databricks, Snowflake, or other Iceberg-compatible catalogs set up with Amazon Simple Storage Service (Amazon S3) as the storage layer.
  • With the growing adoption of Apache Iceberg, catalog federation will become a common standard in the future and simplify metadata unification.

2. Query federation (Reference point 1.1)

  • Direct cross-cloud querying over the public internet to Google BigQuery, Azure SQL, Salesforce, and other platforms.
  • Real-time access to external data sources without replication, and seamless access with AWS analytics services.
  • Provides flexibility, because the catalog federation capability of the Iceberg REST catalog is limited to Iceberg tables only.

2.1. Secured private connectivity to Google Cloud Platform using AWS Interconnect for multi-cloud (Reference points 3.1, 3.2)

The default query federation approach makes the connection and transfers data over the public internet, which has its own latency implications depending on the target platform and the data volume transferred over the internet. During re:Invent 2025, AWS announced the public preview of AWS Interconnect – multicloud, which recently became generally available.

AWS Interconnect – multicloud is a managed service that provides private, high-speed, and secure network connections between Amazon Web Services (AWS) and other cloud providers, starting with Google Cloud Platform (GCP), with Microsoft Azure and Oracle Cloud Infrastructure (OCI) coming later in 2026. You can enable the integration with three steps: 1) specify the target cloud service provider, 2) select the destination Region on the other side, and 3) pick the required bandwidth.

The following architecture represents AWS and GCP integration with AWS Interconnect – multicloud.

High-level architecture of AWS and GCP integration through AWS Interconnect for multi-cloud

On the AWS side, you need an AWS Direct Connect gateway (a global construct that acts as a route reflector), which you can attach to your Amazon Virtual Private Cloud (Amazon VPC) through a virtual private gateway or AWS Transit Gateway, or AWS Cloud WAN. On the GCP side, you need a Google Cloud Router that you attach to your customer VPC. Interconnect – multicloud offers pre-cabled capacity pools at shared Interconnect points of presence (PoPs) in selected Regions, where both AWS and GCP routers are co-located and pre-wired.

Because Interconnect – multicloud primarily routes traffic within the VPC through a private network, to benefit from it you need to keep your query engine or jobs within a customer VPC.

2.2. High network bandwidth with on-premises systems (Reference point 4)

  • AWS Direct Connect for high-bandwidth, low-latency on-premises connectivity.

Data ingestion methods

This section focuses on ways you can use to onboard datasets (complete or subset) to a lakehouse on AWS.

1. Zero-ETL: Data movement to AWS with Zero-ETL ingestion (Reference points 5.1, 5.2)

  • AWS Zero-ETL capabilities for seamless data loading from AWS and non-AWS sources.
  • Flexibility to choose your target as an Amazon S3 based data lake or Amazon Redshift.

2. Extract, transform, load (ETL): Extract data from JDBC or SaaS sources and transform through a batch or stream pipeline (Reference points 3.1, 3.2)

The following architecture expands the flow 1.1 to 1.2 ingestion method that integrates AWS services to onboard data to the Amazon S3 raw layer and then takes it through an ETL pipeline for data cleansing and transformations. It also includes steps to onboard unstructured data to Amazon S3 using Amazon Bedrock Data Automation, and taking the lakehouse data for machine learning development with Amazon SageMaker AI.

Ingestion architecture integrating AWS services to load data into the Amazon S3 raw layer and process it through an ETL pipeline

You can also use AWS Interconnect – multicloud to run Spark jobs (Spark with Amazon EMR on EKS or open source Spark on any compute within a customer VPC) to ingest and transform data from Google Cloud with private connectivity.

3. Accessing data from Google Cloud over a private network

Refer to the preceding data access methods (3.1 and 3.2).

4. Onboarding data from AWS Outposts (S3 on Outposts) (Reference points 9.1 to 9.5)

  • Option to onboard S3 on AWS Outposts data to regional Amazon S3 through AWS DataSync (reference 9.1 to 9.3), which might be a better fit to sync files as-is through a scheduled batch or an event-driven approach.
  • Flexibility to transform the S3 on Outposts data using an Amazon EMR clusters on Outposts job, and then directly write the transformed output to a regional Amazon S3 bucket in the formats you want (including open table formats such as Apache Hudi, Apache Iceberg, and Delta Lake).

Lakehouse foundation with Apache Iceberg

By standardizing on Apache Iceberg, you’re not choosing AWS over your other platforms. You’re choosing interoperability and future flexibility. Your data becomes truly portable across any Iceberg-compatible engine.

  • Open table format: Industry-standard format supported across AWS, Databricks, Snowflake, and other platforms, which eliminates vendor lock-in.
  • ACID transactions: Reliability with full transactional consistency.
  • Time travel and schema evolution: Built-in versioning and flexible schema management.
  • Performance optimization: Advanced features such as hidden partitioning, partition evolution, and metadata management.

Note that lakehouse storage is not limited to the Apache Iceberg format, and you have the flexibility to include other open table formats (for example, Apache Hudi and Delta Lake) or file formats (for example, Apache Parquet and Apache Avro).

Unified governance and access control

AWS governance capabilities transform the lakehouse from a storage layer into a fully governed data platform. This delivers security, compliance, and data quality out of the box, applied consistently across all data sources including federated catalogs. A unified catalog consolidates metadata from AWS and non-AWS sources with generative AI-powered business glossary generation, while automated ML-powered classification identifies sensitive data (for example, PII, PHI, and financial data) across structured and unstructured datasets. AWS Identity and Access Management (AWS IAM) and AWS Lake Formation enforce fine-grained access control at the row, column, cell, and tag level, applied consistently across Amazon Athena, Amazon Redshift Spectrum, Amazon EMR, and federated sources. End-to-end data lineage tracking provides visual data flow graphs, impact analysis, and compliance audit trails. When AI agents explore metadata from the unified catalog and submit a query to Amazon Athena for execution, the Lake Formation fine-grained access control filters data based on the user interacting with the AI agent.

For the foundation model integrated into your AI agents, you can use Amazon Bedrock Guardrails, which implements customized safeguards to block harmful content and minimize hallucinations. Amazon Bedrock AgentCore provides fine-grained policy control over agent actions with real-time enforcement and managed authentication for agents accessing AWS and third-party services.

A comprehensive audit and compliance stack spans Amazon CloudWatch, AWS CloudTrail, AWS IAM, AWS Key Management Service (AWS KMS), AWS Audit Manager, and AWS PrivateLink. This stack makes sure every agent invocation is traceable, every key is managed, and every configuration is automatically mapped to frameworks including ISO, SOC, GDPR, and HIPAA.

When an end user interacts with the AI chat assistant, the layers of security and governance should go through the following.

Layer 1: Who can access?

  • Enable Active Directory and single sign-on integration for user authentication, and a combination of AWS IAM roles for AWS API-level authorization.

Layer 2: What can they see?

  • Integrate an agent profile to define what datasets each agent can access, because not all agents should have access to all datasets.
  • Enable fine-grained access control on the metadata layer using AWS Lake Formation that can filter rows and columns.
  • Enable data masking as applicable while the query responses are served through the query engine.

Layer 3: What can the agent do?

  • Control agent actions by restricting them to read-only, and apply restrictions to INSERT, UPDATE, and DELETE if the agents are supposed to query only.
  • Apply a limit on the number of rows that can be returned from the query, and apply a query scan limit to reduce cost.

Layer 4: What does the agent reveal?

  • Enable output filtering to make sure no PII is included.
  • Apply Amazon Bedrock Guardrails on large language model (LLM) responses to make sure the model does not produce anything inappropriate.
  • In addition, enable audit logging of all queries to make sure future audit and compliance needs can be met.

Comprehensive analytics ecosystem (Reference points 7.1, 7.2, 7.3)

AWS offers a complete analytics ecosystem that includes the following.

  • Amazon Athena: Serverless SQL queries with Iceberg v2 support, including provisioned capacity for consistent performance and workgroups for resource and cost management.
  • Amazon Redshift Spectrum: Federated queries across the data warehouse and Iceberg data lake.
  • Amazon Quick Sight: Enterprise visualization with governed access to all data.
  • AWS Glue and Amazon EMR: Distributed data processing capability for enterprise transformations.

AI-ready architecture (Reference points 8.1 to 8.4)

A consolidated lakehouse architecture helps you make data ready for AI agents that can access the data through readily available MCP servers or through the AWS SDK for Python (Boto3) for Amazon Athena or Amazon Redshift Spectrum. AI agents can integrate the AWS MCP Server to interact with AWS analytics services such as AWS Glue, Amazon Athena, and Amazon S3 Tables, a capability of Amazon S3, to query both data and metadata.

AI agents need context to understand how the catalog tables and their attributes are linked to each other, how users have queried them in the past, or what priorities are defined to understand which one is an authoritative source for a particular natural language question. To enable the AI agent with additional context, we can integrate the AWS Context service that was pre-announced recently at the AWS New York Summit 2026.

Governance integration: AI agents automatically inherit Lake Formation permissions, because the agent can submit the SQL query to be run through Amazon Athena or Amazon Redshift Spectrum. This makes sure they only access data that users are authorized to see. Amazon SageMaker Unified Studio data lineage tracks AI agent queries for full auditability.

The following diagram represents how the AI agent request flow looks.

AI agent request flow through the unified catalog, Lake Formation governance, and Amazon Athena

This architecture delivers value across every layer of the organization. Business teams gain faster time-to-insight by querying data across all platforms without waiting for data movement, while eliminating duplicate storage and reducing transfer costs through federation. The Apache Iceberg open table format ensures data portability and freedom from vendor lock-in. For IT and data teams, a single governance layer across all sources, including federated catalogs, reduces operational complexity, while the flexibility to choose between federation and ingestion for each use case, combined with the elastic AWS infrastructure and the petabyte-scale metadata architecture of Iceberg, delivers both agility and scalability. Data governance and compliance teams benefit from a single point of policy enforcement across all data regardless of location, complete lineage and access logs for audit and compliance reporting, automated sensitive data classification, and policies that are defined once and enforced everywhere, including across federated sources.

Architecture tradeoffs and best practices

The following are a few key trade-offs you need to consider while designing the solution.

Data ingestion and access methods

Use catalog federation (Iceberg REST) when:

  • The source platform supports the Iceberg REST API (Databricks, Snowflake Polaris).
  • Data is already in Iceberg format with Amazon S3 backed storage.
  • You want bidirectional discovery (AWS tables visible in Databricks or Snowflake too).

Use query federation (Amazon SageMaker Lakehouse architecture or AWS Glue connectors) when:

  • The source is BigQuery, SQL Server, or another non-Iceberg platform.
  • Data must stay in the source cloud (sovereignty, contractual, or latency reasons).
  • Real-time access is required without replication lag.

Use ingestion (Zero-ETL, AWS Glue, or Amazon EMR) when:

  • Data is accessed frequently with a low-latency requirement by AI agents or high-concurrency analytics.
  • The business decides to build a data lake and warehouse on AWS.
  • You need full governance, time travel, and performance optimization.

Use AWS Interconnect – multicloud when:

  • You need real-time or near-real-time query federation to GCP data sources (BigQuery, AlloyDB, Cloud Spanner) and latency or security requirements prohibit public internet routing.
  • You have high-volume, recurring data transfers between AWS and GCP where public internet egress costs or bandwidth variability are unacceptable.
  • Your organization has compliance or regulatory requirements mandating that data never traverse the public internet (HIPAA, PCI-DSS, or financial services regulations).
  • You need bidirectional connectivity, such as GCP workloads calling AWS APIs, or AWS workloads calling GCP APIs, both over private paths.

Choosing between federation and ingestion based on use case

Dimension Federation (Query in Place) Ingestion (Move to AWS)
Data freshness Real-time or near-real-time Dependent on ingestion frequency
Query performance Subject to source system latency and network Subject to data volume and operation, avoids cross-cloud network latency
Cost Lower storage cost. Higher per-query cost for cross-cloud egress Higher upfront ingestion cost. Lower ongoing query cost
Governance Partial. Source system retains some control, and a unified catalog can simplify governance for consumers Full. Lake Formation enforces all policies across all AWS analytics services
Data portability Data remains in source Data fully portable in open format
AI readiness Limited. Agents depend on source availability High. Agents query optimized, governed Iceberg tables
Operational complexity Lower initial setup. Harder to debug cross-cloud issues Higher initial setup. Simpler long-term operations

Integrating Amazon Bedrock AgentCore Gateway and Amazon Bedrock AgentCore Runtime based on use case

The following are key differences between AgentCore Gateway and AgentCore Runtime that are relevant for our use case.

Dimension Amazon Bedrock AgentCore Gateway Amazon Bedrock AgentCore Runtime
Timeout 5 minutes (hard limit) 15 min sync / 8 hours async
Statefulness Stateless (per-request) Stateful (session-based)
Best for Lightweight API proxying Long-running data processing
Your lakehouse queries Will time out frequently Handles multi-hour jobs

Because AgentCore Gateway has a 5-minute hard timeout limit, use AgentCore Runtime for data processing jobs.

  • AWS Glue ETL jobs can run for minutes to hours.
  • Amazon Redshift queries on large datasets routinely exceed 5 minutes.
  • Athena federated queries (especially cross-cloud through Interconnect) can be slow.
  • Iceberg table scans on multi-TB datasets take time.

You can use AgentCore Gateway if the scope is limited to Glue Data Catalog interactions to fetch metadata schema, because that won’t run for more than 5 minutes.

Design considerations for production implementation

In practice, there are multiple aspects to consider when deploying the solution for production. The following summarizes a few of the key issues you might encounter and approaches to address them.

Catalog federation: The metadata drift problem

One of the first surprises in production is metadata drift, the state where your federated catalog no longer reflects the actual schema of the source system, because the source system’s metadata changes are not reflected in the unified catalog. The agent continues to generate SQL against the stale schema, producing silent failures that are hard to trace.

The following are a few ways you can address the metadata drift issue.

  • Implement a catalog refresh schedule. Even a daily Glue crawler run against federated sources catches most drift before it causes agent failures.
  • Add schema validation as a pre-query step in your agent tool. Before running SQL, verify that the referenced columns exist in the current catalog metadata.
  • Instead of pulling metadata changes from the source in a scheduled manner, you can design an event-driven system, where the source system triggers a push event to run the schema change in the federated catalog.

Query federation: Latency is non-deterministic

Query federation works well for moderate data volumes, but latency becomes non-deterministic at scale. A query that returns in 3 seconds during testing can take more than 10 seconds in production when the source system is under load, the network path is congested, or the federated connector is cold-starting.

The following are a few approaches you can consider to improve the performance.

  • Set explicit query timeouts in your Athena execution context. Without them, a slow federated query will block your agent indefinitely.
  • Implement query result caching for frequently asked questions. Most business users ask the same questions repeatedly, and caching at the agent layer improves perceived performance.
  • For time-sensitive use cases, consider caching aggregated data in an AWS lakehouse on a schedule rather than querying live. This trades freshness for reliability.

AgentCore memory: Statefulness cost

AgentCore Memory enables stateful conversations, but in production, unbounded memory accumulation creates its own problems. An agent that remembers every conversation eventually starts surfacing stale context. For example, a user who asked about Q3 revenue six months ago gets that context injected into a Q1 query today.

The following are a few ways you can optimize cost and improve relevance.

  • Set explicit memory expiry (we use 30 days as shown in the implementation) and enforce it consistently.
  • Use session-scoped memory for transactional queries and long-term memory only for user preferences and recurring patterns.
  • Implement a memory review step in your LangGraph workflow. Before invoking the model, filter retrieved memories by recency and relevance score rather than injecting all of them.

LangGraph orchestration: When tool calls loop

The conditional routing of LangGraph is powerful, but in production we observed a failure mode where the agent enters a tool call loop. The model repeatedly calls the same tool with slightly different parameters, never reaching a satisfactory answer. This typically happens when the tool returns partial or ambiguous results and the model keeps trying to refine.

What we learned:

  • Add a maximum tool call counter in your LangGraph state. If the agent has called tools more than N times in a single session, force a graceful exit with a summary of what was found.
  • Return structured, unambiguous responses from your tools. Include row counts, column names, and explicit null indicators so the model can reason clearly about completeness.
  • Log every tool invocation with its input and output. This is the single most valuable debugging artifact when diagnosing agent misbehavior in production.

Handling hallucination risks in federated agent architectures

This is the most important section for teams moving from prototype to production. Hallucination in agentic AI systems that query real data is qualitatively different from hallucination in general-purpose LLMs, and it is more dangerous because the outputs look authoritative.

There are three distinct hallucination risk zones in a lakehouse AI agent:

  • SQL generation: The model generates SQL that is syntactically valid but semantically wrong. For example, when asked “What is our revenue growth this quarter?”, the model might generate a query that compares the wrong date ranges, uses the wrong aggregation function, or joins tables on incorrect keys, and then returns a confident, formatted answer with the wrong numbers.
  • Cross-source synthesis: When the agent queries multiple federated sources and synthesizes results, the risk compounds. The model may correctly retrieve customer counts from Amazon S3 and revenue figures from Snowflake, but incorrectly draw conclusions that aren’t supported by either dataset individually.
  • Memory-augmented reasoning: When long-term memory is active, the model may blend historical context with current query results in ways that are factually incorrect. For example, it might apply a business rule that was true six months ago but has since changed.

To improve, before any agent output informs a business decision, apply the following three-step validation framework:

  • Step 1: Source verification. Can you trace the answer back to a specific table, column, and row count? If the agent can’t show you the SQL and the row count, the answer is unverified.
  • Step 2: Reasonableness check. Does the answer fall within expected ranges? A sudden 10x spike in customer count is a signal to investigate.
  • Step 3: Cross-validation. For critical decisions, run the equivalent query directly in Athena or your BI tool and compare. Discrepancies reveal either a model reasoning error or a data quality issue. Resolve both before the answer is trusted.

These lessons don’t diminish the value of the architecture. They make it production-ready. The teams that move fastest with agentic AI are not the ones who skip these guardrails. They’re the ones who build them in from the start and spend less time firefighting in production.

Alternative to the unified catalog approach

In case you face technical and process challenges to unify catalogs across providers, you can let each data producer expose the metadata and data through MCP servers, as represented in the following diagram. In this approach, each producer takes the responsibility of maintaining the MCP servers and exposing them to the context layer. While this approach provides autonomy to data owners to operate independently and with flexibility, it also creates operational overhead to synchronize all metadata in a consistent way.

Alternative architecture where each data producer exposes its metadata and data through its own MCP server to the context layer

What’s next

In Part 2 of this series, we walk through the full implementation step by step, including hands-on scripts to:

  • Load example sales datasets into Databricks and marketing data to Snowflake as Iceberg tables, and federate them into AWS Glue Data Catalog through the Iceberg REST API.
  • Register Google BigQuery as a native federated data source in Amazon SageMaker, instead of a traditional AWS Lambda connector integration.
  • Create a customer master table as a native Iceberg table in Amazon S3.
  • Run a single SQL query in Amazon Athena that joins all four sources across two federation patterns, with no data movement.
  • Deploy an AI agent on Amazon Bedrock AgentCore that can autonomously query the same unified catalog using Amazon Athena and answer complex business questions in natural language queries. In addition, integrate AgentCore Memory to persist user context.

Conclusion

In this post, we summarized how you can unify data access across multiple cloud and ISV providers on AWS with the combination of catalog federation, query federation, and data movement to AWS. We then explained how AWS Glue Data Catalog and Lake Formation help provide unified catalog and access governance, and how AI agents hosted in Amazon Bedrock AgentCore can access it using MCP servers to explore the metadata context, convert user natural language queries to SQL, and use Amazon Athena to run the query across data sources to get the response to the end user. In addition, we provided an overview of different data ingestion methods to build a lakehouse architecture on AWS, including AWS Interconnect – multicloud and where it adds value.

We also provided architecture trade-offs and best practices to integrate the service capabilities. In the next post (Part 2), we will take a specific use case and provide a step-by-step implementation guide to unify the catalog and deploy the agent to Amazon Bedrock AgentCore.


About the author

Sakti Mishra

Sakti Mishra

Sakti is a Principal Data and AI Solutions Architect at AWS, where he helps customers modernize their data architecture and define end-to-end data strategies, including data security, accessibility, governance, and more. He is also the author of Simplify Big Data Analytics with Amazon EMR and AWS Certified Data Engineer Study Guide. Outside of work, Sakti enjoys learning new technologies, watching movies, and visiting places with family. You can connect with Sakti through his LinkedIn profile.

How MAPFRE USA modernized fraud claims with Amazon EMR Serverless

Post Syndicated from Lijan Kuniyil original https://aws.amazon.com/blogs/architecture/how-mapfre-usa-modernized-fraud-claims-with-amazon-emr-serverless/

Insurance fraud remains a significant challenge for the insurance industry. Fraudulent claims can increase loss costs, reduce trust, and consume investigation capacity that could otherwise be focused on serving customers. Traditional fraud detection approaches typically rely on rules-based controls, manual investigation triggers, historical claim patterns, and structured-data-only analysis. These approaches are useful for known fraud patterns, but they can struggle to detect sophisticated fraud rings or hidden relationships across claimants, policies, vehicles, providers, addresses, and prior suspicious activities.

MAPFRE USA is a top-rated auto and home insurer in Massachusetts, serving customers in 11 states nationwide. Our coverage includes auto, home, motorcycle, watercraft, business insurance, and more. As part of MAPFRE Group, we’re a worldwide leader serving over 31.1 million customers in more than 100 countries with a team of 31,000 employees. In collaboration with AWS and Neo4j, MAPFRE USA modernized its fraud prevention capabilities by combining graph-based features with machine learning (ML) models deployed on AWS. This initiative focused initially on Massachusetts auto insurance and later expanded to home insurance. It has delivered significant business impact, exceeding $5 million in net present value (NPV) over five years, with realized savings already outperforming projections.

In this post, we share how MAPFRE USA designed and implemented this solution, highlight the technical architecture running on AWS, specifically the MAPFRE data platform called Atenea, and explore lessons learned that can apply to other industries facing complex fraud challenges.

Business challenge

Fraudulent claims aren’t always isolated events. They often involve hidden networks of policyholders, vehicles, providers, and prior suspicious activities. Detecting these complex relationships requires going beyond traditional structured data analysis.

MAPFRE set out with a clear goal:

  • Goal: Improve fraud detection accuracy and claims handling efficiency.
  • Key performance indicator (KPI): Identify fraudulent claims missed by traditional methods.
  • Approach: Develop several ML models using both traditional structured data and 54 graph-based features derived from claim relationships.
  • Deployment: Integrate with Guidewire Claims, so front-line adjusters automatically receive fraud alerts with explanations.

Each flagged claim exposure generates a Guidewire activity showing the top three model drivers, helping investigators understand why the claim was flagged and act quickly.

Technical solution on AWS (Atenea data platform)

The fraud detection platform is built on a modern data architecture on AWS, designed to scale efficiently and support long-term governance.

At its core, the solution uses Apache Iceberg tables stored on Amazon Simple Storage Service (Amazon S3), with metadata managed through the AWS Glue Data Catalog and access governed through AWS Lake Formation as part of the Atenea lakehouse governance model. The platform feature store is implemented through feature-store-managed Iceberg tables that manage model features, predictions, and Guidewire activities. The implementation is structured across three logical layers:

  • Silver layer: Iceberg tables that contain source data from each of the sources. Used as the initial consumption point of the platform.
  • Gold layer: Iceberg tables storing intermediate data, such as unified Guidewire activity logs, Auto features, and Home features.
  • Platinum layer: Feature Store-managed Iceberg tables containing encoded features and model predictions, making them reusable across models and ensuring strong metadata governance.

Processing pipelines are executed on Amazon EMR Serverless, with orchestration managed by Apache Airflow operators running on Amazon Managed Workflows for Apache Airflow (MWAA). This provides elastic, cost-efficient compute for both batch processing and fast-time scoring, while keeping orchestration, monitoring, and recovery centralized.

For graph enrichment, the platform connects to Neo4j using a dedicated driver, enabling advanced network-based features like suspicious claim linkages, provider fraud ratios, and centrality metrics.

This architecture supports efficient, reliable, and transparent production execution. It uses repeatable Airflow orchestration, environment-based continuous integration and continuous delivery (CI/CD) promotion, centralized monitoring, failure notifications, retry mechanisms, dead-letter queue handling for Guidewire integration, and controlled secret management. At the same time, the layered lakehouse design keeps the platform flexible enough to evolve with new business needs and fraud detection use cases.

Fraud detection architecture on AWS showing data ingestion to Amazon S3, the Silver, Gold, and Platinum Iceberg layers, Neo4j graph enrichment, Amazon EMR Serverless processing, and Guidewire integration

The data sources here are policy, claims, vehicles, and notes (from AS400 and Guidewire), which are structured data. Derived features that capture entity relationships make up the graph data.

Let’s go through the architecture overview:

  1. Data ingestion – Claim batch data is uploaded to Amazon S3. The data is standardized and materialized in Iceberg tables within the Silver layer.
  2. Graph enrichment – Data processed to update Neo4j graph database hosted on AWS.
  3. Model training and scoring – Batch scoring for several ML models.
  4. Model orchestration – Unified orchestration for ingestion, training, and inference using Apache Airflow operators. CI/CD pipelines for promotion across environments.
  5. Execution platform – Amazon EMR Serverless for cost-efficient Spark processing. Migration to Apache Iceberg plus AWS Glue Data Catalog for scalable metadata handling.
  6. Integration with claims systems – Fraud predictions automatically create Guidewire activities, enriched with a description for investigators.
  7. Secrets and security – AWS Secrets Manager securely stores credentials and tokens for Guidewire API integration, with environment-specific and region-specific access controls.
  8. Monitoring and reliability – Amazon CloudWatch and Amazon Simple Notification Service (Amazon SNS) provide visibility into pipeline health and notify teams on failures. Data quality checks are executed at key stages of the pipeline to validate data availability, schema consistency, completeness, and business-rule expectations before outputs are consumed by models or sent to Guidewire.

Guidewire integration with MLOps on AWS

One of the most important parts of MAPFRE’s solution was closing the loop between ML predictions and the claims handling system. This required a resilient integration between the Atenea data platform on AWS and Guidewire Claims.

Integration flow:

  1. When an ML use case finishes scoring, the results are written as JSON files into the S3 path: <bucket_name>/guidewire/.
  2. An S3 event notification triggers the AWS Lambda function LambdaXXXInvokeGuidewireAPI.
  3. This Lambda function:
    • Reads the JSON file.
    • Calls the Guidewire Predictive Model API.
    • Because Guidewire doesn’t support batch requests, the Lambda function sends each JSON payload individually. This keeps the integration compatible with Guidewire and isolates failures at the individual activity level, but it increases the number of API calls and makes retry, throttling, DLQ handling, and monitoring controls important.
  4. If successful, the API responds with HTTP 201 (activity created).
    • If not, the Lambda retries up to two times.
    • Failed requests are sent to an SQS Dead-Letter Queue (DLQ) and an SNS notification is published to an SNS queue for monitoring.
  5. Secrets are stored in AWS Secrets Manager and injected as Lambda environment variables, along with AWS Region-specific URLs for token retrieval and API endpoints.
  6. Example JSON structure for Guidewire integration:
    {
      "method": "createPredictiveActivity",
      "params": [
        {
          "claimNumber": "AUXXXXXXX",
          "exposureNumber": 1,
          "subject": "Fraud alert from ML model",
          "description": "Claim flagged as potential fraud based on graph + ML features",
          "shortSubject": "ML_Fraud_Flag",
          "priority": "high",
          "availableForClosedClaim": true,
          "autoCloseOnExposureClosure": false,
          "targetDays": 4,
          "escalationDays": 6
        }
      ]
    }

Guidewire integration flow from Amazon S3 to an AWS Lambda function that calls the Guidewire API, with an SQS dead-letter queue and Amazon SNS for failures

Key benefits of this integration:

  • Real-time actionability – Fraud predictions automatically create Guidewire activities for front-line adjusters.
  • Resilience – Built-in retries, DLQ handling, and SNS alerts keep failed events from being lost.
  • Security – Secrets and tokens are managed using AWS Secrets Manager, with strict environment separation (dev, pre, pro).
  • Scalability – Any new MLOps use case writes results into the S3 output path, automatically flowing into Guidewire.

This integration shows that fraud models don’t just exist in isolation but actively augment daily claim workflows in production. It connects Atenea’s MLOps pipelines on AWS directly with business decisioning systems, which is critical to realizing the fraud savings impact.

Data quality and resilience

For robustness, we apply data quality checks on ingestion pipelines and graph features. Automated validation detects anomalies early, monitoring dashboards track KPIs and model performance, and standardized recovery and promotion processes run across environments.

Visualization and investigative tools

Neo4j Bloom supports Special Investigations Unit (SIU) workflows by visually exploring entity relationships, such as a provider linked across multiple suspicious claims, accelerating fraud ring identification.

Neo4j Bloom graph visualization showing a provider node linked across multiple suspicious insurance claims

Conclusion

The fraud detection model for auto claims has enhanced MAPFRE USA’s ability to identify fraudulent activity, driving significant savings and improving overall claims efficiency.

During the pilot phase alone, savings exceeded projections by over half a million dollars, and in production the initiative has proven an NPV of more than $5M at current business volumes. These results confirm the business case and highlight the strength of combining structured data with graph-based features to uncover fraud networks that traditional approaches miss.

The results have been compelling:

  • Accuracy gains – detection improved by 50–135 percent compared to baseline methods.
  • Realized value – In 2025, MA Auto and MA Home claim savings reached a combined total of $6.81M, with $6.59M from MA Auto and $225K from MA Home.
  • Proven return on investment (ROI) – the project delivered an NPV of $4.7M at approval, and results are already exceeding expectations.
  • Cross-functional success – the initiative brought together Claims, IT Data, Advanced Analytics, and Neo4j teams in an agile, collaborative model.

Beyond the financial outcomes, several lessons emerged. First, cross-functional collaboration between groups like Claims, Data Engineering, Advanced Analytics, and technology partners like AWS and Neo4j was critical to success. Second, explainability proved essential. By presenting adjusters with the top model drivers directly in Guidewire, we increased trust and adoption of the system substantially. Finally, building resilience into the architecture through monitoring, retries, and data quality processes helped the models operate reliably in production.

Looking ahead, the platform is well-positioned to expand beyond fraud detection. New use cases such as underwriting anomaly detection, customer entity resolution, and retention modeling are already on the roadmap. With a robust architecture built on AWS using Amazon EMR Serverless, Apache Iceberg on Amazon S3 supported by AWS Glue Data Catalog and AWS Lake Formation, a custom-built Feature Store, and Neo4j, MAPFRE now has a scalable foundation to continue driving innovation and business impact.

To start building a similar solution, open the Amazon EMR console and review the AWS Architecture Center for reference patterns you can adapt to your own fraud detection and analytics workloads.


About the authors

Introducing Apache Spark Connect support in AWS Glue interactive sessions

Post Syndicated from Zach Mitchell original https://aws.amazon.com/blogs/big-data/introducing-apache-spark-connect-support-in-aws-glue-interactive-sessions/

When we built AWS Glue interactive sessions, our goal was to make AWS Glue as interactive as running local Python from a notebook. We mostly succeeded. With a straightforward Python package and a Jupyter notebook, you could execute remotely against the AWS Glue ephemeral Spark backend. The Livy-based approach was ahead of its time, but it had limitations from its REST-based protocol. Running local PySpark unlocked powerful integrated development environment (IDE) features such as debugging and linting, so your environment could understand the code and help you develop Spark applications more quickly. Customers would often split their development work. They used local Spark (or Docker containers) to develop in an IDE on a small amount of data, then switched to AWS Glue interactive sessions to validate scaling and tuning against the full dataset.

With modern PySpark releases came a new protocol: Apache Spark Connect. Spark Connect bridges the gap between these two worlds: you develop in local Python, but execute on AWS Glue against actual data. Today, AWS Glue interactive sessions support Spark Connect natively. You can connect from any environment that supports the PySpark remote() API, including VS Code, PyCharm, Amazon SageMaker Unified Studio notebooks, and standalone Python applications. You don’t need to install specialized kernels or manage cluster infrastructure.

What Spark Connect changes

Spark Connect, introduced in Spark 3.4, decouples the Spark client from the server through a lightweight gRPC protocol. Instead of running your driver program on the cluster, your IDE communicates with a remote Spark server through a thin client layer. This architecture unlocks the key workflow improvement: you develop locally and execute remotely.

Spark Connect architecture diagram showing a thin client communicating with a remote Apache Spark server

Spark Connect architecture — thin client with the full power of Apache Spark

With Spark Connect support in AWS Glue interactive sessions, you get:

  • IDE freedom – Use VS Code, PyCharm, JupyterLab, or any Python environment. No kernel installation required.
  • Programmatic access – Build Spark into your Python applications and automation scripts with a standard SparkSession.builder.remote() call.
  • Serverless execution – AWS Glue provisions and manages the Spark cluster. You pay only for the data processing units (DPUs) consumed while your session is active.
  • Spark Connect monitoring – The Spark Live UI now includes a dedicated Connect tab showing active Spark Connect sessions and operations alongside the existing Jobs, Stages, and Executors views.

Getting started with SageMaker Unified Studio

Amazon SageMaker Unified Studio provides the most direct path to Spark Connect on AWS Glue. The notebook environment handles session creation, endpoint retrieval, and token refresh automatically, so no connection boilerplate is required.

Prerequisite: You need an Amazon SageMaker Unified Studio project to use this workflow. If you don’t have one, create a project in your SageMaker Unified Studio domain first.

To connect to an AWS Glue Spark Connect session:

  1. Sign in to SageMaker Unified Studio, choose your project, and create or open a Notebook.

A notebook open in SageMaker Unified Studio

A notebook open in SageMaker Unified Studio

  1. Choose the compute icon in the left toolbar to open the Compute environment panel. Expand the Spark section.

Compute environment panel in SageMaker Unified Studio with the Spark section expanded

The Compute environment panel with the Spark dropdown list

  1. Select a Glue Spark connection. Depending on your SageMaker domain configuration, you will see either default.spark or named connections such as project.spark.compatibility. Select the appropriate Glue (Spark) connection and choose Apply.

Notebook cell showing spark.version returns 3.5.6-amzn-1 after connecting to Glue Spark Connect

Connected to Glue Spark Connect — running spark.version returns ‘3.5.6-amzn-1’

After you make your selection, you’re connected. The spark session object is available natively. No imports or configuration are needed. Start running PySpark immediately:

spark.sql("SHOW DATABASES").show()

The session manages itself in the background, including automatic token refresh.

Using the sagemaker_studio SDK

The sagemaker-studio Python package extends the Spark Connect experience beyond SageMaker Unified Studio notebooks into local IDEs, continuous integration and continuous delivery (CI/CD) pipelines, and any Python environment. The sparkutils module handles session initialization and connection configuration in a single call. You get the same streamlined experience as in the notebook, anywhere you run Python:

from sagemaker_studio import sparkutils

# Initialize a Glue Spark Connect session using your project connection
spark = sparkutils.init(connection_name="default.spark")

# Run queries immediately
spark.sql("SHOW DATABASES").show()

You can also use sparkutils.get_spark_options() to retrieve pre-configured Java Database Connectivity (JDBC) options for reading and writing to data sources through your project connections. Supported sources include Amazon Redshift, Amazon Aurora, and Amazon DocumentDB (with MongoDB compatibility):

# Get connection options for a Redshift connection in your project
options = sparkutils.get_spark_options("my_redshift_connection")

# Read from Redshift via Spark Connect
df = spark.read.format("jdbc").options(**options).option("dbtable", "analytics.orders").load()
df.show()

Within SageMaker Unified Studio, the sagemaker-studio SDK is native to the environment. The spark session and sparkutils are available without installation. For local IDE use, install it with pip install sagemaker-studio and configure credentials through an AWS named profile or boto3 session.

How it works

Spark Connect sessions in AWS Glue use a three-step workflow:

  1. Create a session – Call the CreateSession API with SessionType set to SPARK_CONNECT. The session provisions in approximately 30 seconds.
  2. Retrieve the endpoint – Call GetSessionEndpoint to receive a sc:// gRPC endpoint URL and a time-limited authentication token.
  3. Connect with PySpark – Pass the endpoint and token to SparkSession.builder.remote() and start running Spark operations.

Spark Connect protocol flow from the DataFrame API to a logical plan, sent over gRPC and protobuf, with results streamed back over gRPC and Arrow

Spark Connect protocol flow — DataFrame API translated to logical plan, sent via gRPC/protobuf, results streamed back via gRPC/Arrow

Connecting with the low-level API

Some environments don’t have the sagemaker-studio SDK, such as custom containers, AWS Lambda functions, or non-Python toolchains. In these environments, or if you’re not using SageMaker Unified Studio, you can use the AWS SDK (Boto3) to manage sessions directly. The following example demonstrates the full workflow:

import time, boto3, urllib.parse
from pyspark.sql import SparkSession

glue = boto3.client("glue", region_name="us-east-1")

# 1. Create a Spark Connect session
session_id = "my-spark-connect-session"
glue.create_session(
    Id=session_id,
    Role="arn:aws:iam::123456789012:role/GlueServiceRole",
    Command={"Name": "glueetl"},
    GlueVersion="5.1",
    SessionType="SPARK_CONNECT",
    DefaultArguments={"--enable-spark-live-ui": "true"},
)

# 2. Wait for the session to reach READY
while True:
    status = glue.get_session(Id=session_id)["Session"]["Status"]
    if status == "READY":
        break
    time.sleep(5)

# 3. Get the Spark Connect endpoint
sc = glue.get_session_endpoint(SessionId=session_id)["SparkConnect"]
endpoint_url = sc["Url"]
auth_token = sc["AuthToken"]

# 4. Connect with PySpark
encoded_token = urllib.parse.quote(auth_token, safe="")
connection_string = f"{endpoint_url}:443/;use_ssl=true;x-aws-proxy-auth={encoded_token}"
spark = SparkSession.builder.remote(connection_string).getOrCreate()
spark.sql("SELECT 1 + 1 AS result").show()

Monitoring with Spark Live UI

When you enable the Spark Live UI at session creation, you gain access to a real-time dashboard showing:

  • Jobs and Stages – Track active, completed, and failed jobs with stage-level metrics.
  • Executors – Monitor memory usage, shuffle data, and executor health.
  • SQL – Inspect query plans and execution details.
  • Connect tab – View active Spark Connect sessions and operations (specific to Spark Connect).

Access the dashboard through the GetDashboardUrl API or directly from the AWS Glue console.

import boto3, webbrowser

glue = boto3.client("glue", region_name="us-east-1")
dashboard = glue.get_dashboard_url(
    ResourceId="my-spark-connect-session",
    ResourceType="SESSION",
)
webbrowser.open(dashboard["Url"])

In SageMaker Unified Studio, no API call is needed. Choose Ready in the notebook status bar to open the kernel info popover. From there, open the Spark UI link for the live dashboard or Spark Driver Logs for real-time log output.

Notebook status bar Ready button that opens the Spark UI and Spark Driver Logs links

Image showing “Ready” in the status bar to access Spark UI and Driver Logs directly from the notebook

Token refresh

Authentication tokens expire after 30 minutes. In SageMaker Unified Studio, this is handled automatically. For programmatic use, you can use a background thread to keep the connection alive. The following helper reconnects transparently before the token expires:

import threading, time, boto3, urllib.parse
from pyspark.sql import SparkSession

class GlueSparkConnect:
    """Maintains a SparkSession with automatic token refresh."""

    def __init__(self, session_id, region="us-east-1", refresh_margin=300):
        self.session_id = session_id
        self.glue = boto3.client("glue", region_name=region)
        self.refresh_margin = refresh_margin  # seconds before expiry to refresh
        self._lock = threading.Lock()
        self.spark = self._connect()
        self._start_refresh_loop()

    def _connect(self):
        sc = self.glue.get_session_endpoint(SessionId=self.session_id)["SparkConnect"]
        encoded_token = urllib.parse.quote(sc["AuthToken"], safe="")
        remote_url = f"{sc['Url']}:443/;use_ssl=true;x-aws-proxy-auth={encoded_token}"
        self._token_expiry = sc["AuthTokenExpirationTime"].timestamp()
        return SparkSession.builder.remote(remote_url).getOrCreate()

    def _start_refresh_loop(self):
        def _loop():
            while True:
                sleep_for = max(self._token_expiry - time.time() - self.refresh_margin, 30)
                time.sleep(sleep_for)
                with self._lock:
                    self.spark = self._connect()
        t = threading.Thread(target=_loop, daemon=True)
        t.start()

# Usage
session = GlueSparkConnect("my-spark-connect-session")
session.spark.sql("SELECT 1 + 1 AS result").show()

The background thread sleeps until 5 minutes before token expiry, then transparently reconnects. Because the daemon thread exits when your script ends, there is no cleanup required.

Getting started

To start using Spark Connect with AWS Glue interactive sessions:

  1. Use AWS Glue version 5.1 (Apache Spark 3.5.6).
  2. Install PySpark 3.5.6 locally: pip install pyspark==3.5.6.
  3. Grant your AWS Identity and Access Management (IAM) identity permissions for glue:CreateSession, glue:GetSession, and glue:GetSessionEndpoint.
  4. Create a session with --session-type SPARK_CONNECT and connect from your preferred environment.

VPC note: If you connect to AWS Glue interactive sessions through a virtual private cloud (VPC) endpoint, add the new Spark Connect endpoint (com.amazonaws.{region}.glue.sessions) to your VPC configuration. Existing AWS Glue VPC endpoints don’t cover Spark Connect traffic.

For detailed instructions, see Connecting to a Spark Connect session in the AWS Glue Developer Guide.


About the authors

Zach Mitchell

Zach Mitchell

Zach is a Senior Big Data Architect at AWS Worldwide Specialist Organization for Analytics. He works with customers to design and build data applications on AWS, with a focus on SageMaker Unified Studio, AWS Glue, and AWS Lake Formation. Outside of work, he enjoys building things with code and occasionally writing about it.

Shrey Malpani

Shrey Malpani

Shrey is a Senior Technical Product Manager at AWS Analytics. He is focused on building and scaling data processing, data integration, and data management capabilities across services like AWS Glue, Amazon EMR, and Amazon Redshift that help customers build AI-ready data platforms for their analytics or machine learning workflows.

Vaibhav Naik

Vaibhav Naik

Vaibhav is a Software Engineer at AWS Glue, where he leads the development of enterprise Generative AI managed services and Agentic data systems. He has over a decade of experience designing massive-scale cloud infrastructure and distributed computing platforms.

Tom Olson

Tom Olson

Tom is a Software Development Engineer on the AWS Glue team, focused on Interactive Sessions and operational excellence. He brings over 20 years of software development experience, including government contracting and EC2 Networking at AWS. Outside of work, he enjoys running and playing board games.

Gaurav Krishnan

Gaurav Krishnan

Gaurav is a Software Development Engineer at AWS Glue. He has a deep interest in distributed systems and creating low-friction developer experiences for interactive data workloads on Apache Spark. In his spare time, he enjoys running and trying new restaurants.

How BigBasket uses the Iceberg based lakehouse architecture on AWS to power lightning-fast grocery delivery across India

Post Syndicated from Annie Mattoo original https://aws.amazon.com/blogs/big-data/how-bigbasket-uses-the-iceberg-based-lakehouse-architecture-on-aws-to-power-lightning-fast-grocery-delivery-across-india/

Delivering fresh groceries to millions of customers across India in a few minutes demands a radically modern data architecture and resilient processes to help the business make faster decisions. This is what BigBasket was able to achieve by building a lakehouse architecture on AWS.

In this post, we demonstrate how BigBasket implemented the lakehouse architecture on AWS, including their architecture decisions, implementation approach, and the measurable business results you can expect from a similar modernization. Whether you’re facing scalability challenges or planning your own lakehouse implementation, this blueprint provides actionable insights you can adapt for your organization.

About BigBasket

BigBasket (Innovative Retail Concepts Private Limited) is India’s largest online supermarket, serving millions of customers across over 60 cities. Founded in 2011, the company offers groceries, fresh produce, household items, and personal care products through its mobile app and website, operating subscription services (BBDaily) and quick commerce (bbnow). For BigBasket, the ability to deliver groceries on time isn’t only a competitive advantage. It’s the foundation of customer trust, where every minute counts.

However, rapid business growth brought significant operational challenges:

  • Inability to consistently meet on-time delivery adherence because of high order volumes, extended travel times, and more, directly impacting key metrics like on-time rate (OTR)-10 mins and OTR-15 mins.
  • Struggling to meet on-time delivery targets because of picking inefficiency, high order volumes, and extended travel times, directly impacting key metrics like OTR-10 mins and OTR-15 mins.
  • Delays in stock availability impacting vendor fill-rates, inter-distribution center orders, and warehouse operations.
  • Inaccurate stock forecasting for top-selling stock keeping units (SKUs), assortment variety, event SKUs, store capacity, and buying cycles.
  • Lower dark store productivity across picking, stacking, order processing, and goods receipt notes (GRN).

Behind these business challenges lay a fundamental technology problem: the existing data infrastructure couldn’t keep pace. The company experienced rapid store growth, expanding 4x in a short timeframe, which exposed several limitations within their existing data architecture that needed attention.

Understanding the technical bottlenecks

BigBasket’s initial architecture relied heavily on a single data warehouse built on Amazon Redshift to meet all reporting and dashboarding needs. While this traditional approach had served them well initially, several important limitations emerged:

  • Stale data: Extract, transform, load (ETL) pipelines delivered only day-old (D-1) data, making near real-time analysis impossible for dashboard requirements.
  • Extended recovery times: Pipeline failure recovery processes took several hours, causing significant delays in data availability for business users.
  • Schema rigidity: Schema changes in source databases frequently triggered pipeline failures because of a lack of schema evolution support.
  • Scalability constraints: The infrastructure struggled to handle the sudden load increase from 13,000 to over 35,000 transactions for reports and dashboards with more than 1,000 dataset refreshes.
  • Cost implications: Increasing data volumes demanded additional compute resources, driving up costs.

Diagram of the scalability and cost limitations of BigBasket’s legacy Amazon Redshift data warehouse

It became clear that the existing data infrastructure wasn’t able to meet the evolving business requirements and a redesign of their data architecture is needed.

Why lakehouse architecture?

A modern data lakehouse architecture addresses these issues with near real-time data processing, flexible schema evolution, and scalable analytics, capabilities necessary for fast-moving commerce operations. The lakehouse approach combines the flexibility and cost-effectiveness of data lakes with the performance and governance features of data warehouses, combining the strengths of both. The design of a data lakehouse provides interoperability across storage systems for combined analytics activities.

Solution overview

BigBasket partnered with AWS to implement a comprehensive lakehouse architecture using a combination of AWS native services and open-source technologies.

The following diagram shows an elaborated view of Bigbasket’s modernized architecture on AWS.

Detailed lakehouse data flow across bronze, silver, and gold medallion layers on AWS

Data ingestion: Enabling continuous replication

AWS Database Migration Service (AWS DMS) ingests data from online transaction processing (OLTP) databases running on Amazon Relational Database Service (Amazon RDS) into the lakehouse on AWS.

This method continuously replicates data with minimal latency, so your analytics reflect near real-time business operations.

Storage and governance: Building a solid foundation

The lakehouse is built on Amazon Simple Storage Service (Amazon S3) and Amazon Redshift, which serve as the centralized data lake and warehouse following a medallion architecture.

The architecture persists all analytical data using Apache Iceberg as the open table format. Iceberg provides a robust foundation for large-scale analytics with the following capabilities:

  • ACID transactions: Guarantees data consistency and correctness across concurrent read and write operations.
  • Time travel: Supports querying historical table versions for auditing, troubleshooting, and recovery.
  • Schema evolution: Allows schema changes without disrupting existing queries or downstream pipelines.

The medallion architecture structures data across three logical layers within the lakehouse:

  • Bronze layer: Implements change data capture (CDC)-based source replication using AWS DMS. Raw change events flow into Amazon S3 as Apache Parquet files in their original format from source systems, preserving the complete change history. The data pipeline processes and deduplicates these events using Apache Spark on Amazon EMR to create and maintain Apache Iceberg tables that act as replicated source tables.
  • Silver layer: Represents the conformed data model, where data is cleansed, standardized, and validated with enforced quality checks. This layer contains core dimension and fact tables, modeled for analytical consistency and reuse across domains. Data is stored as Apache Iceberg tables on Amazon S3, making it reliable and performant for downstream analytics and transformations.
  • Gold layer: Provides business-ready data marts and wide tables optimized for reporting, dashboarding, and domain-specific use cases. These datasets are curated to align with business metrics and key performance indicators (KPIs) and are served from Amazon Redshift, using Iceberg-backed tables to deliver fast, scalable analytics for business intelligence (BI) tools and end users.

This layered approach maintains a clear separation of concerns across raw ingestion, analytical modeling, and business consumption, while supporting scalability and flexibility across the organization. AWS Lake Formation enforces fine-grained data access controls, and the AWS Glue Data Catalog centrally manages metadata across Amazon S3 and Amazon Redshift, ensuring consistent data discovery and governance across the analytics ecosystem.

Data processing: Flexibility and performance

For data processing and transformations, BigBasket uses Amazon EMR with Apache Spark and dbt, orchestrated by Apache Airflow running on Amazon Elastic Kubernetes Service (Amazon EKS) as the core compute layer of the lakehouse. Apache Spark on Amazon EMR handles large-scale distributed processing, including CDC deduplication, incremental transformations, and complex data reshaping. Apache Iceberg serves as the open table format, which provides several critical capabilities.

dbt is used to define and execute transformation logic using SQL, managing the build of data models such as staging, intermediate, and final tables on top of the raw data. dbt uses the dbt-Trino adapter to run these transformations using the Trino engine, materializing the results as Apache Iceberg tables in Amazon S3. This approach provides a simple, modular, and governed way to manage transformations while taking advantage of Iceberg’s transactional guarantees.

These features are necessary for production lakehouse implementations and help you avoid vendor lock-in while maintaining enterprise reliability.

Online analytical processing (OLAP) and analytics: Hybrid approach for cost optimization

The analytics layer uses a hybrid approach that you can adapt based on your query patterns:

  • Amazon Redshift: For querying of active, frequently accessed data from the Gold layer.
  • Amazon Athena: For ad-hoc queries on historical data.
  • Apache Trino: For federated queries across multiple data sources while powering dbt-driven transformations directly on Apache Iceberg tables.

This hybrid strategy optimizes costs by keeping frequently accessed data in Amazon Redshift while querying historical data directly from Iceberg tables in Amazon S3. Amazon Redshift data sharing supports a multi-warehouse architecture for cross-team collaboration, allowing different teams to access shared datasets without data duplication.

Orchestration: Managing complex workflows

Apache Airflow running on Amazon EKS orchestrates and schedules data pipelines across the entire environment, providing visibility and control over complex workflows. This gives you a unified view for monitoring and managing your data operations.

Machine learning integration

Amazon SageMaker AI powers machine learning workloads for predictive analytics and model training directly on lakehouse data, from demand forecasting to delivery optimization. This tight integration means your data scientists can work with the same governed data that powers your analytics.

Visualization: Making insights accessible

Amazon Quick Sight provides data visualization and business intelligence reporting capabilities, making insights accessible to business users across the organization without requiring technical expertise.

Special focus: Clickstream data processing

BigBasket implemented a sophisticated dual-path architecture for processing clickstream data from mobile apps and web interactions:

  • Real-time path: Data flows through Scala stream collectors on Amazon Elastic Compute Cloud (Amazon EC2) (behind Elastic Load Balancing) to Amazon Kinesis Data Streams and Amazon OpenSearch Service for immediate insights into customer behavior. This path is necessary when you need to react to user actions within seconds, for example detecting fraud or personalizing experiences in real time.
  • Batch path: The batch path validates data, stores it in Amazon S3, processes it through Amazon EMR, and loads it into Amazon Redshift for comprehensive historical analysis. This path handles data quality checks, enrichment, and aggregation for long-term analytics.

The trade-off between these approaches is latency versus completeness. Real-time processing gives you speed but may sacrifice some data quality checks, while batch processing provides accuracy but introduces delay. This dual approach achieves both immediate operational insights and deep analytical capabilities, letting you optimize for different use cases.

The following diagram shows how the clickstream data is handled and effectively processed today.

BigBasket’s dual-path clickstream processing architecture with real-time and batch paths on AWS

The results: measurable business impact

The data platform transformation achieved significant results across multiple dimensions:

Technical improvements

  • Near real-time data: Achieved near real-time data availability for dashboards within 3–5 minutes, replacing previously day-old data.
  • Rapid failure recovery: Pipeline failure re-runs now complete in minutes instead of hours.
  • Comprehensive governance: Full control over data governance with robust observability, lineage, data accuracy, and consistency.
  • Enhanced scalability: Successfully handling over 35,000 reports and dashboards with over 1,000 dataset refreshes.

Business outcomes

  • On-time delivery: Improved monitoring with real-time insights on low-performing stores.
  • Stock availability: Reduced operational issues with visibility into key bottlenecks.
  • Stock forecasting: Improved accuracy and availability of top-selling SKUs.
  • Dark store productivity: Enhanced productivity of warehouse executives across all operations.

Key takeaways: lessons for modern data platforms

BigBasket’s journey offers valuable insights for organizations facing similar challenges:

  1. Quick commerce needs quick observability. In the fast-paced world of quick commerce, faster decision-making directly improves business metrics. Real-time data isn’t a luxury. It’s a necessity.
  2. Embrace ELT for real-time needs. Shifting from traditional ETL to an extract, load, transform (ELT) pattern within a lakehouse architecture is important to unlock near real-time analytics capabilities.
  3. A lakehouse delivers speed and governance. Modern lakehouse architectures don’t force trade-offs. You can achieve both fast data availability and comprehensive control, lineage, and accuracy.
  4. Focus on operational resilience. Designing for rapid failure recovery (re-runs in minutes, not hours) is necessary for maintaining data availability and business trust, especially in customer-facing operations.
  5. Incremental migration. You don’t need to rebuild everything. Evolve your current Amazon S3 data lake or reuse your existing investments in Amazon Redshift to build the data lakehouse capabilities.

The road ahead

BigBasket continues to innovate, now moving to adopt Amazon SageMaker Unified Studio to access all lakehouse components in a simplified manner across the enterprise. This next evolution will further streamline data access and accelerate insights across teams.

The company’s transformation demonstrates that with the right architecture and AWS services, organizations can turn data infrastructure challenges into competitive advantages, delivering not only better analytics but better customer experiences.

As you plan your own lakehouse implementation, use these patterns and lessons learned to accelerate your journey and avoid common pitfalls.


About the authors

Naga Sandeep Grandhi

Naga Sandeep Grandhi

Sandeep is an engineering leader at BigBasket, driving data platform and cloud architecture initiatives, including the next-gen data lake built for scale, reliability, and real-time insights.

Vikram Kumar

Vikram Kumar

Vikram is a Principal Engineer at BigBasket, where he leads the data engineering team. He specializes in designing and scaling modern data platforms on AWS, enabling BigBasket to process large-scale data efficiently and power data-driven decision-making across the organization.

Annie Mattoo

Annie Mattoo

Annie is a Sr. Analytics Specialist at AWS, bringing over 15+ years of expertise in helping customers with their DATA & AI journeys. She has successfully led customer teams to seamlessly adopt AWS Data & AI services and has worked with Fortune 500 customers across the globe in her previous roles.

Vineet Thapliyal

Vineet Thapliyal

Vineet is an Enterprise Account Manager at Amazon Web Services (AWS) in Bengaluru, India, where he manages strategic cloud and generative AI engagements across some of India’s largest conglomerates spanning energy, retail, and technology. He is passionate about helping enterprises unlock business value through AI/ML, cloud modernization, and industry-specific innovation — from renewable energy analytics to retail transformation at scale.

Anirudh Chawla

Anirudh Chawla

Anirudh is an Analytics Solution Architect at AWS. He helps organization empowers businesses to harness their data effectively through AWS’s analytics platform. His interest lies in building highly available distributed systems.

Accelerating log analytics at scale with AWS Glue and Apache Iceberg materialized views

Post Syndicated from Shinu Tharol original https://aws.amazon.com/blogs/big-data/accelerating-log-analytics-at-scale-with-aws-glue-and-apache-iceberg-materialized-views/

Managing high-volume application logs at scale presents challenges from slow query performance and difficulty running complex aggregations to maintaining real-time analytics on streaming data. Apache Iceberg materialized views with AWS Glue, Amazon Data Firehose, and AWS Lambda address these challenges by accelerating log analytics through pre-computed query results.

In this post, you learn how to build an application log pipeline for production use with Amazon CloudWatch Logs, AWS Lambda, Amazon Data Firehose, AWS Glue, and Apache Iceberg materialized tables. You then use materialized views to accelerate query performance. This solution helps you achieve faster query response times on large-scale log data without requiring you to manage continuous data lake refresh.

Solution overview

This solution accelerates log analytics by pre-computing query results through Apache Iceberg materialized views. By querying pre-aggregated results instead of scanning raw log data for every request, you can help reduce query response times. For example, queries that previously took minutes scanning terabytes of raw data may return in seconds from the compact materialized view. Results update automatically as new logs arrive, helping you handle high-volume log streams while maintaining fast analytics performance.

Architecture overview

The architecture consists of AWS services working together to create a data pipeline:

  • Amazon CloudWatch Logs receives application logs and system events, then routes them to downstream targets using CloudWatch Logs subscription filters. CloudWatch Logs has a built-in retry mechanism. If the destination service returns a retryable error, CloudWatch Logs automatically retries delivery for up to 24 hours.
  • AWS Lambda serves as the transformation layer, parsing log messages, enriching data, and preparing records for storage.
  • Amazon Data Firehose buffers incoming data and handles the technical requirements of writing to Apache Iceberg tables (an open-source data table format), including batch optimization, schema validation, and automatic retry logic for failed writes.
  • Apache Iceberg tables stored in Amazon Simple Storage Service (Amazon S3) provide ACID transaction support, schema evolution capabilities, and efficient query performance. Materialized views are managed tables in the AWS Glue Data Catalog that store precomputed query results in Apache Iceberg format.
  • AWS Glue runs a one-time job during stack creation to provision the Iceberg database, base table, and materialized view structure in the Data Catalog. A second scheduled Glue job refreshes the materialized view by recomputing aggregations from the base table on a configurable interval helping downstream queries through Amazon Athena return up-to-date, pre-aggregated results without scanning raw data.

This architecture is designed to support automatic scaling, serverless infrastructure, error handling that routes failed records to Amazon S3 for analysis and replay, capture of failed Lambda invocations for automatic retry, and real-time monitoring through Amazon CloudWatch metrics.

Prerequisites

Before you deploy the solution, review the following prerequisites.

  • AWS account with necessary permissions to execute an AWS CloudFormation template, run AWS Glue jobs, run queries to verify Iceberg table data using Amazon Athena.
  • Basic familiarity with Boto3 to understand Python code. Foundational understanding of Apache Iceberg concepts.

Solution deployment

The following deployment steps guide you through implementing this solution in your AWS account.

Step 1: Deploy the AWS CloudFormation pipeline stack

You can deploy this solution using an AWS CloudFormation stack. The template handles creating Amazon S3 buckets, uploading AWS Glue and Lambda scripts, provisioning IAM roles, configuring the Firehose delivery stream, and running the Glue job to create the Iceberg database, base table, and materialized view.

Launch the stack in the AWS CloudFormation console. Review the parameters marked REQUIRED and adjust the toggle options (CreateScriptBucket, EnableLakeFormation, CreateSubscriptionLogGroup) based on your environment. Other parameters include preconfigured defaults that you should review for your environment. Choose the CloudFormation stack to deploy resources using the AWS CloudFormation console.

Pipeline stack required parameters view in the AWS CloudFormation console.

Additional pipeline stack required parameters in the AWS CloudFormation console.

Step 2: Test the end-to-end pipeline

Send sample log events matching the Iceberg table schema (for example, id, customer_name, amount, and order_date) to the CloudWatch log group. The subscription filter triggers the Lambda, which forwards records to Firehose for delivery into the Iceberg table.

git clone https://github.com/aws-samples/sample-log-analytics-iceberg-mv.git
cd sample-log-analytics-iceberg-mv
python3 scripts/send_test_logs.py
Terminal output showing the test log event script sending sample records to the CloudWatch log group

Execution of test events.

Verify data delivery and refresh the materialized view

Allow approximately 30 seconds (learn more in Buffer data for dynamic partitioning) for the Firehose buffer to flush. After the buffer flushes, run the following query in Amazon Athena to verify that data has been successfully delivered to the base table.

Query result using Amazon Athena.

Automated materialized view refresh

In this example, the AWS CloudFormation stack provisions a Glue job configured to run the materialized view (MV) refresh once daily at midnight UTC, meaning the MV reflects data up to the previous day. You can adjust the trigger’s cron schedule to match common MV refresh requirements such as hourly, every 15 minutes, or on demand.

The Glue job performs a full recomputation of the aggregations from the base Iceberg table and writes the results to the MV. Downstream consumers querying through Athena read from this pre-aggregated view, delivering faster performance. This is especially critical in real production scenarios where the base table contains millions of records and numerous columns. Computing aggregations directly from raw data at query time would degrade downstream application performance.

Job scheduled view in the AWS Glue console.

In a production environment, the base Iceberg table stores every individual order event, potentially millions of rows with dozens of columns growing daily. When dashboards or downstream applications need aggregated insights like daily revenue per customer or monthly order counts by region, querying the base table directly forces Athena to scan terabytes of raw data on every request. This results in slow response times and high costs at scale. The materialized view solves this by pre-computing these business-level aggregations once during the scheduled refresh, storing the results in a compact, purpose-built table with far fewer rows and columns. This means a dashboard query that would scan millions of raw records now reads from a pre-aggregated table, designed to reduce query response time. The base table remains your source of truth for granular, row-level lookups, while the materialized view serves as the performance layer for repeated analytical queries with embedded business logic.

Materialized View query result using Amazon Athena

Alternative: Amazon S3 Tables

This solution can also be implemented using Amazon S3 Tables, which provides a fully managed Apache Iceberg experience with native support for materialized views. In this post, we use the Glue-based approach to demonstrate the underlying mechanics and provide full flexibility to customize refresh logic for your specific requirements. To learn more, see Getting started with S3 Tables.

Clean up

To avoid incurring future charges, delete the resources you created as part of this exercise if you are not planning to use them further. Delete the stacks created in the previous steps, then empty and delete the Amazon S3 buckets.

Conclusion

This solution shows how to build a scalable application log data pipeline that delivers log events from Amazon CloudWatch Logs to Apache Iceberg tables using AWS Lambda and Amazon Data Firehose. This architecture uses fully managed AWS services to minimize operational overhead while providing high availability and consistent performance.

Key strengths include serverless infrastructure designed to support automatic scaling, error handling designed to route failed records to Amazon S3 for troubleshooting and replay, and analytics capabilities through Apache Iceberg’s ACID transactions and query performance optimizations. As you move this solution into production, we recommend that you implement data quality checks in Lambda and configure encryption at rest and in transit for your data. You can also establish data retention policies and explore partitioning strategies for better query performance.

You now have a log analytics pipeline built for production use that scales with your workload.

Additional resources


About the author

Shinu Tharol

Shinu Tharol

Shinu is a Technical Account Manager at AWS, delivering technical guidance and strategic support to enterprise customers. His expertise includes cloud operations, artificial intelligence, data analytics, and cloud cost optimization, enabling customers to maximize their AWS investments while maintaining operational excellence.

Deploy modern data platforms in minutes with MDAA

Post Syndicated from Sudeshna Dash original https://aws.amazon.com/blogs/big-data/deploy-modern-data-platforms-in-minutes-with-mdaa/

Modern Data Architecture Accelerator (MDAA) is an open source framework that replaces infrastructure code with concise YAML configuration, so your team can deploy a governed, production-ready data architecture, reducing deployment time from months to weeks (depending on complexity and team experience).

Organizations building modern data architecture on AWS face a critical challenge: deploying production-ready, governed infrastructure traditionally requires 6–12 months of custom development, thousands of lines of infrastructure code, and continuous remediation cycles to maintain security and compliance. Governance is often added incrementally, treated as an afterthought that creates compliance gaps and engineering rework.

MDAA addresses this by replacing infrastructure code with concise YAML configuration, achieving up to 97.6 percent code reduction (from approximately 1,800 lines of AWS CloudFormation to 45 lines of MDAA YAML) while embedding governance from the start. The complete Governed Lakehouse Starter Kit deploys 491 AWS resources across 12 stacks from approximately 450 lines of YAML configuration, representing a 66x verbosity ratio where each line automatically expands into production-ready infrastructure.

In this post, we explore how MDAA transforms data architecture development from months of manual coding to production-ready deployment through configuration-driven infrastructure and embedded governance, examine a real customer transformation, and provide a clear implementation pathway for your own data modernization journey.

Customer use case and challenge

A university system office needed to modernize its analytics architecture across 17 campuses while managing sensitive educational data. Their third-party dependency created bottlenecks that slowed feature implementation from weeks to months, and their IT team lacked the cloud skillsets to build modern infrastructure independently.

With MDAA, they achieved:

  • 95 percent reduction in time-to-value for dashboard and feature implementation (from weeks to hours).
  • 17 campuses integrated into a unified, secure architecture.
  • 7.2TB of data and over 8,000 dashboards migrated successfully.
  • Significant cost savings by removing third-party dependencies and reducing license costs.
  • Enhanced security posture for external stakeholders accessing sensitive educational data.

The team used MDAA to implement a modernization strategy with continuous integration and continuous delivery (CI/CD) for automated deployment. The architecture now supports rapid response to stakeholder requests while maintaining strict data governance through AWS Lake Formation.

Their transformation demonstrates what becomes possible when governance is embedded from launch rather than added incrementally, moving from months-long manual development to weeks of production-ready deployment through configuration-driven infrastructure.

Solution: MDAA and its value propositions

MDAA’s capabilities stem from its modular, composable architecture. The accelerator provides over 40 pre-built modules that encapsulate AWS best practices for security, governance, and operational excellence. Organizations describe the outcomes they want in MDAA-specific YAML configuration files (not CloudFormation or Terraform YAML) and the accelerator automatically translates these configurations into AWS Cloud Development Kit (AWS CDK) constructs, which then deploy via CloudFormation with embedded governance.

Configuration over code. The MDAA framework takes a fundamentally different approach: describe the outcomes you want in YAML, and the accelerator deploys production-ready infrastructure with embedded governance. Consider deploying a governed data lake where fraud detection teams need write access to transaction data, while marketing analytics teams require read-only access to customer behavior data. Traditional approaches require over 1,800 lines of CloudFormation across Amazon Simple Storage Service (Amazon S3) buckets, AWS Key Management Service (AWS KMS) keys, AWS Identity and Access Management (IAM) policies, and Lake Formation permissions. With MDAA, the same governed data lake is expressed in 45 lines of configuration, a 97.6 percent reduction, while helping you apply encryption, least-privilege access, and cross-account governance as built-in defaults.

The configuration deploys multi-zone S3 storage with KMS encryption, Lake Formation permissions with tag-based access control (TBAC) enabled, Amazon SageMaker Unified Studio for data product discovery, and encrypted AWS Glue Data Catalog with automated crawlers. All permissions flow through Lake Formation rather than individual IAM policies.

Embedded governance from day one. Governance is declared in YAML and deployed alongside infrastructure from the first run. Fine-grained access controls, encrypted data catalogs, data quality validation, audit trails, and sensitive data classification are all part of the same configuration. MDAA’s Governed Lakehouse starter kit defines an entire governed data architecture in roughly 450 lines of YAML, which produces approximately 29,700 lines of CloudFormation across 12 stacks (a 98.5 percent reduction in infrastructure code).

Modular, composable architecture. Each module is purpose-built to handle a specific capability within the data architecture. Modules communicate through AWS Systems Manager Parameter Store, passing resource identifiers (Amazon Resource Names (ARNs), IDs, and names) between stacks. This approach removes hardcoded dependencies. A KMS key created in one module can be referenced by another through parameter resolution, with all dependencies resolved automatically at deployment time.

The diagram illustrates the deployed architecture and team-level access flow that MDAA generates from the 45-line configuration.

Progressive architecture patterns. MDAA provides four reference architecture patterns that align to progressive stages of data infrastructure maturity:

  • Basic Data Lake deploys a governed data lake with built-in security controls, data quality checks, centralized metadata management using AWS Lake Formation and AWS Glue.
  • Data Science Platform extends the data lake with Amazon SageMaker notebooks, feature stores, and machine learning (ML) pipelines so data science teams can experiment and train models on governed data.
  • SageMaker Unified Studio adds a single interface for analytics and ML collaboration, connecting data engineers, analysts, and data scientists in one workspace.
  • Generative AI Platform layers Amazon Bedrock and Retrieval Augmented Generation (RAG) capabilities on top of your existing data foundation, so teams can build generative AI applications grounded in enterprise data.

Each pattern builds the one before it. You can start with the Basic Data Lake and adopt additional patterns as your team’s needs grow. MDAA’s modular design means you add capabilities without rearchitecting what you already deployed.

The infrastructure is versioned through GitHub, repeatable across environments, and auditable through comprehensive AWS CloudTrail logging. Data engineers focus on data pipelines and business logic while MDAA manages infrastructure complexity and governance integration. This represents the fundamental shift: from writing infrastructure code to describing the outcomes you want through configuration, with governance embedded from the start.

Use case of MDAA: Governed data architecture

DataOps teams spend significant time on governance tasks, including permissions management, compliance validation, and access control, rather than building pipelines and analytics. These aren’t data problems, they’re governance problems that consume engineering capacity meant for higher-value work. MDAA addresses this at the architectural level. Governance is declared in YAML and deployed alongside infrastructure from the first run.

The following sections walk through how each governance module works in practice.

Publish, discover, subscribe, and consume data products between business units: SageMaker Unified Studio

Amazon SageMaker Unified Studio provides a governed data catalog where data producers publish data products, and consumers discover and subscribe to them. Your deployment with MDAA includes a pre-configured domain, blueprints (managed and custom), projects, and environment profiles, all defined in a single configuration file:

# sagemaker.yaml --- 16 lines that deploy 114 CloudFormation resources
domains:
  domain1:
    dataAdminRole:
      id: ssm:/{{org}}/govern1/generated-role/data-admin/id
    description: SMUS Domain 1
    userAssignment: MANUAL

    tooling:
      vpcId: '{{context:vpc_id}}'
      subnetIds:
        - '{{context:private_subnet_id1}}'
        - '{{context:private_subnet_id2}}'

    groups:
      team1:
        ssoId: '{{context:team1-group-sso-id}}'
      team2:
        ssoId: '{{context:team2-group-sso-id}}'

Behind this configuration, MDAA deploys an Amazon SageMaker Unified Studio domain with dedicated KMS keys, execution and provisioning roles, and single sign-on group profiles for team access. Data producers tag and publish assets with metadata, ownership, and classification. Consumers browse a searchable catalog, see only authorized assets, and request access through a governed workflow. Cross-account and cross-business-unit data sharing flows through a subscription model, ensuring every access grant is tracked, auditable, and revocable.

Use case of MDAA: Restricting access to cardholder data using Lake Formation

AWS Lake Formation provides fine-grained access control at database and table levels, removing manual IAM policy management. MDAA deploys AWS Lake Formation with pre-configured settings that disable IAMAllowedPrincipals, the critical governance setting that ensures all permissions flow through centralized governance:

# lakeformation-settings.yaml --- 6 lines that deploy 25 CloudFormation resources
lakeFormationAdminRoles:
  - id: generated-role-id:data-admin
createCdkLFAdmin: true
createDataZoneAdminRole: true
iamAllowedPrincipalsDefault: false

That last flag is the single most important governance setting in the platform. Without it, an IAM principal with glue:GetTable can read tables in the catalog, bypassing the entire access control model. Most manual setups miss this or defer it.

With the data lake configuration, you declare roles and access policies in YAML where admins get full control, engineers get read access to curated data, extract, transform, and load (ETL) roles get scoped write access, and MDAA compiles them into the correct S3 bucket policies and Lake Formation registrations.

Use case of MDAA: Ensuring data integrity with AWS Glue Data Quality

AWS Glue Data Quality runs automated validation rulesets continuously as part of the pipeline, not as periodic batch checks. MDAA’s data quality module supports over 15 built-in rule types, from completeness and uniqueness checks to statistical thresholds and data freshness validation:

# data-quality.yaml
projectName: example-project

rulesets:
  customer-data-quality:
    description: Validate customer data completeness and uniqueness
    targetTable:
      databaseName: project:databaseName/customer-data
      tableName: customers
    ruleset:
      - ruleType: IsComplete
        column: customer_id
      - ruleType: Uniqueness
        column: email
        comparisonOperator: ">"
        threshold: 0.95
      - ruleType: RowCount
        comparisonOperator: ">"
        value: 100

Quality metrics flow into Amazon CloudWatch for real-time alerting. If anomalies are detected, automated workflows quarantine affected records and alert data engineering teams before issues reach downstream consumers.

Protecting metadata at rest: AWS Glue Data Catalog encryption

Table schemas, column names, and partition structures can reveal sensitive information about an organization’s data architecture, even without access to the underlying data. AWS Glue Catalog Encryption secures metadata at rest using AWS KMS-managed keys. MDAA configures catalog encryption by default, so schema definitions and connection passwords are encrypted from initial deployment without requiring manual key management setup. Access to catalog metadata follows the same Lake Formation governance controls applied to the data itself, so teams see only the schemas that they’re authorized to query.

Auditing every data access event: CloudTrail integration

Every data access event must be logged and attributable to a specific identity. Without a complete audit trail, demonstrating compliance during a regulatory review becomes a manual, error-prone process. AWS CloudTrail captures API-level activity across the data infrastructure, recording who accesses what data, when, and from which service. MDAA configures CloudTrail integration by default, so audit logging is active from initial deployment rather than added retroactively. Log data flows into a centralized, tamper-resistant store, giving compliance teams a single location to query access history across all business units and accounts.

Identifying sensitive data automatically: Macie integration

In large environments, sensitive information spreads across dozens of S3 buckets through pipelines, transforms, and ad hoc data drops, and self-reporting data owners consistently produce gaps. Amazon Macie uses machine learning to automatically discover and classify sensitive data in S3, surfacing findings at the object level without manual tagging. MDAA configures Macie across your S3 buckets during deployment, routing findings to Amazon EventBridge where automated workflows can alert owners or trigger remediation.

Together, these controls form a layered defense: Lake Formation governs access to cataloged data, Glue Data Quality validates integrity on arrival, and Macie identifies sensitive data that lands outside governed pipelines to reduce compliance risk.

Multi-account data mesh

MDAA provides extensive support for multi-account data mesh setups, with decentralized data ownership across business units and centralized governance. The data mesh starter kit supports cross-account data product publishing and consumption, allowing organizations to scale data sharing while maintaining consistent security and compliance controls.

Technical implementation

Ready to deploy your modern data architecture? Here are the resources to get started:

MDAA Implementation Guide provides detailed instructions for deploying all starter packages, including architecture patterns, configuration examples, security best practices, and troubleshooting guidance.

MDAA Hands-on Workshop offers step-by-step guided implementation with AWS experts. The workshop covers configuration management best practices, implementation patterns, hands-on labs with real-world scenarios, and cleanup instructions.

GitHub Repository and Documentation provide source code, module reference, and comprehensive documentation.

Organizations approach MDAA from different starting points. Some modernize existing data architectures, migrating from on-premises infrastructure or legacy cloud architectures. Others build new architectures for artificial intelligence and machine learning (AI/ML) initiatives or generative AI applications. Financial services organizations require PCI-DSS compliance from day one. Healthcare organizations need controls that can help support HIPAA. Each journey benefits from MDAA’s configuration-driven approach and embedded governance.

Conclusion

MDAA transforms data architecture development from months of manual coding to production-ready deployment. Configuration-driven infrastructure reduces development time by 40–60 percent while embedding governance from the start. The university system’s 95 percent reduction in time-to-value demonstrates the outcome: organizations deploy secure, compliant, governed data architectures in weeks rather than months.

Financial services organizations can deploy architectures to help them align with PCI-DSS compliance requirements using Lake Formation access controls, Glue Data Quality validation, SageMaker Unified Studio data discovery, comprehensive CloudTrail audit trails, and automated Macie data classification, all inherited from configuration rather than built manually.

Data architecture journeys need not follow six-month timelines with governance added incrementally. MDAA provides an alternative: describe the outcomes you want through YAML configuration, inherit pre-validated security controls, and deploy production-ready infrastructure with comprehensive governance from initial deployment.

Security and compliance is a shared responsibility between AWS and the customer. For more information, see the AWS Shared Responsibility Model.

Need help or have questions? Contact AWS ProServe for personalized guidance on selecting the right package and deployment strategy for your organization.


About the author

Sudeshna Dash

Sudeshna Dash

Sudeshna is a Data Scientist at AWS Professional Services based in Berlin, Germany. She specializes in data architecture, generative AI, and agentic AI systems on AWS. Sudeshna is a contributor to the Modern Data Architecture Accelerator (MDAA) open-source project and helps customers design and deploy governed, production-ready data and AI/ML architectures on AWS.

John Reynolds

John Reynolds is a Principal Engineer with AWS Professional Services based in Seattle, Washington. He leads the architecture and development of Modern Data Architecture Accelerator (MDAA), focusing on turning proven delivery patterns into reusable, production-ready foundations that customers can adopt and extend at scale.

Autonomous troubleshooting for Medallion Architecture with AWS DevOps Agent and Apache Spark Troubleshooting Agent

Post Syndicated from Mohammad Sabeel original https://aws.amazon.com/blogs/big-data/autonomous-troubleshooting-for-medallion-architecture-with-aws-devops-agent-and-apache-spark-troubleshooting-agent/

Every minute of data processing pipeline downtime delays business decisions, stalls downstream analytics, drives revenue loss, and erodes stakeholder confidence. Teams that run Medallion Architecture pipelines—a common data lakehouse pattern where data flows through bronze, silver, and gold layers with increasing quality—face cascading failures that impact revenue-critical reporting and machine learning workloads. As you scale these multi-stage pipelines with Amazon Managed Workflows for Apache Airflow (MWAA), AWS Glue, and Amazon Redshift, troubleshooting failures becomes increasingly complex. When a mission-critical job fails, an engineer must sift through gigabytes of logs across interconnected systems. This means spending hours on incident investigations, examining execution timelines and resource metrics, and cross-referencing findings with Amazon CloudWatch and recent deployment changes to find the root cause. This requires deep familiarity with the underlying technologies, expertise not every team member has. When the right engineer is unavailable during off-hours, pipeline downtime extends and downstream consumers wait. The cycle of detect, investigate, fix, and repeat is costly and entirely reactive. A proactive operational model moves issue identification upstream, catching and addressing problems before they disrupt your data pipelines.

In this post, we show you how to diagnose multi-layer Medallion Architecture pipeline failures in minutes using AWS DevOps Agent with Apache Spark Troubleshooting Agent integrated as an MCP server.

What is AWS DevOps Agent and Apache Spark Troubleshooting Agent?

AWS DevOps Agent is an autonomous investigation agent powered by AI that automatically diagnoses operational issues across your AWS environment. When a failure occurs, the agent independently gathers evidence from logs, metrics, and configurations across interconnected services, identifies the root cause, and delivers actionable remediation steps, all without human intervention. It integrates with your existing workflows through webhooks and delivers findings directly to communication channels like Slack. With AWS DevOps Agent, you can replace the reactive cycle of detect, investigate, fix, and repeat with autonomous, proactive troubleshooting. The agent acts as your always-on, on-call engineer, starting its investigation the moment a failure occurs, whether during business hours or in the middle of the night.

Apache Spark Troubleshooting Agent is an AI-powered, fully managed Model Context Protocol (MCP) server that data engineers can use to diagnose Spark application failures across Amazon EMR, AWS Glue, and Amazon SageMaker AI Notebooks using natural language. It automatically correlates Spark History Server data, distributed executor logs, and configuration patterns to identify root causes and deliver actionable recommendations. This removes hours of manual investigation across multiple consoles and log files.

Use case

The following sections walk through a common Medallion Architecture failure scenario and show how autonomous troubleshooting resolves it.

The scenario

Consider this scenario: a gold layer AWS Glue job fails with “Missing data for not-null field.” The logs don’t reveal the actual problem. The root cause is a subtle data quality issue introduced upstream in the silver layer, a job that succeeded without errors. Without autonomous troubleshooting, you would manually trace data lineage across Amazon Simple Storage Service (Amazon S3), Amazon Redshift, and multiple AWS Glue job logs to find the source.

The solution

When integrated with the Apache Spark Troubleshooting Agent, AWS DevOps Agent identifies the gold layer Amazon Redshift write failure, traces it back to silver layer data corruption, and provides detailed root causes and actionable recommendations. The investigation typically completes within 3 to 5 minutes.

Solution overview

The following diagram shows the Medallion Architecture data flow across bronze, silver, and gold layers.

Medallion Architecture data flow showing the bronze layer in Amazon S3, the silver layer in Amazon S3 and Amazon Redshift, and the gold layer in Amazon Redshift, with Amazon MWAA orchestrating AWS Glue jobs and AWS DevOps Agent investigating failures

The architecture flow includes the following steps:

  1. Amazon MWAA triggers the Medallion pipeline directed acyclic graph (DAG), orchestrating three AWS Glue jobs sequentially: bronze layer, silver layer, and gold layer.
  2. The bronze layer job generates 50,000 synthetic ecommerce order records and writes raw Parquet files to Amazon S3.
  3. The silver layer job reads bronze data from Amazon S3, applies transformations, and writes the results to two destinations in parallel: Amazon S3, and Amazon Redshift (filtered, cleaned, and augmented data in the silver_ecommerce table). This job silently introduces data corruption in approximately 8 percent of total_amount values.
  4. The gold layer job reads from the Amazon Redshift silver_ecommerce table, performs aggregation, and attempts to write business-level aggregates back to the Amazon Redshift gold_ecommerce_summary table. If upstream data corruption introduces NULL values, this job fails with “Missing data for not-null field” because those NULL values violate the NOT NULL constraint.
  5. When the gold layer job enters a FAILED state, Amazon EventBridge captures the AWS Glue Job State Change event and invokes an AWS Lambda function. The Lambda function retrieves webhook credentials from AWS Secrets Manager, constructs an HMAC-signed event payload containing the job name, run ID, and error details, and sends it to AWS DevOps Agent.
  6. AWS DevOps Agent receives the HTTP POST request to the webhook and starts an autonomous investigation. It authenticates with Amazon Cognito using the OAuth 2.0 client credentials flow, then sends an MCP request through Amazon Bedrock AgentCore Gateway. The AgentCore Gateway invokes a Signature Version 4 (SigV4) Proxy Lambda, which signs the request and forwards it to the Apache Spark Troubleshooting Agent MCP Server. The MCP Server analyzes Spark event logs, executor metrics, and error stack traces for the failed gold job.
  7. AWS DevOps Agent delivers the investigation to your configured Slack channel. The delivery includes root cause analysis, upstream data lineage back to the silver layer corruption, and step-by-step remediation recommendations.

Walkthrough

In the following sections, you deploy a three-layer Medallion Architecture pipeline that processes ecommerce order data. Complete the steps to get started with autonomous troubleshooting using AWS DevOps Agent.

Prerequisites

Before you begin, verify that you have the following:

  • An AWS account. Your AWS Identity and Access Management (IAM) user or role must have the following permissions:
    • iam:CreateRole, iam:AttachRolePolicy, iam:PutRolePolicy
    • lambda:CreateFunction, lambda:AddPermission
    • glue:CreateJob, glue:StartJobRun
    • redshift:CreateCluster, redshift:GetClusterCredentials
    • airflow:CreateEnvironment
    • events:PutRule, events:PutTargets
    • sqs:CreateQueue
    • secretsmanager:CreateSecret
    • kms:CreateKey
    • ec2:CreateVpc, ec2:CreateSubnet, ec2:CreateSecurityGroup
    • cloudformation:CreateStack, cloudformation:DescribeStacks
    • Alternatively, you can use the AdministratorAccess managed policy for simplicity in a dev/test environment.
  • AWS Command Line Interface (AWS CLI) version 2.30.0 or later, installed and configured with appropriate credentials.
  • (Optional) A Slack workspace if you want investigation results delivered to a channel.

Set up AWS DevOps Agent

In this section, you configure AWS DevOps Agent to receive and investigate pipeline failure events. This involves three tasks: creating an Agent Space (your investigation workspace), optionally connecting a Slack channel for notifications, and generating a webhook endpoint that your pipeline uses to send failure alerts to the agent.

Create an Agent Space

  1. Open the AWS DevOps Agent console.
  2. Choose Create Agent Space.
  3. Enter a name (for example, medallion-troubleshooting).
  4. Choose Create.

Connect Slack integration (optional)

If you use Slack for internal communication, you can configure it to receive investigation results.

  1. In the AWS DevOps Agent console, go to Agent Spaces, select medallion-troubleshooting and then Communications.
  2. Choose Add integration and choose Slack.
  3. Choose Next to allow AWS DevOps Agent to access your Slack workspace, and choose Allow.
  4. Provide the Slack workspace and the Channel ID where you want investigation results delivered, then choose Next.
  5. Enter the following command in your channel chat to complete the integration: /invite @AWS DevOps Agent.
    • While running this command, when prompted, choose the correct region where the Agent Space is provisioned.

Create a webhook

  1. In your Agent Space, go to Webhooks.
  2. Choose Add webhook and choose Next on the two following pages.
  3. Choose Generate URL and secret key, and give the webhook a name (for example, medallion-failure-webhook).
  4. After creation, copy and save the Webhook URL (HTTPS endpoint) and Secret Key. You can also choose Download .csv to save this information to a secure location. Select the checkbox labeled I’ve saved and stored my URL and secret key, then choose Add.

Note the Webhook URL and Secret Key for later. You provide them as parameters when you create the AWS CloudFormation stack.

Deploy the AWS CloudFormation stack

The AWS CloudFormation template deploys the full Medallion Architecture pipeline. This includes an Amazon Virtual Private Cloud (Amazon VPC) with private subnets, an Amazon Redshift cluster (ra3.xlplus, single-node), and three AWS Glue jobs. It also creates an Amazon MWAA environment, Amazon EventBridge rules, AWS Lambda functions, and an AgentCore Gateway with Amazon Cognito OAuth authentication.

You can deploy the stack using one of two methods. Use Option A if you prefer a visual, guided experience through the AWS Management Console. Use Option B if you prefer working from the command line or need to integrate the deployment into a script or automation workflow.

Before you start, download the CloudFormation template from GitHub.

Option A: AWS Management Console (recommended)

  1. Open the AWS CloudFormation console and choose Create stack → With existing resources (import resources) or Upload a template file.
  2. Choose Choose file, select the downloaded blog-medallion-stack.yaml, then choose Next.
  3. For Stack name, enter medallion-troubleshooting.
  4. Fill in the parameters:
    • For WebhookUrl, enter your AWS DevOps Agent webhook URL (from Agent Space settings).
    • For WebhookSecret, enter the webhook secret for authentication.
  5. Choose Next, select I acknowledge that AWS CloudFormation might create IAM resources with custom names, then choose Submit.

Option B: AWS CLI

aws cloudformation create-stack \
    --stack-name medallion-troubleshooting \
    --template-body file://blog-medallion-stack.yaml \
    --parameters \
        ParameterKey=WebhookUrl,ParameterValue=<YOUR-WEBHOOK-URL> \
        ParameterKey=WebhookSecret,ParameterValue=<YOUR-WEBHOOK-SECRET> \
    --capabilities <CAPABILITY_NAMED_IAM> \
    --region <YOUR-REGION>

Replace the placeholder values:

  • YOUR-WEBHOOK-URL – Your AWS DevOps Agent webhook URL (from Agent Space settings).
  • YOUR-WEBHOOK-SECRET – The webhook secret for authentication.
  • YOUR-REGION – The AWS Region.

Wait for the stack status to show CREATE_COMPLETE. In our testing, this took approximately 30–40 minutes.

Retrieve Amazon Cognito client credentials

After the stack is deployed, it creates an Amazon Cognito user pool with an OAuth 2.0 client for AWS DevOps Agent authentication. Retrieve the client secret using the command below. The --user-pool-id  and CognitoClientId needs to be copied from the stack outputs.

aws cognito-idp describe-user-pool-client \
    --user-pool-id <UserPoolId-from-outputs> \
    --client-id <CognitoClientId-from-outputs> \
    --query UserPoolClient.ClientSecret \
    --output text --region <YOUR-REGION>

Replace YOUR-REGION with the actual AWS Region value, and save this value for the MCP Server registration in the following step.

Register the Spark Troubleshooting MCP Server

The Spark Troubleshooting MCP Server gives AWS DevOps Agent the ability to analyze Apache Spark event logs, executor metrics, and error stack traces from your AWS Glue jobs. By registering this server, you connect the agent to the diagnostic tooling it needs to autonomously investigate pipeline failures.

To register the MCP Server in AWS DevOps Agent, complete the following steps:

  1. In the AWS DevOps Agent console, go to Agent Spaces, select medallion-troubleshooting and then Capabilities.
  2. In the MCP Servers section, choose Add or Add Source.
  3. Find New MCP Server Registration and choose Register.
  4. For Name, enter sparkagent.
  5. For Endpoint URL, enter the AgentCoreGatewayUrl value from the stack outputs.
  6. For Description, enter Apache Spark Troubleshooting MCP Server via AgentCore Gateway.
  7. Leave Enable Dynamic Client Registration cleared.
  8. Leave Connect to endpoint using a private connection cleared, then choose Next.Registration page for the Apache Spark Troubleshooting MCP Server in the AWS DevOps Agent console, showing endpoint URL and description fields
  9. Under Authorization Flow, select OAuth Client Credentials, and choose Next.
  10. For Client ID, enter the CognitoClientId value from the stack outputs.
  11. For Client Secret, enter the value you retrieved in the preceding step.
  12. For Exchange URL, enter the CognitoTokenEndpoint value from the stack outputs.
  13. For Add Scope, enter <stack-name>-mcp-proxy/invoke. For example, medallion-troubleshooting-mcp-proxy/invoke.
  14. Choose Next, review your configuration, and choose Add.
  15. Once you choose Add, on the following screen, click on the checkbox next to the spark___analyze_spark_workload. This is the root cause analysis tool which provides detailed troubleshooting for failed Apache Spark workloads.
    Selecting the tool within the AWS Managed Apache Spark Troubleshooting MCP server
  16. Choose Save as a last step. You will see the MCP Server associated successfully message on the top.
    Confirmation showing the successful Integration of AWS DevOps Agent Space with Apache Spark Troubleshooting MCP Server

See AWS DevOps Agent in action

Now that you have completed the prerequisites, you can see AWS DevOps Agent in action. Go to the Amazon MWAA Airflow Environments UI and click on Open Airflow UI under Airflow UI. It will open in a new browser tab. In the Airflow console, locate and manually trigger the medallion_architecture_pipeline DAG.

Amazon MWAA Airflow console showing the medallion_architecture_pipeline DAG with the Trigger DAG action selected

Amazon MWAA Airflow UI showing the medallion_architecture_pipeline DAG with bronze, silver, and gold tasks listed sequentially

The DAG runs three AWS Glue jobs sequentially:

  1. Bronze layer – This job generates 50,000 ecommerce order records and writes them to Amazon S3 as Parquet files.
  2. Silver layer – This job applies transformations and loads the results to both Amazon S3 and Amazon Redshift. It also silently injects approximately 8 percent of total_amount values with $ prefix strings, introducing hidden data corruption.
  3. Gold layer – This job reads from Amazon Redshift, casts total_amount to numeric (producing NULL values for the $-prefixed strings), and attempts to write aggregated results to the Amazon Redshift target table. It fails because the NULL values violate the NOT NULL constraint on revenue_total.

Amazon MWAA DAG run showing the bronze task succeeded, the silver task succeeded, and the gold task failed

With the components deployed and connected, the autonomous troubleshooting pipeline is ready to respond to failures. In this walkthrough, the silver layer job deliberately introduces data corruption to simulate a real-world data quality issue. This causes the gold layer job to fail, giving you the opportunity to see how AWS DevOps Agent responds.

As soon as the gold layer job fails, AWS DevOps Agent starts an autonomous investigation and uses the Apache Spark Troubleshooting MCP Server where needed.

Go to the AWS DevOps Management console and choose the medallion-troubleshooting under Agent Spaces. Next, select the Operator Access button. This will redirect you to Operator Console where you will see that the incident investigation automatically started in 1-2 minutes post Gold layer job failure.

After the investigation completes, AWS DevOps Agent presents its findings within the incident analysis. The results are organized into two sections.

Root cause identified by AWS DevOps Agent

The agent identifies the underlying cause of the failure, tracing the gold layer write error back to data corruption introduced in the upstream silver layer AWS Glue job.

Root cause analysis from AWS DevOps Agent showing the gold layer write error traced back to silver layer data corruption

Mitigation plan generated by AWS DevOps Agent

On choosing Generate Mitigation Plan, the agent provides step-by-step remediation recommendations to resolve the issue and prevent recurrence.

Mitigation plan from AWS DevOps Agent listing remediation steps to fix the silver layer data corruption and prevent recurrence

AWS DevOps Agent sends a notification to Slack

Slack channel showing the AWS DevOps Agent investigation summary with root cause identification and upstream data lineage trace

Typically, within 3–5 minutes, the agent delivers a detailed investigation in Slack that includes root cause identification, upstream data lineage tracking, and an actionable recommendation.

You have deployed an autonomous troubleshooting pipeline for Medallion Architecture data pipelines. The pipeline runs using AWS Glue, Amazon Redshift, and Amazon MWAA, with AWS DevOps Agent providing autonomous investigation. The agent traced a gold layer Amazon Redshift write failure back to a silver layer data quality issue. This type of diagnosis would typically require hours of manual investigation by an engineer with deep expertise in Apache Spark, Amazon Redshift, and data pipeline architecture. AWS DevOps Agent completed it autonomously within minutes.

If you need human assistance, you can use the Ask for human support feature within AWS DevOps Agent to open a case with AWS Support, automatically populated with relevant investigation context.

Enhanced investigations with AWS DevOps Agent Skills

AWS DevOps Agent autonomously investigates failures out of the box. You can enhance its diagnostic depth using Skills, a feature that provides the agent with domain-specific guidance tailored to your environment.

For Medallion Architecture pipelines, you can create Skills that instruct the agent to check for data type mismatches between pipeline layers when Amazon Redshift COPY errors occur, cross-reference silver layer data quality metrics with gold layer aggregation failures, or follow your internal runbook for escalating data quality issues to the upstream data engineering team.

To configure Skills, go to your Agent Space in the AWS DevOps Agent console and choose the Skills tab.

Clean up

To avoid incurring future charges, delete the resources you created during this walkthrough promptly after you finish testing.

To clean up resources, complete the following steps:

  1. Deregister the MCP Server. In the AWS DevOps Agent console, go to your Agent Space and choose the Capabilities tab. In the MCP Servers section, choose the sparkagent server, then choose Deregister.
  2. Delete the webhook. In your Agent Space, go to the Webhooks tab. Choose the medallion-failure-webhook, then choose Delete.
  3. Empty the Amazon S3 buckets. Open the Amazon S3 console. Locate the buckets created by the stack (their names start with medallion-troubleshooting). For each bucket, choose Empty, enter permanently delete to confirm, and choose Empty.
  4. Delete the AWS CloudFormation stack. Open the AWS CloudFormation console. Choose the medallion-troubleshooting stack, then choose Delete. Alternatively, run the following command:
aws cloudformation delete-stack \
    --stack-name medallion-troubleshooting \
    --region <your-region>

Wait for the stack deletion to complete.

  1. Delete any retained Amazon S3 buckets. Some Amazon S3 buckets might have a DeletionPolicy of Retain and aren’t automatically deleted with the stack. Return to the Amazon S3 console, locate any remaining buckets created by the stack, empty them using the process in the preceding step, and then choose Delete for each bucket.

Conclusion

In this post, you deployed an autonomous troubleshooting pipeline for Medallion Architecture data pipelines using AWS Glue, Amazon Redshift, Amazon MWAA, and AWS DevOps Agent. The agent traced a gold layer Amazon Redshift write failure back to a silver layer data quality issue—a diagnosis that would typically require hours of manual investigation by an engineer with deep expertise across multiple services.

As your data pipelines grow in complexity, so does the challenge of diagnosing failures that span multiple layers and services. AWS DevOps Agent reduces your mean time to resolution by autonomously investigating incidents the moment they occur, whether during business hours or at 2 AM. Your on-call engineers spend less time sifting through logs and more time building reliable data infrastructure. By shifting from reactive firefighting to autonomous, proactive troubleshooting, you can improve pipeline reliability, protect downstream analytics and machine learning workloads, and maintain stakeholder confidence in your data platform.

To learn how to structure Agent Spaces for investigation accuracy, scope resource access, and use infrastructure as code to streamline deployment, see Best practices for deploying AWS DevOps Agent in production. To learn how to evaluate and choose the right lakehouse pattern for your needs, see Navigating architectural choices for a lakehouse using Amazon SageMaker. For more about Apache Spark Troubleshooting Agent, see Introducing the Apache Spark Troubleshooting Agent for Amazon EMR and AWS Glue.

Next steps

Now that you have set up autonomous troubleshooting for your Medallion Architecture pipeline, consider exploring the following:


About the authors

Mohammad Sabeel

Mohammad Sabeel

Mohammad is a Senior Technical Account Manager (TAM) at Amazon Web Services (AWS) with over 14 years of experience in Information Technology (IT). As a member of the Technical Field Community for Analytics team, he is a subject matter expert in Analytics services including AWS Glue, Amazon Managed Workflows for Apache Airflow (MWAA), and Amazon Athena. Sabeel provides strategic guidance and proactive technical support to enterprise and ISV customers, helping them optimize their data analytics solutions, build resilient architectures, and accelerate cloud adoption. With deep subject matter expertise, he enables organizations to build scalable, efficient, and cost-effective data processing pipelines.

Ishan Gaur

Ishan Gaur

Ishan is a Principal Cloud Engineer at AWS. He has worked in the Analytics domain for the last 17 years, now focused on data analytics, AI/ML operations, and proactive cloud optimization. He works with enterprise customers to design resilient data pipelines, automate incident response, and adopt GenAI-powered services and operational tools. He’s passionate about turning reactive support patterns into proactive, self-healing architectures.

Modernizing financial analytics with Amazon SageMaker Unified Studio

Post Syndicated from Umang Aggarwal original https://aws.amazon.com/blogs/architecture/modernizing-financial-analytics-with-amazon-sagemaker-unified-studio/

Avanse Financial Services is one of India’s leading education loan providers. Their Data Engineering Team had built a data lake on AWS using Amazon Simple Storage Service (Amazon S3), Amazon Athena, and AWS Glue for data ingestion and processing. However, their analytics and reporting layer ran on an external analytics application that wasn’t integrated with AWS. Data had to be copied from Amazon S3 into this external application before analysts could run any report, its license consumed a significant portion of their budget despite low utilization, and every integration with AWS services required custom-built pipelines.

After evaluating their options, Avanse migrated to a cloud-native lakehouse architecture using Amazon SageMaker Unified Studio, which unified their data engineering, analytics, and artificial intelligence (AI) workflows in a single governed environment on AWS. In this post, we walk through their migration journey so you can adapt their approach to your own environment.

Why Avanse chose to modernize

The separation between their AWS data lake and their external analytics application created five problems:

  1. Daily data synchronization bottleneck. Every report required a 4-hour batch copy from Amazon S3 into the external analytics application before analysts could query it. Business decisions were based on data that was at least a day old.
  2. Fixed licensing costs disconnected from usage. The external analytics application charged an annual fee regardless of how many queries analysts ran. Avanse needed usage-based pricing that matched what they actually consumed, not a fixed fee for capacity they weren’t using.
  3. Limited auditability. The external analytics application ran on a shared server where different business units (risk, collections, portfolio management) shared the same resources. It lacked granular audit trails, making it difficult to trace who accessed what data and when, or to allocate costs per team.
  4. No centralized data discovery. Although AWS Glue Data Catalog managed schema metadata for the data lake, the external analytics application couldn’t access it. Analysts working in that application relied on folder structures and manual documentation to find the right datasets, slowing onboarding and increasing the risk of using outdated data.
  5. Disconnected from AWS services. The external analytics application couldn’t query data in Amazon S3 or use AWS Glue catalogs natively. Every data flow required connectors and custom-built pipelines, adding maintenance overhead.

Additionally, some datasets were stored on Network File System (NFS) storage outside of Amazon S3, creating another data silo that needed to be consolidated.

Avanse chose Amazon SageMaker Unified Studio because it addressed all five challenges: direct querying of data in Amazon S3 avoiding synchronization, usage-based compute through Amazon Athena and Amazon EMR Serverless, project-based isolation with per-project billing, lineage tracking with AWS IAM Identity Center, and native integration with their existing AWS services.

Solution overview

The core architectural change was moving from a two-application model to a single integrated stack:

Previous architecture
Avanse’s data ingestion and processing ran on AWS (Amazon S3, AWS Glue, Athena), but analytics and reporting ran on an external analytics application. Data had to be batch-copied from Amazon S3 into this external application daily before analysts could query it. Each system had its own access controls, and there was no shared catalog or lineage tracking between them.
New architecture
Analytics now run directly against data in Amazon S3 through Amazon SageMaker Unified Studio. There’s no data copy step. Analysts query the same data that the ingestion pipelines produce, using Athena for SQL and EMR Serverless for large-scale processing. Governance, access control, and lineage are centralized through IAM Identity Center and SageMaker Catalog.

The following diagram illustrates the target architecture. It follows a lakehouse pattern, storing data in open formats on Amazon S3 while maintaining ACID transaction support for the consistency financial regulators expect.

Three-layer lakehouse architecture for Avanse on AWS, showing the data layer with Amazon S3 and AWS Glue Data Catalog, the compute layer with Amazon SageMaker Unified Studio, AWS Glue ETL, AWS Lambda, Amazon EMR Serverless, Amazon SageMaker AI, and Amazon Bedrock, and the governance layer with AWS IAM Identity Center, SageMaker Catalog, and Amazon DataZone

The architecture has three layers:

  1. Data layer – Amazon S3 stores data in open formats (Parquet, Delta Lake) with S3 Intelligent-Tiering for automatic cost optimization. AWS Glue Data Catalog maintains schema metadata, making data discoverable across tools.
  2. Compute layer – Amazon SageMaker Unified Studio provides project-based workspaces organized by business function. Collections uses the built-in SQL Query Editor powered by Athena, Risk Reporting uses JupyterLab for interactive analysis, and MIS runs large-scale Spark jobs through Amazon EMR Serverless. AWS Glue ETL handles data transformations and AWS Lambda provides event-driven triggers for report generation. For machine learning (ML) workloads, Amazon SageMaker AI supports model training and deployment, with Amazon Bedrock available for generative AI capabilities such as enhancing risk narratives.
  3. Governance layer – IAM Identity Center provides SSO and audit logging across workspaces. SageMaker Catalog serves as the business glossary with data lineage tracking and access controls. Amazon DataZone connects components through a common metadata layer.

Migration journey

Avanse followed a five-phase approach. The timelines can be adapted to your environment, but the systematic progression from validation through production deployment is key.

Phase 1: Technical validation (72-hour workshop)

Avanse started with a focused 72-hour workshop using isolated SageMaker environments where developers could experiment without impacting production. Their team tested SQL analytics against existing Athena tables and validated that Python and PySpark could replicate their existing analytics workflows.

The team confirmed that querying data directly in Amazon S3 addressed their synchronization bottleneck entirely. The 4-hour daily data copy was no longer necessary, which validated the migration approach.

Phase 2: Data migration and storage optimization

Avanse migrated datasets from NFS storage and legacy analytics formats into Amazon S3, consolidating the data into a single location. They implemented S3 Intelligent-Tiering, which automatically moves data between access tiers based on usage patterns, optimizing costs without impacting retrieval performance.

They replaced legacy analytics connectors with native Athena workgroups within SageMaker Unified Studio, avoiding data synchronization entirely. Source data remained in Amazon S3, queryable by both Athena SQL and SageMaker notebooks, establishing a single source of truth.

Phase 3: Compute modernization

Avanse moved from a shared analytics server to project-based isolation in SageMaker Unified Studio. Each business function (Risk Reporting, Collections, MIS) received its own project with dedicated compute spaces running JupyterLab. Project-specific IAM execution roles provided access controls and cost allocation per business unit.

A single browser-based URL with multi-factor authentication (MFA) now provides access to SQL analytics using the built-in query editor, ML development in JupyterLab notebooks, and big data processing through Amazon EMR Serverless. This replaced the need for local analytics client installations.

Phase 4: Governance implementation

Avanse deployed SageMaker Catalog as their central business data catalog. Analysts now discover approved datasets through semantic search rather than navigating folder structures or relying on manual documentation. They mapped technical Athena table names to business terms. For example, analysts search for “collection efficiency” and find the relevant tables with descriptions, schemas, and lineage.

Lineage capture traces each metric in risk reports back to source tables, transformations, and intermediate datasets. Every action (notebook execution, SQL query, data access) is tied to IAM Identity Center users, creating the comprehensive audit trail their compliance team needed.

Phase 5: Use case migration

Rather than attempting a big-bang migration, Avanse moved critical workflows one at a time:

Portfolio MIS (Monthly/Fortnightly)
Previously required the daily 4-hour data copy from Amazon S3 into the external analytics application before report generation could begin. Avanse avoided the data synchronization step entirely and now generates MIS reports by querying existing Athena tables directly in Amazon S3. Because the source data was already on AWS, there was no need to involve the external application for this activity. Report generation dropped from hours to under 30 minutes.
Collection Efficiency and Bounce Calculation
Ported complex legacy analytics procedures for calculating metrics like collection efficiency and bounce rates to event-driven processing using AWS Glue ETL, AWS Lambda, and PySpark jobs for high-volume data aggregation. The serverless execution model charges only for compute time consumed.
EDW Risk Reporting
Large-scale regulatory joins of Enterprise Data Warehouse assets previously ran as legacy scheduled procedures. These now run as SQL queries in the SageMaker Unified Studio query editor, where analysts execute them on-demand or schedule them through Athena workgroups. The distributed query engine handles complex multi-table joins spanning millions of rows.
Scorecard Generation
Model building shifted from the external analytics application to SageMaker AI workflows. Data scientists use JupyterLab with Python libraries and deploy models directly to SageMaker endpoints, avoiding data movement between separate environments.

Overcoming technical challenges

One technical challenge was code migration. Avanse’s analytics code base contained years of accumulated proprietary scripts and procedures. Direct line-by-line translation was not practical. Instead, they took a pragmatic approach: basic data transformations moved to SQL in Athena, complex business logic was rewritten in PySpark for scalability, and statistical procedures were replaced with Python libraries like pandas and scikit-learn. The approach was to focus on what the code accomplishes, then implement it using cloud-native patterns.

The other technical challenge was performance validation. The team needed to confirm that querying data in Amazon S3 would deliver acceptable performance compared to the external analytics application’s in-memory processing. Queries against Parquet-formatted data in Amazon S3 using Athena delivered comparable performance for standard reporting workloads, while avoiding the 4-hour daily data synchronization step entirely. For large-scale regulatory joins spanning millions of rows, Amazon EMR Serverless provided distributed Spark processing that completed in minutes rather than the hours required in the external application.

Key outcomes

Area Result
Licensing costs Avoided external analytics application fees entirely
Storage costs Reduced through S3 Intelligent-Tiering, which automatically moves data between access tiers based on usage patterns
Report generation From over 4 hours (including data synchronization from Amazon S3 to the external analytics application) to under 30 minutes with direct Amazon S3 querying
Compliance audits From weeks of manual investigation to days with automated lineage reports
Compute costs Usage-based serverless model replaced always-on external analytics infrastructure
Collaboration Unified browser-based environment for data scientists, analysts, and engineers

“By adopting SageMaker Unified Studio, we as the Data Team eliminated legacy licensing costs, reduced storage and compute expenses with a serverless, usage-based model, and accelerated our periodic report generation. At the same time, we transformed compliance and collaboration by cutting audit timelines while unifying our teams in a single, efficient data environment.” – Komal Thakkar, AVP – Lead, Data Engineering, Avanse Financial Services

Best practices

Based on their experience, Avanse recommends:

  • Start with a workshop. Validate your specific use cases in a 72-hour technical validation before committing to full migration.
  • Migrate use cases, not code. Focus on what your analytics accomplish, then implement using cloud-native patterns rather than translating legacy scripts line by line.
  • Invest in governance early. Implement the data catalog and lineage tracking from day one.
  • Embrace project-based isolation. Organize around business functions for clear cost allocation and security boundaries.
  • Document business logic. Use migration as an opportunity to capture undocumented knowledge in the business glossary and dataset descriptions.

Conclusion

Avanse’s migration from an external analytics application to Amazon SageMaker Unified Studio consolidated their analytics stack into a single integrated environment on AWS. By querying data directly in Amazon S3 instead of copying it into the external application, they alleviated their biggest operational bottleneck. Project-based isolation replaced a shared server model, giving each business unit independent compute and clear cost visibility. And centralized governance through SageMaker Catalog and IAM Identity Center gave their compliance team the audit trails they had been missing.

The serverless, usage-based model means Avanse no longer pays for idle capacity. The lakehouse architecture supports new analytics patterns as they emerge, and native integration with AWS services, including generative AI through Amazon Bedrock, positions them to adopt new capabilities as their needs evolve.

Next steps

Start your analytics modernization journey by scheduling a 72-hour technical validation workshop. Contact your AWS account team to discuss your migration approach.

For more information, see:

Beyond JSON blobs: Implementing the VARIANT data type in Apache Iceberg V3

Post Syndicated from Arun Shanmugam original https://aws.amazon.com/blogs/big-data/beyond-json-blobs-implementing-the-variant-data-type-in-apache-iceberg-v3/

Apache Iceberg V3 introduces the VARIANT data type. VARIANT provides data engineers with a high-performance, native solution for managing semi-structured data within the data lake. Consider a massive fleet of IoT sensors: street-level temperature probes, air quality monitors, and vehicle telemetry. Each device emits data in unique JSON structures that constantly evolve with firmware updates.

Historically, engineers were forced to store these payloads as STRING blobs. This legacy approach mandates expensive CPU-intensive parsing at runtime and inflates storage costs with redundant raw text. VARIANT solves these inefficiencies by employing a shredded, binary-encoded format. This allows query engines to skip irrelevant data and access specific nested fields with columnar speed, effectively bridging the gap between the flexibility of JSON and the performance of a structured schema.

VARIANT is stored in Parquet as a three-part group: binary metadata (type and dictionary info), a binary value (the full variant for fallback), and a typed_value group where individual JSON fields are shredded into separate Parquet columns. When you query a specific field, Spark prunes the typed_value group to include only the requested sub-columns. It always retains metadata and the value fallback, so it avoids reading the entire document. This approach delivers two concrete benefits:

  • Reduced query processing time: Queries access only the fields they need without deserializing entire JSON documents. This reduces the amount of data scanned and the time spent on deserialization.
  • Lower storage footprint: Binary encoding compresses more efficiently than raw text, reducing storage costs.

Fields inside the JSON become individually accessible columns under the hood. A query that needs one value out of a deeply nested document no longer must read and deserialize the entire thing. You maintain schema flexibility while gaining the performance characteristics of structured columnar storage.

This post is part 1 of a two-part series. We walk through the basics: creating an Iceberg V3 table with a VARIANT column, inserting semi-structured data, and querying it with variant_get(). In Part 2, we scale to millions of rows and benchmark VARIANT against traditional string storage. We measure the difference in query performance and storage footprint.

Solution overview

This walkthrough demonstrates an end-to-end workflow for working with semi-structured data using the VARIANT data type in Apache Iceberg V3 on Amazon EMR Serverless. Raw JSON payloads are ingested and converted to binary VARIANT format using parse_json(). The data is stored in an Iceberg V3 table where the engine shreds the structure into columnar Parquet sub-columns. You can then query the data efficiently using variant_get() to extract specific fields without deserializing the entire document. AWS Glue Data Catalog manages the table metadata. Amazon Simple Storage Service (Amazon S3) provides the underlying storage.

Note: Check the Apache Iceberg documentation for the latest information on specification status and engine compatibility. Additionally, Fine-Grained Access Control (FGAC) through AWS Lake Formation is not currently supported for the VARIANT data type.

How VARIANT works

When you insert a JSON document into a VARIANT column, Spark converts it from a JSON string into the Variant binary format. During writes, the engine can shred the structure. It extracts individual fields and stores them as native Parquet-typed sub-columns within the VARIANT column’s typed_value group. Fields that are not shredded remain in the binary value column as a fallback. This is conceptually similar to how a columnar table stores each column independently. The difference is that the sub-columns live within a single VARIANT column, and the engine handles the shredding schema automatically.

At query time, when you ask for a specific field using variant_get(), Spark reads only the sub-column that contains that field. It does not need to load or parse the rest of the document. For workloads that repeatedly query a handful of fields out of large, complex JSON payloads, this can significantly reduce the amount of data scanned. It also reduces the time spent deserializing it.

The variant_get() function uses JSON path syntax to navigate the structure. You can extract scalar values with an explicit type (optional), access nested objects, and reach into arrays by index. The function signature is the following.

variant_get(column, '$.path.to.field', 'type')

Where column is the VARIANT column name, the second argument is a JSON path expression, and the optional third argument specifies the expected return type (such as 'string', 'int', or 'double'). When the type argument is omitted, the function returns a VARIANT value that preserves the original encoding.

Running Iceberg V3 on Amazon EMR Serverless

Amazon EMR Serverless 8.0 ships with Apache Spark 4.0.1, which includes native support for Iceberg V3 and the VARIANT data type. You do not need to install additional libraries or configure custom JARs. Amazon EMR Serverless manages the compute infrastructure and scales resources up and down based on workload demand. You can focus on the data rather than the cluster.

While this post uses Amazon EMR Serverless, Iceberg V3 VARIANT support is also available on Amazon EMR on EC2 and Amazon EMR on EKS. You can choose the deployment model that fits your environment.

Getting started

The following walkthrough creates an Iceberg V3 table with a VARIANT column, inserts a set of IoT sensor events, and runs queries to extract fields from the semi-structured payload. Each step includes the code you need to run it on Amazon EMR Serverless.

Prerequisites

Before you begin, verify you have the following:

  • An AWS account with permissions to create Amazon EMR Serverless applications and access Amazon Simple Storage Service (Amazon S3).
  • An Amazon S3 bucket for storing Iceberg table data and scripts.
  • AWS Glue Data Catalog configured for metadata management.
  • An IAM execution role with permissions for Amazon EMR Serverless, Amazon S3, AWS Glue, and Amazon CloudWatch Logs.
  • AWS Command Line Interface (AWS CLI) installed and configured.Note: Running this solution in your AWS account might incur charges for Amazon EMR Serverless, Amazon S3, and AWS Glue. Refer to the respective pricing pages for cost details.

Step 1: Initialize a Spark session with Iceberg V3

Start by creating a Spark session configured to use the Iceberg catalog backed by AWS Glue. The key settings are the Iceberg Spark extensions and the AWS Glue catalog implementation. Replace <YOUR_S3_BUCKET> with your bucket name.

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, parse_json

spark = SparkSession.builder \
    .appName("IcebergV3VariantDemo") \
    .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.warehouse",
            "s3://<YOUR_S3_BUCKET>/warehouse/") \
    .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") \
    .getOrCreate()

When running on Amazon EMR Serverless, some Spark configurations might be set at the application or job level. The configuration shown here is included in the script for completeness. Depending on your Amazon EMR Serverless application settings, you might not need to specify all these properties in the script.

Step 2: Create an Iceberg V3 table with a VARIANT column

Create a namespace and table. The format version must be set to 3 for VARIANT data type support. The following table models IoT sensor events with a few standard columns and a VARIANT column for the semi-structured payload.

spark.sql("CREATE NAMESPACE IF NOT EXISTS glue_catalog.iceberg_v3_demo")

spark.sql("""
CREATE TABLE IF NOT EXISTS glue_catalog.iceberg_v3_demo.sensor_events (
    event_id STRING,
    device_id STRING,
    event_timestamp TIMESTAMP,
    event_data VARIANT
)
USING iceberg
TBLPROPERTIES (
    'format-version' = '3'
)
""")

The event_data column is declared as VARIANT. Iceberg stores it in Parquet as a binary-encoded VARIANT structure (metadata, value, and optional shredded sub-columns) rather than as a plain text string.

Step 3: Insert semi-structured data

To insert JSON data into a VARIANT column, use the parse_json() function. This converts a JSON string into the binary VARIANT format at write time. The following example creates a small DataFrame of IoT events and appends them to the table.

import json
from pyspark.sql.functions import current_timestamp
from pyspark.sql.types import StructType, StructField, StringType

# Sample IoT events with nested JSON payloads
events = [
    ("evt_001", "sensor_001", json.dumps({
        "device": {"manufacturer": "SensorTech", "model": "ST-200",
                   "firmware_version": "3.1.4"},
        "sensors": {"temperature": 22.5, "humidity": 61.3,
                    "air_quality": {"pm25": 12.4, "co2": 415}},
        "network": {"connection": "WiFi", "latency_ms": 42},
        "alerts": [{"severity": "low", "message": "Calibration due"}]
    })),
    ("evt_002", "sensor_002", json.dumps({
        "device": {"manufacturer": "IoTCorp", "model": "IC-500",
                   "firmware_version": "2.8.1"},
        "sensors": {"temperature": 34.1, "humidity": 78.9,
                    "air_quality": {"pm25": 142.7, "co2": 1850}},
        "network": {"connection": "LTE", "latency_ms": 210},
        "alerts": [{"severity": "critical",
                    "message": "Temperature threshold exceeded"},
                   {"severity": "high",
                    "message": "Poor air quality detected"}]
    })),
    ("evt_003", "sensor_003", json.dumps({
        "device": {"manufacturer": "SmartDevices", "model": "SD-100",
                   "firmware_version": "1.5.9"},
        "sensors": {"temperature": 18.7, "humidity": 45.2,
                    "air_quality": {"pm25": 8.1, "co2": 390}},
        "network": {"connection": "Ethernet", "latency_ms": 5},
        "alerts": []
    })),
]

schema = StructType([
    StructField("event_id", StringType(), False),
    StructField("device_id", StringType(), False),
    StructField("event_data", StringType(), False),
])

df = spark.createDataFrame(events, schema)
df = df.withColumn("event_timestamp", current_timestamp())

# Convert JSON string to VARIANT using parse_json
df = df.withColumn("event_data", parse_json(col("event_data")))

df.writeTo("glue_catalog.iceberg_v3_demo.sensor_events").append()
print("Data inserted successfully.")

The parse_json() call is the key step. It takes the raw JSON string and encodes it into the binary VARIANT format before writing to the Iceberg table.

Step 4: Query VARIANT data with variant_get()

Once the data is in the table, you can extract individual fields from the VARIANT column using variant_get(). The following queries demonstrate three common patterns: simple field extraction, deep nested access with filtering, and array element access.

The following queries are shown as raw SQL for readability. To run them in your PySpark script, wrap each query in a spark.sql() call. For example: spark.sql("SELECT ...").show().

Query 1: Simple field extraction

Extract top-level sensor readings from the payload.

SELECT
    event_id,
    device_id,
    variant_get(event_data, '$.sensors.temperature', 'double') AS temperature,
    variant_get(event_data, '$.sensors.humidity', 'double') AS humidity
FROM glue_catalog.iceberg_v3_demo.sensor_events

This query reads only the temperature and humidity sub-columns from the VARIANT data. It does not parse or load the rest of the JSON document.

Query 2: Deep nested access with filtering

Reach into nested objects and filter on a value buried inside the structure.

SELECT
    device_id,
    variant_get(event_data, '$.sensors.air_quality.pm25', 'double') AS pm25,
    variant_get(event_data, '$.sensors.air_quality.co2', 'int') AS co2_level,
    variant_get(event_data, '$.device.manufacturer', 'string') AS manufacturer
FROM glue_catalog.iceberg_v3_demo.sensor_events
WHERE variant_get(event_data, '$.sensors.air_quality.pm25', 'double') > 100.0

The WHERE clause filters directly on a nested VARIANT field. Spark evaluates the predicate against the shredded sub-column without deserializing the full payload.

Query 3: Array element access

Access elements inside a JSON array stored within the VARIANT column.

SELECT
    event_id,
    device_id,
    variant_get(event_data, '$.alerts[0].severity', 'string') AS first_alert_severity,
    variant_get(event_data, '$.alerts[0].message', 'string') AS first_alert_message
FROM glue_catalog.iceberg_v3_demo.sensor_events
WHERE variant_get(event_data, '$.alerts[0].severity', 'string') = 'critical'

Array indexing uses standard bracket notation in the JSON path. This query finds events where the first alert has critical severity and returns the alert details.

Query results showing simple field extraction, nested access with filtering, and array element access from the VARIANT column

Figure 1: Query results showing simple field extraction, nested access with filtering, and array element access from the VARIANT column.

Submitting the job to Amazon EMR Serverless

To run this on Amazon EMR Serverless, save the preceding code as a single PySpark script (for example, iceberg_v3_variant_demo.py), upload it to Amazon S3, and submit it as a job. Replace the placeholder values with your own.

Before submitting the job, make sure you have created an Amazon EMR Serverless application. For instructions, see Getting started with Amazon EMR Serverless in the Amazon EMR documentation.

# Upload script to S3
aws s3 cp iceberg_v3_variant_demo.py \
    s3://<YOUR_S3_BUCKET>/scripts/ \
    --region <REGION>

# Submit the job
aws emr-serverless start-job-run \
    --application-id <APPLICATION_ID> \
    --execution-role-arn arn:aws:iam::<ACCOUNT_ID>:role/EMRServerlessExecutionRole \
    --job-driver '{
        "sparkSubmit": {
            "entryPoint": "s3://<YOUR_S3_BUCKET>/scripts/iceberg_v3_variant_demo.py"
        }
    }' \
    --configuration-overrides '{
        "monitoringConfiguration": {
            "cloudWatchLoggingConfiguration": {
                "enabled": true,
                "logGroupName": "/aws/emr-serverless/applications/<APPLICATION_ID>"
            }
        }
    }' \
    --region <REGION>

Use cases

VARIANT fits naturally into workloads where the data is semi-structured and the schema is not fully known in advance. Some use cases include the following:

  • IoT and sensor data: Device fleets produce telemetry in varying JSON formats that evolve with firmware updates. VARIANT stores these payloads without requiring a fixed schema, and queries can extract specific readings without scanning the entire document.
  • Clickstream analytics: User behavior events on websites and mobile apps carry different attributes depending on the action. Page views, clicks, form submissions, and purchases each have their own structure. VARIANT accommodates these data types in a single column.
  • Log analytics: Application logs, infrastructure metrics, and audit trails often arrive as unstructured or loosely structured JSON. VARIANT lets you ingest them as is and query specific fields on demand, without defining a schema up front.

Clean up

To avoid ongoing charges, delete the resources you created:

  • Drop the Iceberg table and namespace using Spark SQL.
    spark.sql("DROP TABLE IF EXISTS glue_catalog.iceberg_v3_demo.sensor_events")
    spark.sql("DROP NAMESPACE IF EXISTS glue_catalog.iceberg_v3_demo")

  • Stop and delete the Amazon EMR Serverless application.
    aws emr-serverless delete-application --application-id <APPLICATION_ID> --region <REGION>

  • Delete the S3 objects and bucket used for table data, scripts, and logs.
    aws s3 rm s3://<YOUR_S3_BUCKET>/warehouse/ --recursive
    aws s3 rm s3://<YOUR_S3_BUCKET>/scripts/ --recursive

Conclusion

Apache Iceberg V3’s VARIANT type provides an efficient way to store and query semi-structured data in your data lake. Columnar storage and shredding reduce storage costs, and direct field access through variant_get() removes the need to parse JSON strings at query time. On Amazon EMR Serverless, you get this capability without managing infrastructure.

In Part 2 of this series, we scale to millions of rows and benchmark VARIANT against traditional string storage. We measure query performance and storage footprint under realistic workloads.

To learn more about Apache Iceberg on AWS, see Apache Iceberg on AWS prescriptive guidance. For more information about Amazon EMR Serverless, see the Amazon EMR Serverless documentation.


About the authors

Arun Shanmugam

Arun Shanmugam

Arun is a Senior Analytics Solutions Architect at AWS, with a focus on building modern data architecture. He has been successfully delivering scalable data analytics solutions for customers across diverse industries. Outside of work, Arun is an avid outdoor enthusiast who actively engages in CrossFit, road biking, and cricket.

Suthan Phillips

Suthan Phillips

Suthan is a Senior Analytics Architect at AWS, where he helps customers design and optimize scalable, high-performance data solutions that drive business insights. He combines architectural guidance on system design and scalability with best practices to provide efficient, secure implementation across data processing and experience layers. Outside of work, Suthan enjoys swimming, hiking, and exploring the Pacific Northwest.

Ron Ortloff

Ron Ortloff

Ron Ortloff is a Principal Product Manager at AWS, where he focuses on Apache Iceberg, S3 Tables, and open data lakehouse solutions. He has over 15 years of experience building and leading data platform initiatives, including launching Azure Synapse Analytics at Microsoft and leading Iceberg and data lake strategy at Snowflake. When he’s not building data platforms, Ron can be found cheering on his favorite football and hockey teams.

Xiaoxuan Li

Xiaoxuan Li

Xiaoxuan is a Software Development Engineer at AWS, working on the performance and scalability of Apache Iceberg in large-scale data lakehouse systems. Her interests span query optimization, storage-efficient architectures, and distributed data processing. Outside of work, she explores AI systems for creative storytelling and tooling for writers and content creators.