Tag Archives: Intermediate (200)

Best practices for Lambda durable functions using a fraud detection example

Post Syndicated from Debasis Rath original https://aws.amazon.com/blogs/compute/best-practices-for-lambda-durable-functions-using-a-fraud-detection-example/

AWS Lambda durable functions extend the Lambda programming model to build fault-tolerant multi-step applications and AI workflows using familiar programming languages. They preserve progress despite interruptions and execution can suspend for up to one year, for human approvals, scheduled delays, or other external events, without incurring compute charges for on-demand functions.

This post walks through a fraud detection system built with durable functions. It also highlights the best practices that you can apply to your own production workflows, from approval processes to data pipelines to AI agent orchestration. You will learn how to handle concurrent notifications, wait for customer responses, and recover from failures without losing progress. If you are new to durable functions, check out the Introduction to Durable Functions blog post first.

Fraud detection with human-in-the-loop

Consider a credit card fraud detection system, which uses an AI agent to analyze incoming transactions and assign risk scores. For ambiguous cases (medium-risk scores), the system needs human approval before authorizing a transaction. The workflow branches based on risk:

  • Low risk (score < 3): Authorize immediately
  • High risk (score ≥ 5): Send to the fraud department immediately
  • Medium risk (score 3–4): Suspend transaction, send SMS and email to cardholder, wait up to 24 hours for confirmation (wait time is customizable)
Figure 1. Agentic Fraud Detection with durable Lambda functions

Figure 1. Agentic Fraud Detection with durable Lambda functions

With human-in-the-loop workflows, response times can vary from minutes to hours. These delays introduce the need to durably preserve the state without consuming compute resources while waiting. With financial systems, we must also implement idempotency to guard against duplicate messages (invocations) and recover from failures without reprocessing completed work. To address these requirements, developers implement polling patterns with external state stores like Amazon DynamoDB or Amazon Simple Storage Service (Amazon S3) to manage idempotency, pay for idle compute while waiting for callbacks, introduce external orchestration components, or build asynchronous message-driven systems to handle long-processing tasks.

Lambda durable functions provide a new alternative to address these challenges through durable execution, a pattern that uses checkpoints (saved state snapshots) to preserve progress and replays from saved state to recover from failures or resume after waiting. With checkpointing capabilities, you no longer need to pay Lambda compute charges while waiting, whether for callbacks, scheduled delays, or external events. Learn how to implement durable functions using the complete fraud detection implementation at this GitHub repository. You can deploy it to your AWS account and experiment with the code as you read. The repository includes deployment instructions, sample data, and helper functions for testing.

As we walk through the code, we’ll focus on best practices for designing workflows with durable execution and how to apply these patterns correctly in production workflows.

Design steps to be idempotent

Durable execution is designed to preserve progress through checkpoints and replay, but that reliability model means step logic can execute more than once. When steps retry, how do you prevent duplicate actions like charges to the credit card or repeated customer SMS or email notifications?

Durable functions use at-least-once execution by default, executing each step at least one time, potentially more if failures occur. When a step fails, it retries. There are two strategies to design idempotent steps that prevent duplicate side effects: using external API idempotency keys and using the at-most-once step semantics built into durable functions.

Strategy A: External API Idempotency Keys

// Strategy A: Use external API idempotency keys
await context.step(`authorize-${tx.id}`, async () => {
  return payment.charges.create({
    amount: tx.amount,
    currency: 'usd',
    idempotency_key: `tx-${tx.id}`, // Prevents duplicate charges
    description: `Transaction ${tx.id}`
  });
});

Notice the configuration:

  • idempotency_key in API call: If the step retries, the payment processor recognizes it’s a duplicate request and returns the original result
  • Defense in depth: Two layers of protection: Lambda checkpointing and external API idempotency

Each layer provides independent protection. If Lambda’s checkpoint fails, the external API prevents duplicate charges. For legacy systems without idempotency support, where it’s critical that an operation is not executed more than once, use at-most-once semantics:

Strategy B: Use At-Most-Once Semantics

For legacy systems without idempotency support, use at-most-once execution, a delivery feature that executes each step zero or one time, never more:

// Strategy B: At-most-once step semantics
await context.step("charge-legacy-system", async () => {
  return await legacyPaymentSystem.charge(tx.amount);
}, {
  semantics: StepSemantics.AtMostOncePerRetry,
  retryStrategy: createRetryStrategy({ maxAttempts: 0 })
});

This checkpoints before step execution, preventing the step from re-execution on retries. The tradeoff? If the step fails, you must decide whether to retry (risking duplicates) or fail the entire workflow.

Use idempotency for critical side effects like payment processing, database writes, external API calls, state transitions, and resource provisioning. Read more about idempotency here.

Prevent duplicate executions with DurableExecutionName

Idempotent steps prevent duplicate side effects within a single execution, but what about duplicate workflow executions running concurrently? For example, duplicate messages in the queue, users clicking “Submit” multiple times in the UI, or the same event arriving via multiple channels like webhook and API. Without protection, each invocation creates a separate durable execution, potentially running the fraud check multiple times, sending duplicate notifications, and creating confusion about which execution is authoritative. Durable functions provide DurableExecutionName to help ensure only one concurrent execution per unique name.

// Invoke fraud detection function with execution name
await lambda.invoke({
  FunctionName: 'fraud-detection',
  InvocationType: 'Event',
  DurableExecutionName: `tx-${transactionId}`,
  Payload: JSON.stringify({
    id: transactionId,
    amount: 6500,
    location: 'New York, NY',
    vendor: 'Amazon.com'
  })
});

Notice the configuration:

  • DurableExecutionName: tx-${transactionId}: Uses the transaction ID as a unique execution identifier
  • InvocationType: ‘Event’: Asynchronous invocation supports long-running workflows beyond 15 minutes
  • One execution per transaction: If three invocations arrive with the same transaction ID, only the first creates an execution. Subsequent requests with the same execution name and payload receive an idempotent response returning the existing execution’s ARN, rather than creating a new execution.

Lambda durable functions work with Lambda event sources, including event source mappings (ESM) such as Amazon Simple Queue Service (Amazon SQS), Amazon Kinesis, and DynamoDB Streams. ESMs invoke durable functions synchronously and inherit Lambda’s 15-minute invocation limit. Therefore, like direct Request/Response invocations, durable functions executions using event source mappings cannot exceed 15 minutes.

For workflows exceeding 15 minutes, use an intermediary Lambda function between the event source mapping and durable function:

// Intermediary function for SQS -> Durable function
export const handler = async (event) => {
  for (const record of event.Records) {
    const transaction = JSON.parse(record.body);
    await lambda.invoke({
      FunctionName: process.env.FRAUD_DETECTION_FUNCTION,
      InvocationType: 'Event',
      DurableExecutionName: `tx-${transaction.id}`,
      Payload: JSON.stringify(transaction)
    });
  }
};

This removes the 15-minute limit, allows executions up to one year, and enables custom execution name parameters for idempotency. Use Powertools for AWS Lambda to prevent duplicate invocations of the durable function when the event source mapping retries the intermediary function. Additionally, configure failure handling for your event source to capture failed invocations for future redrive or replay. For example, dead letter queues for SQS, or on-failure destinations for other event sources.

Match timeouts to invocation type

One important configuration detail ties these patterns together: matching your timeout settings to your invocation type. Lambda synchronous invocations (RequestResponse) have a hard 15-minute timeout limit. If you configure a durable execution to run for 24 hours but invoke it synchronously, the synchronous invocation fails immediately with an exception. Durable functions support workflows up to one year when invoked asynchronously.

// Lambda function configuration
{
  FunctionName: 'fraud-detection',
  Timeout: 300,
  MemorySize: 512,
  DurableConfig: {
    ExecutionTimeout: 90000
  }
}

And invoke asynchronously:

// Async invocation for long-running workflow
await lambda.invoke({
  FunctionName: 'fraud-detection',
  InvocationType: 'Event',
  DurableExecutionName: `tx-${transactionId}`,
  Payload: JSON.stringify(transaction)
});

Notice the configuration:

  • Timeout: 300: Lambda function timeout (5 minutes in this example, up to a maximum of 15 minutes). This defines the maximum duration for each active execution phase, including the initial invocation and any subsequent replays. Set this to cover the longest expected active processing time in your workflow.
  • ExecutionTimeout: { hours: 25 }: Durable execution timeout covers the workflow’s expected total duration including suspension periods. Set this slightly above the longest wait timeout to avoid edge cases.
  • InvocationType: ‘Event’: Asynchronous invocation removes the 15-minute limit and enables executions up to one year.

The Lambda function timeout applies to active execution phases (AI calls, notification sending). During suspension (waiting for callbacks), the function isn’t running, so this timeout doesn’t apply. Setting the durable execution timeout to a meaningful boundary prevents workflows from running longer than expected. Without an explicit timeout, executions can run up to the maximum lifetime of one year.

Synchronous (RequestResponse) Asynchronous (Event)
Total duration Under 15 minutes Up to 1 year
Caller needs result Yes No
Idempotency support Yes Yes
Waits with suspension Yes Yes

Execute Concurrent Operations with context.parallel()

In the fraud detection workflow, the system notifies the cardholder through multiple channels such as SMS and email. Preserving business logic when executing parallel workflows introduces code complexities such as managing execution state across branches, handling synchronization, and coordinating branch completion. Durable functions simplify parallel workflow implementation using context.parallel(), which executes branches concurrently while maintaining durable checkpoints for each branch and provides configurable options to handle partial completions. By checkpointing and managing the state internally, durable functions help make sure that the state is preserved even if there are retries or failures. Note that context.parallel() manages the internal execution state for each branch. If your branches interact with a shared external state (such as a database), you’re responsible for managing concurrent access to that external state.

// Human-in-the-loop: verify via email AND SMS (first response wins)
let verified = await context.parallel("human-verification", [
  (ctx) => ctx.waitForCallback("SendVerificationEmail",
    async (callbackId) => sendCustomerNotification(callbackId, 'email', tx)
  ),
  (ctx) => ctx.waitForCallback("SendVerificationSMS",
    async (callbackId) => sendCustomerNotification(callbackId, 'sms', tx)
  )
], {
  maxConcurrency: 2,
  completionConfig: {
    minSuccessful: 1 // Continue after 1 success
  }
});

Notice the configuration:

  • maxConcurrency: 2: Both notifications sent at the same time
  • minSuccessful: 1: We only need one channel to succeed, whichever responds first wins

Each parallel branch waits for its callback independently, and the durable execution checkpoints each branch as part of the execution state. Using the minSuccessful parameter, you control the minimum number of successful branch executions required for the parallel operation to complete. In this example, only one of the two branches needs to succeed. Verifications through SMS or email are both valid, and the workflow resumes as soon as either channel completes successfully. We call this the first-response-wins pattern. This pattern works well when you only need a single successful result from any parallel branch and want the remaining branches to stop blocking progress.

But what happens if neither channel responds? Without timeouts, this workflow could remain suspended for up to the configured execution lifetime.

Always configure callback timeouts

Let’s add timeout protection to the parallel verification from the previous section. context.waitForCallback() accepts a timeout option that bounds how long each branch waits before throwing an exception. By wrapping the parallel call in a try/catch, you can implement fallback logic when users don’t respond in time.

// Enhanced: parallel verification with timeout and error handling
let verified;
try {
  verified = await context.parallel("human-verification", [
    (ctx) => ctx.waitForCallback("SendVerificationEmail",
      async (callbackId) => sendCustomerNotification(callbackId, 'email', tx),
      { timeout: { days: 1 } }  // Wait up to 1 day for email response
    ),
    (ctx) => ctx.waitForCallback("SendVerificationSMS",
      async (callbackId) => sendCustomerNotification(callbackId, 'sms', tx),
      { timeout: { days: 1 } }  // Wait up to 1 day for SMS response
    )
  ], {
    maxConcurrency: 2,
    completionConfig: {
      minSuccessful: 1
    }
  });
} catch (error) {
  const isTimeout = error.message?.includes("timeout");
  if (isTimeout) {
    context.logger.warn("Customer verification timeout", { error, txId: tx.id });
    // Fallback: escalate to fraud department
    return await context.step("sendToFraudDepartment", async () =>
      sendToFraudDepartment(tx, true)
    );
  }
  throw error; // Re-throw non-timeout errors
}

Notice what changed from the previous section:

  • timeout: { days: 1 }: Each callback branch now has a maximum wait time of 1 day. If neither the email nor SMS callback arrives within that window, a timeout exception is thrown.
  • try/catch with timeout detection: The catch block distinguishes between timeout errors and other exceptions. When a timeout occurs, the workflow implements fallback logic by escalating the transaction to the fraud department, while non-timeout errors are re-thrown to be handled by the durable execution retry mechanism.

Without this error handling, the entire execution fails unhandled. The timeout also works with the minSuccessful configuration: if one branch times out but the other succeeds, the parallel operation still completes successfully since only one successful result is required.

For advanced use cases where the callback handler performs long-running work, you can also configure a heartbeatTimeout to detect stalled callbacks before the main timeout expires. See the Lambda Developer Guide for details.

Use callback timeouts for human approvals, external API callbacks, asynchronous processing, and third-party integrations.

Putting it all together: complete fraud detection implementation

Now let’s see how all the best practices work together in the complete fraud detection workflow:

import { withDurableExecution } from "@aws/durable-execution-sdk-js";
import { BedrockAgentCoreClient, InvokeAgentRuntimeCommand } from "@aws-sdk/client-bedrock-agentcore";

const agentRuntimeArn = process.env.AGENT_RUNTIME_ARN;
const agentRegion = process.env.AGENT_REGION || 'us-east-1';
const client = new BedrockAgentCoreClient({ region: agentRegion });

export const handler = withDurableExecution(async (event, context) => {
  const tx = {
    id: event.id,
    amount: event.amount,
    location: event.location,
    vendor: event.vendor
  };

  // AI fraud assessment with error handling
  tx.score = await context.step("fraudCheck", async () => {
    try {
      const payloadJson = JSON.stringify({ input: { amount: tx.amount } });
      const command = new InvokeAgentRuntimeCommand({
        agentRuntimeArn: agentRuntimeArn,
        qualifier: 'DEFAULT',
        payload: Buffer.from(payloadJson, 'utf-8'),
        contentType: 'application/json',
        accept: 'application/json'
      });
      const response = await client.send(command);
      const responseText = await response.response.transformToString();
      const result = JSON.parse(responseText);
      return result?.output?.risk_score ?? 5;  // Default to high-risk if score unavailable
    } catch (error) {
      context.logger.error("Fraud check failed", { error, txId: tx.id });
      return 5;
    }
  });

  // Route based on AI decision
  if (tx.score < 3) {
    // Best Practice: Idempotent authorization
    return await context.step(`authorize-${tx.id}`, async () =>
    authorizeTransaction(tx, { idempotency_key: `tx-${tx.id}` })
    );
  }

  if (tx.score >= 5) {
    return await context.step(`sendToFraudDepartment-${tx.id}`, async () =>
      sendToFraudDepartment(tx)
    );
  }

  // Medium risk: need human verification
  await context.step(`suspend-${tx.id}`, async () => suspendTransaction(tx));

  // Best Practice: Concurrent operations with timeout configuration
  let verified;
  try {
    verified = await context.parallel("human-verification", [
      (ctx) => ctx.waitForCallback("SendVerificationEmail",
        async (callbackId) => sendCustomerNotification(callbackId, 'email', tx),
        { timeout: { days: 1 } }
      ),
      (ctx) => ctx.waitForCallback("SendVerificationSMS",
        async (callbackId) => sendCustomerNotification(callbackId, 'sms', tx),
        { timeout: { days: 1 } }
      )
    ], {
      maxConcurrency: 2,
      completionConfig: {
        minSuccessful: 1
      }
    });
  } catch (error) {
    const isTimeout = error.message?.includes("timeout");
    context.logger.warn(
      isTimeout ? "Customer verification timeout" : "Customer verification failed",
      { error, txId: tx.id }
    );
    return await context.step(`timeout-escalate-${tx.id}`, async () =>
      sendToFraudDepartment(tx, true)
    );
  }

  // Idempotent final step with idempotency key
  return await context.step(`finalize-${tx.id}`, async () => {
    const action = !verified.hasFailure && verified.successCount > 0
      ? "authorize"
      : "escalate";
    if (action === "authorize") {
      return authorizeTransaction(tx, true, { idempotency_key: `finalize-${tx.id}` });
    }
    return sendToFraudDepartment(tx, true);
  });
});

Notice how the best practices work together: context.parallel() sends SMS and email concurrently, resuming when either channel responds. Both callbacks configure 1-day timeouts with try/catch handling that escalates on timeout. The DurableExecutionName: tx-${transactionId} parameter (specified at invocation time, shown in the following CLI example) provides execution-level deduplication, while idempotency keys in the authorization steps prevent duplicate charges at the application layer. Asynchronous invocation (InvocationType: 'Event') enables the 24-hour wait period.

Once deployed, invoke the function asynchronously with a sample transaction to see it in action:

transactionId="123456789"
aws lambda invoke \
  --function-name "fraud-detection:$LATEST" \
  --invocation-type Event \
  --durable-execution-name "tx-${transactionId}" \
  --cli-binary-format raw-in-base64-out \
  --payload "{\"id\": \"${transactionId} \", \"amount\": 6500, \"location\": \"New York, NY\", \"vendor\": \"Amazon.com\"}" \
  --region us-east-2 \
  response.json

Upon successful invocation, you can view the execution state in the Lambda console’s durable operations view. The execution shows a suspended state, waiting for customer response:

Figure 2: Suspended execution state

Figure 2: Suspended execution state

Notice the fraudCheck and suspendTransaction steps show as succeeded with checkpointed results. The human-verification parallel operation shows that both SMS and email branches started. The timeline shows the function in a suspended state. Simulate a customer response by sending a callback success through the console, AWS Command Line Interface (AWS CLI) or Lambda API:

aws durable-lambda send-durable-execution-callback-success \
  --callback-id <CALLBACK_ID_FROM_EMAIL_OR_SMS> \
  --result '{"status":"approved","channel":"email"}' \
  --cli-binary-format raw-in-base64-out
Figure 3: Completed execution with customer approval

Figure 3: Completed execution with customer approval

After receiving the customer’s approval, the durable execution resumes from its checkpoint, authorizes the transaction, and completes. The execution spanned hours but consumed only seconds of compute time.

Conclusion

With durable functions, Lambda extends beyond single-event processing to power core business processes and long-running workflows, while retaining the operational simplicity, reliability, and scale that define Lambda. You can build applications that run for days or months, survive failures, and resume where they left off, all within the familiar event-driven programming model.

Deploy the fraud detection workflow from our GitHub repository and experiment with human-in-the-loop patterns in your own account. For core concepts, see Introduction to AWS Lambda Durable Functions. For comprehensive documentation, see the Lambda Developer Guide. Browse Serverless Land for reference architectures and discover where durable execution fits in your designs.

Share your feedback, questions, and use cases in the SDK repositories or on re:Post.

Simplifying Kafka operations with Amazon MSK Express brokers

Post Syndicated from Mazrim Mehrtens original https://aws.amazon.com/blogs/big-data/simplifying-kafka-operations-with-amazon-msk-express-brokers/

In this post, we show you how Amazon Managed Streaming for Apache Kafka (Amazon MSK) Express brokers brokers streamline the end-to-end activities for Kafka administration. Apache Kafka has become the de facto standard for real-time data streaming, powering mission-critical applications across industries worldwide. Its popularity stems from its ability to handle high-throughput, fault-tolerant data pipelines at scale. Given its central role in modern data architectures, managing Apache Kafka with high resilience and reliability is essential for business success.

To maintain this level of resilience, administrators need to handle several important operational tasks. Apache Kafka is a distributed stateful system, whose state management requires constant communication and data movement in dynamic cloud environments. Administrators need to carefully size clusters by calculating complex compute, storage, and network requirements. They must provision storage volumes upfront and monitor utilization constantly to avoid disruptions. When workloads grow, scaling the cluster requires hours or days of effort using multiple tools to provision capacity and rebalance load.

With these operational requirements in mind, many administrators ask: is there an easier way to manage Apache Kafka at scale while maintaining the high resilience their applications demand?

Amazon MSK Express addresses these challenges directly. In this post, we show you how MSK Express brokers streamline the end-to-end activities for Kafka administration, including:

  • Sizing Kafka clusters for optimal performance and cost
  • Scaling cluster storage up and down with workload changes
  • Scaling cluster compute in and out over time
  • Monitoring cluster health
  • Managing cluster security
  • Ensuring high availability with fast and automatic broker recovery

What are Amazon MSK Express brokers?

Amazon MSK Express brokers are a transformative breakthrough for customers needing high-throughput Kafka clusters that scale faster and cost less. Express brokers reimagine Kafka’s compute and storage, decoupling to unlock performance and elasticity benefits. Express brokers deliver performance improvements that directly impact your operations:

  • Up to 3x more throughput per broker, allowing you to handle more data with fewer resources and lower costs
  • Rebalance partitions across brokers 180x faster, reducing scaling from hours to minutes
  • Scale up to 20x faster, enabling you to respond to demand spikes without lengthy planning cycles
  • Recover 90% quicker compared to standard Apache Kafka brokers, minimizing workload disruption and maintaining business continuity

To learn more about the technical details, see Express brokers for Amazon MSK: Turbo-charged Kafka scaling with up to 20 times faster performance. For a comprehensive overview of Express broker capabilities, see the MSK Express brokers documentation.

Let’s explore how MSK Express brokers simplify Apache Kafka management.

Sizing an Express cluster

Sizing a traditional Apache Kafka cluster is complex. Working backwards from your ingress and egress load, you need to consider every dimension of your cluster compute, storage, and network limitations. Each node must be carefully sized to handle:

  • Ingress and egress traffic from your clients
  • Internal Kafka operations like replication and rebalancing (the process of redistributing partitions across brokers to maintain balance)
  • High availability with node and Availability Zone failures
  • Client operations like backfill procedures when reading historical data

These activities impact your cluster storage I/O limits, network ingress/egress limits, and CPU and memory constraints. Beyond this, you need to consider the number of partitions required and determine whether your cluster can scale to handle partition management for your use case.

MSK Express brokers simplify this calculus. Rather than considering these complex variables, you can focus on what matters:

  • Your ingress throughput
  • Your egress throughput
  • Your partition needs

MSK documents the Express broker throughput throttle and partition limits by broker size. MSK pre-calculates these to consider all cluster limits. They include multi-Availability Zone high availability to handle rare events like node failures or AZ impairment.

Notice we did not discuss storage in sizing an Express cluster. That is because storage in Express scales nearly infinitely. You pay for storage as you go rather than sizing storage up front.

Scaling Express cluster storage

With sizing simplified by focusing on throughput and partitions, storage management becomes the next operational consideration.

Normally, Apache Kafka clusters need storage volumes pre-provisioned to handle all retained data. You must allocate all storage up-front and pay for that storage no matter what your actual data retention is.

Example: If you store 7 days of data at 1 MB/sec ingress, that’s 600+ GB of storage. This does not include data replication across nodes and buffers for growth and workload variability. This workload requires over 3 TB of storage, allocated up-front, to handle replicas and storage buffers.

As your workload evolves, careful monitoring of storage utilization becomes essential. Adding storage capacity prevents workload disruptions. Often, you cannot reclaim this storage. Once you increase the volume size, you continue paying for additional storage even if your workload scales down and no longer requires additional capacity.

With Express brokers, there is no need for sizing and provisioning storage volumes. You pay for what you use with no provisioning: the data ingested to the cluster and data stored in the cluster per-GB-per-hour. All data stored in the cluster is replicated across 3 Availability Zones for high availability. This pay-as-you-go model eliminates wasted capacity costs and reduces your total infrastructure spend.

  • As workloads scale up, the cluster uses more storage with no changes needed from you
  • When workloads scale down, the cluster uses less storage, reducing storage charges automatically
  • Storage management for Apache Kafka becomes simpler with Express. You focus on ensuring that your per-topic retention is right-sized for each use case. That is the only consideration. Once you set up topic retention, MSK Express automatically manages and cost-optimizes storage on your behalf.

Storage management in MSK Express brokers is far simpler than in a traditional Apache Kafka cluster. So is scaling the compute capacity for an Express-based cluster.

Scaling Express cluster compute

Just as storage scales automatically with your workload, compute capacity can also adapt to changing demands.

As your workload grows and changes, you may find that you exceed your initial sizing estimates. For a traditional Apache Kafka cluster, scaling the cluster capacity is a significant event. Scaling takes effort to provision capacity and rebalance load, it requires using multiple tools to manage the scaling process (compute, storage, DNS, rebalancing, client configs, and more). The scaling process can take hours or days to complete, which can exacerbate application impact. This means you need to plan well ahead to ensure your Kafka cluster is prepared for any load changes.

With MSK Express clusters, this process becomes much simpler and requires little to no upfront planning. It has near zero disruption to your existing workload, allowing your team to focus on building features rather than managing infrastructure.

To scale up an MSK Express cluster, you simply add brokers to the cluster. Once new brokers come online, Express Intelligent Rebalancing automatically rebalances topic partitions to the new nodes. Thanks to the Express storage architecture, the new nodes automatically have almost all the data they need. There is no significant inter-broker communication for rebalancing. This causes no disruption to existing brokers.

The cluster then elects new broker leaders for each partition, enabling producers to direct traffic to the new nodes. The same applies to consumer groups.

Express broker DNS design keeps this in mind. Express broker connection strings abstract away from the nodes themselves. Clients connect to the active broker nodes with one connection string. No changes to DNS, load balancing, or client configurations are needed.

Deciding when to scale in an Express cluster is also simpler than in a traditional Apache Kafka cluster. The simplified Express architecture means less to monitor and manage for long-term cluster operations.

Monitoring Express clusters

With simplified scaling decisions comes simplified monitoring. Express brokers reduce the number of metrics you need to track for cluster health. The below image demonstrates a dashboard which highlights the key metrics for monitoring MSK Express broker health.

Dashboard with key Amazon MSK Express brokers metrics

In a traditional Apache Kafka cluster, you need to consider dozens of metrics to understand overall cluster health. Express brokers simplify this operational process. They highlight ingress and egress throughput as two critical metrics for workload sizing and scaling. This streamlined monitoring approach reduces the expertise required to operate Kafka clusters and allows smaller teams to manage larger deployments effectively.

Other factors, like poorly designed clients, can incur additional overhead on a cluster. This can cause symptoms such as high CPU utilization without high ingress throughput. It is still important to monitor a variety of metrics with MSK Express brokers.

For Express brokers, the following table shows the critical metrics you must monitor and alert on for cluster health:

Metric Name Description Recommended Alarm
BytesInPerSec Ingress throughput to the cluster When > broker limit for > 5 minutes
BytesOutPerSec Egress throughput to the cluster When > broker limit for > 5 minutes
CpuUser + CpuSystem CPU utilization percentage When greater than 60% for 15 minutes
NetworkProcessorAvgIdlePercent Network processor thread idle time When less than 0.5 for > 5 minutes
RequestHandlerAvgIdlePercent Request processor thread idle time When less than 0.4 for > 15 minutes
FetchThrottleByteRate Consumer fetch throttling rate When < 0 for > 15 minutes
ProduceThrottleByteRate Producer ingress throttling rate When < 0 for > 15 minutes

For more information on monitoring Amazon MSK, see Monitoring Amazon MSK with Amazon CloudWatch.

Managing Express cluster access

Beyond monitoring, cluster management is another area where MSK Express brokers reduce operational complexity.

Express brokers simplify the internal management of Kafka clusters. In a traditional Kafka environment, you use schemes like SASL/SCRAM (username and password-based authentication) or mutual TLS (certificate-based authentication) for client authentication. Once authenticated, you configure complex Kafka ACLs (Access Control Lists—permissions that define who can access which topics) inside the Kafka cluster to authorize client access to topics and data.

These paradigms require you to manage all topics, authentication, and authorization inside Apache Kafka. This includes credential management, rotation, and other operational activities surrounding cluster access.

MSK simplifies this process by integrating with AWS Identity and Access Management (IAM) for access control. Clients can use IAM Roles that clearly specify cluster access boundaries. They also provide topic-level authorization to read and write data to a cluster with Kafka APIs.

Finally, clients can use MSK APIs to directly manage Kafka cluster configurations and Kafka topics, including creating new topics, updating topic configurations and partition counts, and deleting topics. Configurations and topics can be managed with the AWS Console, AWS CLI, and AWS SDK. For more information, refer to Amazon MSK simplifies Kafka topic management with new APIs and console integration.

You can focus only on your existing enterprise standards for IAM access controls, and your existing AWS CloudFormation and AWS CDK automation to manage your cluster with Infrastructure as Code (IaC). This integration reduces the operational overhead of cluster management and accelerates your time to production by leveraging existing security infrastructure.

MSK also supports using SASL/SCRAM and mutual TLS authentication modes alongside IAM access control. This gives you the flexibility to authorize applications outside of AWS. You can also provide access to legacy applications without the need for code changes.

For more information, see IAM access control for Amazon MSK and Security in Amazon MSK.

Building highly available Express brokers

With security simplified through IAM integration, high availability is the final piece of the operational puzzle.

Many of the same considerations we discussed in scaling Express cluster compute align with high availability considerations for MSK Express brokers.

Based on internal testing, MSK Express broker storage improvements enable faster recovery when broker nodes fail—90% faster than standard brokers. The new node can simply start up with almost no disruption to the rest of the cluster without needing to perform significant rebalancing. This contrasts with standard Kafka clusters, where the cluster needs to rebalance partitions to new nodes after recovery.

In addition to these improvements, MSK Express brokers are highly available by default. The service manages critical cluster and topic configurations for high availability and performance on your behalf. This eliminates the need for managing most cluster configurations.

Express fully manages configurations like min.insync.replicas, num.io.threads, and others described in Express brokers’ read-only configurations. This gives you a highly available and performant cluster out of the box.

You no longer need to worry about most cluster-level configurations of an Apache Kafka cluster. You can simply:

  • Start an MSK Express cluster
  • Configure topics and retention
  • Proceed without the fine tuning normally needed to ensure a highly available cluster

Conclusion

In this post, we showed how MSK Express brokers simplify cluster operations for Apache Kafka clusters. They lower the Total Cost of Ownership (TCO) of running an Apache Kafka cluster by simplifying sizing, storage management, compute management, high availability, and access control, while providing high performance, reliability, and cost-efficiency. These simplifications reduce the specialized expertise needed for cluster administration and accelerate your deployment timeline.

With this in mind, we recommend MSK Express brokers for almost all MSK workloads. If you are starting out with a new Kafka cluster or optimizing an existing one, MSK Express brokers provide a strong combination of simplicity, performance, and cost-efficiency.

Ready to simplify your Kafka operations? Get started using Amazon MSK to create your first Express cluster today. You can provision a fully managed, highly available Kafka cluster in minutes and start experiencing the operational benefits immediately. For pricing details, see Amazon MSK pricing.

For comprehensive information about Amazon MSK capabilities and features, visit the Amazon MSK product page and the Amazon MSK Developer Guide.


About the authors

Mazrim Mehrtens

Mazrim Mehrtens

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

Sai Maddali

Sai Maddali

Sai is a Senior Manager Product Management at AWS who leads the product team for Amazon MSK. He is passionate about understanding customer needs, and using technology to deliver services that empowers customers to build innovative applications. Besides work, he enjoys traveling, cooking, and running.

Filter catalog assets using custom metadata search filters in Amazon SageMaker Unified Studio

Post Syndicated from Ramesh H Singh original https://aws.amazon.com/blogs/big-data/filter-catalog-assets-using-custom-metadata-search-filters-in-amazon-sagemaker-unified-studio/

Finding the right data assets in large enterprise catalogs can be challenging, especially when thousands of datasets are cataloged with organization-specific metadata. Amazon SageMaker Unified Studio now supports custom metadata search filters. You can filter catalog assets using your own metadata form fields like therapeutic area, data sensitivity, or geographic region rather than relying only on free-text search. Custom metadata forms are structured templates that define additional attributes that can be attached to catalog assets.

In this post, you learn how to create custom metadata forms, publish assets with metadata values, and use structured filters to discover those assets. We explore a healthcare and life sciences use case. A research organization catalogs metrics in Amazon SageMaker Catalog using custom metadata forms with fields such as Therapeutic Area and Sample Size. Researchers building Machine learning models can now search datasets based on custom filters across hundreds of cataloged assets to identify the best datasets to train their models.

Key capabilities

Custom metadata search filters in SageMaker Unified Studio offer the following key capabilities:

  • Custom metadata form filters – You can filter search results using any custom metadata form fields defined in their catalog. For example, a researcher can filter by Therapeutic Area = Oncology and Data Sensitivity = Confidential to locate specific datasets.
  • Name and description filters – You can add filters that target asset names or descriptions using a text search operator, enabling targeted discovery without scanning full search results.
  • Date range filters – You can filter assets by date using on, before, after, and between operators, making it straightforward to locate recently updated or historically relevant assets.
  • Combinable filters – You can combine multiple filters to construct precise queries. For example, filtering by AWS Region = US AND Classification = PII AND Updated after 2026-01-01 returns only assets matching all three criteria.
  • Persistent filter selections – You can filter configurations stored in your browser and are not shared across devices or other users. You can later return to the catalog and find your previously defined filters.

Solution overview

In the following sections, we demonstrate how to set up custom metadata forms, publish assets with metadata values, and use custom metadata search filters to discover those assets.We complete the following three steps for the demonstration.

  1. Create a custom metadata form
  2. Create and publish assets with metadata
  3. Use custom metadata search filters

Prerequisites

To follow along with this post, you should have:

For instructions on setting up a domain and project, see the Getting started guide.

To create a custom metadata form

Complete the following steps to create a custom metadata form with filterable fields:

  1. In SageMaker Unified Studio, choose Project overview from the navigation pane.
  2. Under Project catalog, choose Metadata entities.
  3. Choose Create metadata form.
  4. To create a new metadata form ‘research_metadata’ use the following details, then choose Create metadata form.
  5. Define the form fields. For this demo, we add the following fields:

    Create first field Therapeutic Area (String) – Mark as Searchable


    Create second field Subject Count (Integer) – Mark as Filterable by range

  6. Mark the form as ‘Enabled’ so the form is visible and can be used.

Create and publish with metadata

In this section, you create a custom asset and attach the research_metadata form created in the previous step.

  1. Under Project catalog in the navigation pane, choose Metadata entities. Choose the ‘ASSET TYPES’ tab and select “CREATE ASSET TYPE’.
  2. Create a new asset type and attach the metadata form that we created in the previous step.

    A new asset type ‘metric’ is created.
  3. Next, we will create two metrics. Under Project catalog in the navigation pane, choose Assets. On the Asset page, choose CREATE, and then choose Create asset from the menu.
  4. In this demo, you create two metrics.

For the first metric ‘drug_1_treatment’, provide the following asset name and description.

Add the following values for the metadata form.

Validate all fields and choose CREATE.

Publish the asset to the catalog.

Next, we will create the second metric ‘drug_1_treatment’. Repeat the steps from the previous procedure and enter the values shown.

  • Subject Count = 450
  • Therapeutic Area = Oncology

Use custom metadata search filters

After publishing assets with custom metadata, go to the Browse Assets page to use the filters.

To browse assets and view filters

  1. In SageMaker Unified Studio, choose Discover from the navigation bar, then select Catalog, Browse Assets.
  2. The search page displays with the filter sidebar on the left. You can see the existing system filters (Data type, Glossary terms, Asset type, Owning project, Source Region, Source account, Domain unit) along with the new Date range and Add Filter sections.

Add a custom filter

  1. Choose + Add Filter at the bottom of the filter sidebar. For Filter type, select Metadata form. For Metadata form, select research_metadata and add a filter as shown in the following image. Choose Apply when you’re done.

    The search results update to show only assets where ‘subject_count’ is greater than 50.

To combine multiple filters

  1. Choose + Add Filter again. For Filter type, select Metadata form. For Metadata form, select research_metadata and add a filter as shown in the following image. Choose Apply when you’re done.

Manage custom filters

Filter configurations are stored in the user’s browser and are not shared across devices or users.

To customize search, you could:

  • Toggle filters – Use the checkboxes next to each custom filter to enable or disable them without deleting.
  • Edit or delete – Choose the kebab menu (⋮) next to any custom filter to edit its values or delete it.
  • Clear all – Choose CLEAR next to the Custom filters header to deselect all custom filters at once.
  • Persistence – Your custom filters persist across browser sessions. When you return to the Browse Assets page, your previously defined filters are still listed in the sidebar, ready to be activated.

Using the SearchListings API

To search catalog assets programmatically, you can use the SearchListings API in Amazon DataZone, which supports the same filtering capabilities as the SageMaker Unified Studio UI. The following example filters assets where a custom string field contains a specific value and a numeric field is within a range:

aws datazone search-listings \
    --domain-identifier "dzd_your_domain_id" \
    --filters '{ "and": [
        { "filter": { "attribute": "research_metadata.TherapeuticArea", "value": "Oncology", "operator": "TEXT_SEARCH" } },
        { "filter": { "attribute": "research_metadata.SubjectCount", "intValue": 100, "operator": "GT" } }
    ] }'

For more details, see the SearchListings API documentation in the Amazon DataZone API Reference.

Best practices

Consider the following best practices when using custom metadata search filters:

  • Define your metadata forms before publishing assets at scale. If you publish assets before the forms are finalized, you might need to re-tag existing assets, which is a time-consuming process in large catalogs.
  • Define metadata forms aligned with your organization’s discovery needs (therapeutic areas, data classifications, geographic regions) before publishing assets at scale.
  • Use specific, consistent values in metadata fields to get precise filter results. For example, use standardized values (for example, use “Oncology” consistently rather than “oncology” or “Onc”) across all assets.
  • Combine multiple filters to narrow results efficiently rather than scanning through broad result sets.
  • Use the date range filter alongside custom metadata filters to locate assets within specific time windows.

Clean up resources

For instructions on deleting the added assets, see Delete an Amazon SageMaker Unified Studio asset.
For instructions on deleting the metadata forms, see Delete a metadata form in Amazon SageMaker Unified Studio.

Conclusion

Custom metadata search filters in Amazon SageMaker Unified Studio give data consumers the ability to find exact assets using structured filters based on their organization’s own metadata fields. By combining multiple filters across custom metadata forms, asset names, descriptions, and date ranges, data consumers can construct precise queries that surface the right datasets without scanning through broad search results. Filter persistence across browser sessions further streamlines repeated discovery workflows.

Custom metadata search filters are now available in AWS Regions where Amazon SageMaker is supported.

To learn more about Amazon SageMaker, see the Amazon SageMaker documentation. To get started with this capability, refer to the Amazon SageMaker Unified Studio User Guide.


About the authors

Ramesh Singh

Ramesh Singh

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

Pradeep Misra

Pradeep Misra

Pradeep is a Principal Analytics and Applied AI Solutions Architect at AWS. He is passionate about solving customer challenges using data, analytics, and Applied AI. Outside of work, he likes exploring new places and playing badminton with his family. He also likes doing science experiments, building LEGOs, and watching anime with his daughters.

Alexandra von der Goltz

Alexandra von der Goltz

Alexandra is a Software Development Engineer (SDE) at AWS based in New York City, on the Amazon SageMaker team. She works on the catalog and data discovery experiences within the Unified Studio.

How Vanguard transformed analytics with Amazon Redshift multi-warehouse architecture

Post Syndicated from Alex Rabinovich original https://aws.amazon.com/blogs/big-data/how-vanguard-transformed-analytics-with-amazon-redshift-multi-warehouse-architecture/

This is a guest post by Alex Rabinovich, Anindya Dasgupta, and Vijesh Chandran from Vanguard, Financial Advisor Services division, in partnership with AWS.

Vanguard stands as one of the world’s leading investment companies, serving more than 50 million investors globally. The company offers an extensive selection of low-cost mutual funds and ETFs with over 450 funds/ETFs along with comprehensive investment advice and related financial services. With a workforce of approximately 20,000 crew members, Vanguard has built its reputation on providing low-cost, high-quality investment solutions that help investors achieve their long-term financial goals.

Within this massive organization, Vanguard’s Financial Advisor Services (FAS) division stands as one of the most prominent B2B operations in the financial services industry. Operating at an extraordinary scale, FAS oversees a broad range and diverse range of assets through the intermediary channel while supporting a vast network of advisory firms and financial advisors across the country. This division delivers a full suite of investment products, model portfolios, research capabilities, and technology-driven support services designed to help financial advisors serve their clients more effectively.

Business use cases and initial architecture

The scale and complexity of FAS operations generate enormous amounts of data that require sophisticated analytics capabilities to drive business insights, regulatory compliance, and operational efficiency. To address this, Vanguard launched the FAS 360 initiative. This initiative aims to empower Financial Advisor Services (FAS) with a centralized cloud data warehouse that integrates both internal and external data sources into a unified, intelligent system.

Key business use cases:

  1. Business operations – Enables sales goal setting, tracking, and compensation management to drive operational excellence. It delivers insights on product usage patterns across financial advisor clients.
  2. Data science – Powers customer segmentation models and call transcription analytics to drive strategic insights. It also supports marketing campaign preparation and customer insights for sales call preparation.
  3. Exploratory analytics – Enables ad-hoc leadership questions, what-if scenario analysis, and sales trend analysis for channel managers competitor comparative analysis.

By consolidating these use cases into a centralized system, FAS 360 enables consistent reporting and data-driven decision-making across Vanguard’s Financial Advisor Services division.

Centralized data warehouse FAS 360:

Vanguard’s first wave of modernization established FAS 360 as a centralized enterprise data warehouse, migrating from a fragmented “data swamp” of Parquet files on Amazon Simple Storage Service (Amazon S3) to a structured, unified system.

The following architecture diagram leverages Amazon S3 for raw data storage with Amazon Redshift serving as the core processing engine, providing integrated access for BI tools, analyst exploration, and data science workloads.

Here are the key benefits achieved with this architecture:

  • Single source of truth – Consolidated fragmented data sources into a unified system, minimizing multiple versions of truth and establishing consistent reporting practices across the organization
  • 10x faster query performance – Dramatically improved query response times compared to the previous solution, helping enhance analyst productivity and enabling more complex analytical workloads
  • Seamless data lake integration – Maintained connectivity with the broader data lake environment while providing structured warehouse capabilities
  • Enhanced business agility – Increased trust in metrics and unlocked new use cases that were previously untenable, directing the new migration efforts toward the FAS360 system

This centralized architecture successfully addressed the limitations of Vanguard’s previous approach, where data was scattered across individuals with limited governance, and established a foundation for their subsequent architectural evolution.

Significant growth and expanding use cases

Vanguard FAS experienced remarkable growth in their data analytics requirements over a two-year period, demonstrating the rapid evolution of modern data needs:

Initial State:

  • 20 AWS Glue ETL jobs processing daily data loads
  • Approximately 100 tables in their data warehouse
  • 20 Tableau dashboards serving business users
  • Around 60 analysts accessing the system

Two Years Later:

  • 20 TB in data volume in Amazon Redshift and another 150 TB in S3 data lake
  • 600+ AWS Glue ETL jobs (a 30x increase) handling complex data transformations
  • 300+ tables (3x growth) storing diverse business data
  • 250+ Amazon Redshift materialized views optimizing query performance
  • Over 500 Tableau dashboards (25x expansion) serving various business functions
  • 500,000+ user queries/months

This exponential growth reflected FAS’s increasing reliance on data-driven decision making across the business functions, from risk management and compliance to client service optimization and operational efficiency improvements.

Resource contention and performance bottlenecks

As Vanguard FAS’s data environment expanded, their initial architecture, a single Amazon Redshift provisioned cluster with 2 nodes (ra3.4xlarge), began experiencing severe performance challenges that threatened business operations:

ETL performance issues:

  • Frequent ETL SLA failures disrupting critical business processes
  • Tableau extract failures resulting in stale dashboard data
  • Resource conflicts between data ingestion and transformation workloads

End-user experience degradation:

  • Poor query performance during peak usage periods
  • Table and object locking issues preventing concurrent access
  • Frustrated analysts unable to perform deep data exploration
  • Limited ability to run long-running analytical queries

Operational challenges:

  • Resource contention between ETL workloads and interactive analytics
  • Inability to scale compute resources independently for different workload types
  • Single point of failure affecting the data operations
  • Difficulty in workload prioritization and resource allocation

These challenges were fundamentally limiting FAS’s ability to leverage their data assets effectively, impacting everything from daily operational reporting to strategic business analysis.

Solution overview

To address these critical challenges, Vanguard FAS implemented following multi-warehouse architecture that leverages the advanced data sharing capabilities of Amazon Redshift for workload isolation and independent scaling.

Producer – Amazon Redshift Provisioned Cluster

The central hub consists of the original Amazon Redshift provisioned cluster with RA3 nodes, optimized for consistent, predictable workloads:

  • Dedicated ETL processing: Handles data ingestion, transformation, and loading operations
  • Write workload optimization: Manages data writes and updates without interference
  • Cost optimization: Utilizes reserved instances for predictable, steady-state workloads
  • Data governance: Serves as the single source of truth for the enterprise data

Consumer – Amazon Redshift Serverless Workgroups

Multiple Amazon Redshift Serverless instances serve as specialized consumer endpoints which auto-scales compute resources based on demand:

  • Analyst Exploration: Dedicated environment for analyst data discovery and experimentation
  • BI Tools: Instance optimized specifically for Tableau dashboard and visualization workloads
  • Data Science: For complex and long running machine learning workloads in completely isolated environment

The solution leverages the native data sharing capabilities of Amazon Redshift to enable secure connectivity between the producer and consumers instances. Consumer clusters can access live data from the producer without data movement, providing real-time access to the most current information available. This zero-copy sharing approach alleviates the need for data duplication or complex synchronization processes, helping reduce both storage costs and operational complexity.

Results

The implementation of the multi-warehouse architecture delivered significant improvements across the key performance indicators:

Predictable Performance

Nightly ETL cycles now consistently complete before the 9 AM SLA, eliminating the previous SLA failures that disrupted business operations and ensuring fresh data is available for morning business activities. Dashboards and reports now reflect the most current data available, providing teams with up-to-date insights for decision-making.

Improved Analyst Productivity and Experience

The new architecture removed the restrictive 10-minute query timeout that previously prevented deep ad hoc exploratory queries. Analysts can now run complex analytical workloads exceeding 30 minutes in a fully isolated environment without impacting other users or ETL processes. This change, combined with significantly faster query response times, has led to higher analyst satisfaction and productivity across the team.

New Analytical Capabilities

The architecture introduced a dedicated “Data Lab” environment where analysts have write access to experiment with data using CREATE TABLE AS SELECT (CTAS) commands. Each workload type can now scale independently based on demand, with different consumer clusters optimized for specific use cases, enabling more sophisticated analytical approaches.

Operational Excellence

The separation of workloads enabled efficient utilization of compute resources across different patterns, leading to better cost control through appropriate sizing, serverless pay-as-you-go pricing, and reserved instance usage. The cleaner separation of concerns between ETL and analytics workloads has simplified overall management of the data platform.

Ongoing modernization: Evolution toward data mesh architecture

As Vanguard’s data environment matured and their success with the multi-warehouse architecture enabled broader adoption across the organization, they recognized an opportunity to evolve their architecture to match their organizational growth. The expanding portfolio of data products and increasing number of teams leveraging the system created new opportunities for innovation.

As Vanguard’s data environment grew, three key challenges emerged:

  1. Centralized ownership bottleneck – Single-team data ownership couldn’t scale with the growing number of data products
  2. Write workload contention – Resource contention persisted for write operations on shared endpoints
  3. Cross-domain dependencies – Data object interdependencies across business domains slowed data product development

Rationale for Data Mesh

Vanguard’s decision to adopt Data Mesh was driven by the need to:

  • Decentralize data ownership by establishing data domains with dedicated stewards
  • Remove write contention by isolating each domain’s data loads to separate endpoints
  • Enable autonomous development allowing stewards to own the complete data product lifecycle and governance
  • Leverage modern data lake capabilities using AWS Glue and Apache Iceberg format for data product curation

This evolution supports Vanguard’s ability to scale organizationally while building on the technical foundation and operational excellence achieved with their multi-warehouse architecture. Building on the success of their Amazon Redshift multi-warehouse implementation, Vanguard FAS is now exploring on the next phase of their data architecture evolution, implementing following data mesh approach.

This new data mesh architecture has several key components that work together to enable scalable, domain-oriented data management.

Domain-Oriented Data Ownership

Vanguard is establishing distinct data domains aligned with business functions and assigning dedicated data stewards to each domain for clear ownership and accountability. This strategy shifts from centralized data management to a decentralized model where data ownership and responsibility can be distributed across business domains, enabling teams closer to the data to make informed decisions about their domain-specific needs.

Distributed Data Architecture

The new architecture isolates domain-specific data loads to separate compute endpoints and creates independent data processing pipelines for each domain. This approach helps reduce cross-domain dependencies and conflicts that previously slowed development cycles, allowing teams to iterate and deploy changes without waiting for coordination across the entire organization.

Data Product Approach

Vanguard is curating data products on the data lake using Apache Iceberg format and leveraging AWS Glue for metrics computation and data lake integration. This approach treats data as products with defined SLAs and quality metrics, helping facilitate reliable, high-quality data delivery that downstream consumers can depend on with confidence.

Self-Service Analytics

The implementation enables domain teams to manage their complete data product lifecycle independently while maintaining enterprise governance standards. Vanguard provides comprehensive tools and systems for independent data management, allowing teams to innovate quickly without compromising data quality or security, ultimately accelerating time-to-insight across the organization.This evolution represents a natural progression from centralized data warehouse to multi-warehouse architecture, and finally to a fully distributed, domain-oriented data mesh that can scale with Vanguard’s continued growth.

Conclusion

Vanguard Financial Advisor Services’ journey demonstrates that scaling analytics is no longer about scaling a single warehouse bigger, but about architecting for workload isolation, independent scaling, and organizational growth.

By evolving from a single 2-node RA3 provisioned cluster to a multi-warehouse architecture using Amazon Redshift Serverless and Provisioned, Vanguard achieved measurable, production-grade outcomes:

  • 500,000+ monthly queries supported without ETL or dashboard contention
  • 100% ETL SLA adherence, with nightly pipelines completing before 9 AM
  • 25x growth in BI consumption (20 → 500+ Tableau dashboards) without performance degradation
  • 8x growth in analyst population (60 → 500+) enabled through workload isolation
  • 30x increase in ETL pipelines (20 → 600+) without re-architecting ingestion logic
  • Zero-copy Amazon Redshift data sharing across producer and consumer warehouses, minimizing data duplication and synchronization costs
  • Removal of 10-minute query limits, unlocking advanced exploratory and long-running analytics

Critically, these gains were not achieved by over-provisioning compute, but by right-sizing and specializing compute per workload, reserving capacity where demand was predictable (ETL) and using Amazon Redshift Serverless auto-scaling where demand was bursty (BI and ad-hoc analysis).

As Vanguard now progresses toward a domain-oriented data mesh, their experience reinforces a key lesson: Multi-warehouse architecture is a foundational enabler for organizational scale, data product ownership, and autonomous analytics.For organizations experiencing exciting growth in their data analytics requirements, Vanguard’s approach showcases the tremendous possibilities that await. With the right architecture and the help of AWS services, organizations can transform their data infrastructure to achieve remarkable improvements in performance, significant cost reductions, and unlock powerful new analytical capabilities that accelerate business value creation.

AWS encourages you to connect with your AWS Account Team to engage an AWS analytics specialist who can provide expert architectural guidance and tailored recommendations to help you achieve your data transformation goals.

© 2026 The Vanguard Group, Inc. and Amazon Web Services, Inc. All rights reserved. This material is provided for informational purposes only and is not intended to be investment advice or a recommendation to take any particular investment action.


About the authors

Alex Rabinovich

Alex Rabinovich

Alex is a Director of Data Engineering at Vanguard, aligned to Financial Advisory Services division. In this role, he leads large‑scale data engineering platforms and modernization initiatives, focusing on building reliable, scalable, and high‑performance data systems in the AWS cloud.

Anindya Dasgupta

Anindya Dasgupta

Anindya is a solutions architect in Vanguard’s Financial Advisor Services Technology division. He has over 25 years of experience building enterprise technology solutions to address complex business challenges. His work focuses on architecting and designing scalable, cloud‑native and data‑driven systems, with hands‑on contributions across application development, system integration, and proof‑of‑concept initiatives.

Vijesh Chandran

Vijesh Chandran

Vijesh is Head of Solution Design, overseeing the architecture and design of enterprise technology solutions that support critical business outcomes. His background spans data architecture on cloud‑native platforms, and data‑driven systems, with a strong focus on aligning technology design to business strategy. He plays a hands‑on role in guiding solution direction, integration patterns, and proof‑of‑concept initiatives.

Raks Khare

Raks Khare

Raks is a Senior Analytics Specialist Solutions Architect at AWS based out of Pennsylvania. He helps customers across varying industries and regions architect data analytics solutions at scale on the AWS platform. Outside of work, he likes exploring new travel and food destinations and spending quality time with his family.

Poulomi Dasgupta

Poulomi Dasgupta

Poulomi is a Senior Analytics Solutions Architect with AWS. She is passionate about helping customers build cloud-based analytics solutions to solve their business problems. Outside of work, she likes travelling and spending time with her family.

Amazon Redshift DC2 migration approach with a customer case study

Post Syndicated from Satoru Ishikawa original https://aws.amazon.com/blogs/big-data/amazon-redshift-dc2-migration-approach-with-a-customer-case-study/

This is a guest post by Satoru Ishikawa, Solutions Architect at Classmethod in partnership with AWS.

In April 2025, AWS announced the deprecation of Amazon Redshift DC2 instances, guiding users to migrate to either Redshift RA3 instances or Redshift Serverless. Redshift RA3 instances and Serverless adopt a design that separates storage and compute, offers new features such as data sharing, concurrency scaling for writes, zero-ETL , and cluster relocation.

In this post, we share insights from one of our customers’ migration from DC2 to RA3 instances. The customer, a large enterprise in the retail industry, operated a 16-node dc2.8xlarge cluster for business intelligence (BI) and ETL workloads. Facing growing data volumes and disk capacity limitations, they successfully migrated to RA3 instances using a Blue-Green deployment approach, achieving improved ETL query performance and expanded storage capacity while maintaining cost efficiency.

Amazon Redshift architecture types

Amazon Redshift offers two deployment options: Provisioned mode, where you choose the instance type and number of nodes and manage resizing as needed, and Redshift Serverless, which automatically provisions data warehouse capacity and intelligently scales the underlying resources. The following diagram compares these two architecture types.

Provisioned clusters require you to determine cluster size in advance, but you can optimize costs by purchasing Reserved Instances (RI) or scheduling pause and resume actions. Serverless automatically provisions resources as needed, with a pay-per-use model where you only pay for compute resources consumed. Both services support migration between each other and offer the same features including SQL, zero-ETL, and Federated Query capabilities. For specific pricing details, see Amazon Redshift pricing.

Provisioned clusters are suitable for large-scale, predictable workloads and offer automatic scaling based on queuing. Serverless provides management-free automatic scaling for variable workloads with AI-driven optimization that scales based on workload complexity and data volumes. For more details, refer to Comparing Amazon Redshift Serverless to an Amazon Redshift provisioned data warehouse.

Customer case study: Migration from DC2 instances

This section describes the customer’s migration from Amazon Redshift DC2 to RA3 instance types. The migration used a Blue-Green deployment approach that minimized downtime while achieving both cost optimization and performance improvement.

The customer’s workload had the following characteristics:

Use cases

The customer had the following key use cases for their Amazon Redshift deployment:

  1. Query via BI tool during business hours
    1. High volume of read queries
    2. Peak access during Mondays and beginning of months
  2. Data processing in early morning
    1. Concentrated write queries for data loading and transformation
  3. Steady-state workload characteristics
    1. Run queries more than 16 hours daily

Requirements

The customer had the following key requirements for their Amazon Redshift migration:

  1. Performance
    1. Use auto-scaling (such as concurrency scaling) during peak access periods
  2. Data size
    1. Disk capacity expansion needed
  3. Cost Management
    1. Easy budget prediction and management
    2. Utilize discount services for long-term usage
  4. Compatibility
    1. Maintain compatibility with existing applications and BI tools
    2. Avoid endpoint changes
  5. Availability
    1. Maximum downtime of 8 hours acceptable during migration
  6. Network
    1. Do not modify the existing 2-Availability Zone (AZ) subnet configuration
  7. When to migrate
    1. To be conducted during low-load days and hours
    2. Planned downtime possible within 8 hours

Key considerations in system design, implementation, and operation included extended operation hours, ease of budget prediction and management, cost optimization through Reserved Instances (RI), and maintaining compatibility with existing systems (avoiding endpoint changes). The customer evaluated Amazon Redshift Serverless, which offered attractive features such as a pay-per-use model, automatic scaling capabilities, and the potential for better price performance for variable workloads. While both Redshift Serverless and provisioned clusters could effectively support their workload patterns, the customer chose the provisioned model with RA3 nodes, leveraging their years of operational experience with provisioned environments, existing RI strategy, and established capacity planning approach.

Features of RA3 instance type

Built on the AWS Nitro System, RA3 instances with managed storage adopt an architecture that separates computing and storage, allowing independent scaling and separate billing for each component. These instances use high-performance SSDs for hot data and Amazon S3 for cold data, providing ease of use, cost-effective storage, and fast query performance. For more details, refer to Amazon Redshift RA3 instances with managed storage.

Migration prerequisites

The customer had the following migration prerequisites in place:

  • The customer used a Redshift cluster with 16 nodes of dc2.8xlarge configuration.
  • The customer chose a Blue-Green deployment approach for migration, where they would restore from a snapshot to RA3 instance type, enabling quick rollback if necessary.
  • The customer implemented cluster switching and rollback through endpoint switching using cluster identifier rotation.
  • Additionally, to improve performance with high concurrency, they transitioned the transaction isolation level from SERIALIZABLE ISOLATION to SNAPSHOT ISOLATION.

Cluster migration methods

There were two migration options available: Elastic Resize and Classic Resize.

Amazon Redshift’s Classic Resize functionality had been enhanced, for resizing to RA3 instance types, significantly reducing the write-unavailable period. Based on PoC testing, after initiating the resize, the cluster’s status was modifying for 16 minutes before it became available. Based on these results, the customer proceeded with the Classic Resize approach.

Cluster sizing

Sizing involved determining the instance type and number of nodes for the migration target. Sizing points considered workload characteristics such as CPU-intensive (queries using high CPU), I/O-intensive (queries with high data read/write), or both.When migrating from DC2 instance types, additional nodes might be required depending on workload requirements. Nodes were added or removed based on the computing requirements for necessary query performance.

Comparing configurations with similar cluster costs in terms of instance size and count, for a dc2.8xlarge 16-node cluster, the recommended configuration was 8 nodes of ra3.16xlarge. The following was the cost comparison in the Tokyo Region:

  1. Recommended: dc2.8xlarge 16-node cluster => ra3.16xlarge * 8-node cluster
    1. $97.52/h (6.095/h * 16 nodes) => $122.776/h (15.347/h * 8 nodes)
  2. Cost-focused: dc2.8xlarge 16-node cluster => ra3.16xlarge * 6-node cluster
    1. $97.52/h (6.095/h * 16 nodes) => $92.082/h (15.347/h * 6 nodes)

For this migration, the customer proceeded with a cost-efficient 6-node ra3.16xlarge cluster to stay within existing budget constraints. However, since this node count could face throughput limitations during certain times, they enabled concurrent scaling for the RA3 instance type to handle spike access.

Concurrency scaling provides up to 1 hour of free credits per day for each active cluster, accumulating up to 30 hours. On-demand usage fees apply when exceeding this free tier.While the customer chose to implement concurrency scaling, Elastic Resize to temporarily increase nodes during peak loads was also considered but rejected due to on-demand costs for additional nodes and the brief disconnection period during switching.

Managed storage cost

RA3 instances use Redshift Managed Storage (RMS), which is charged at a fixed GB-month rate. The customer’s approximately 2 TB of data required including storage costs in the estimates. For pricing details, see Amazon Redshift pricing.

Migration step from DC2 to RA3

After creating an RA3 cluster from the DC2 cluster’s snapshot, the customer swapped the cluster identifiers. The following diagram shows this process.

  1. Take a snapshot of the current DC2 cluster.
  2. Restore RA3 cluster from the snapshot with a different cluster identifier (Classic Resize)
  3. Swap the cluster identifiers between the current DC2 cluster and the new RA3 cluster.

If any issues arise after the cluster switch, you can quickly roll back by returning the original DC2 cluster to its original cluster identifier.

Note: Restore from a snapshot

Running the restore operation using CLI commands is recommended to minimize operational errors and ensure reproducibility. The following is a sample command.

aws redshift restore-from-cluster-snapshot \
--cluster-identifier for-ra3-20250207 \
--snapshot-identifier cm-cluster-for-ra3-20250207 \
--cluster-subnet-group-name cm-cluster \
--vpc-security-group-ids sg-1234567a sg-2345678b sg-3456789c \
--cluster-parameter-group-name cm-cluster \
--node-type ra3.16xlarge \
--number-of-nodes 6 \
--port 5439 \
--no-publicly-accessible \
--enhanced-vpc-routing \
--availability-zone ap-northeast-1a \
--preferred-maintenance-window sat:17:00-sat:17:30 \
--automated-snapshot-retention-period 14 \
--iam-roles 'arn:aws:iam::123456789012:role/AmazonRedshift-CommandsAccessRole' 'arn:aws:iam::123456789012:role/AmazonRedshift-Spectrum' \
--maintenance-track-name current

Production migration duration

The time required for the restore and classic resize steps can vary significantly depending on data volume and target cluster specifications. The customer conducted a rehearsal beforehand to measure the actual required time.

Test results

Before the production migration, the customer created a test cluster by restoring a snapshot to the RA3 instance type. While Redshift Test Drive is typically useful for workload testing, this customer faced unique constraints: enabling audit logging in their production cluster would require configuration changes, cluster restarts, and complex approval processes under their strict change management policies. To address this, they developed a custom load testing tool that captured workload patterns using Amazon Redshift system views (SYS_QUERY_HISTORY and SYS_QUERY_TEXT), which maintain 7 days of query history. The tool replayed 55,755 historical queries with 50-way parallelism against both DC2 and RA3 clusters, comparing metrics including query execution time, CPU utilization, and disk I/O. Query result caching was disabled during testing to ensure accurate comparisons.

BI query performance

BI queries were tested using the custom load testing tool. The results represent the average execution time from 15 test runs of 55,755 queries executed with 50-way parallelism. Without concurrency scaling, the dc2.8xlarge 16-node cluster averaged 45.82 seconds per query, while the ra3.16xlarge 6-node cluster averaged 91.30 seconds. This indicated that RA3 instances showed longer execution times for short and medium queries in a direct migration without optimizations. However, enabling concurrency scaling improved RA3 performance progressively. With concurrency scaling enabled at maximum 2 clusters, the ra3.16xlarge 6-node cluster achieved an average of 72.48 seconds per query, a 21% improvement over the non-scaled configuration.

Node Type / Number of nodes Average Query Time
ra3.16xlarge 6-node cluster 72.48 seconds

ETL query performance comparison

For long-running ETL queries (execution time greater than 10 minutes), the RA3 cluster demonstrated better performance than DC2. These results represented a direct migration of the customer’s workload with no optimizations applied.

  • For the Large-scale data load workload 1, the ra3.16xlarge cluster completed the query 28% faster than the dc2.8xlarge cluster (41 minutes vs. 57 minutes).
  • For the Complex transformation workload 1, the ra3.16xlarge cluster was 23% faster (1 hour 1 minute vs. 1 hour 20 minutes).

These results indicated that the RA3 node type was more performant for time-intensive data loading and transformation tasks. The higher CPU utilization values for RA3 suggested more effective compute resource usage.

Node Type / Number of nodes Average Query Time MAXCPU%
ra3.16xlarge 6-node cluster 41 mins 09 seconds 11:45
dc2.8xlarge 16-node cluster 57 mins 07 seconds 10:85
Node Type / Number of nodes Average Query Time MAXCPU%
ra3.16xlarge 6-node cluster 1 hour 01 mins 33 seconds 74:23
dc2.8xlarge 16-node cluster 1 hour 20 mins 36 seconds 53:58

Performance tuning

Based on the test results, the customer identified that RA3 showed longer execution times for short and medium BI queries but faster performance for long-running ETL queries compared to DC2. To optimize overall performance, they focused on identifying slow queries and frequently referenced tables, prioritizing optimizations with the highest impact.

Performance tuning strategy

The customer considered several optimization strategies to leverage RA3’s architectural advantages. One key strategy involved pre-processing ad-hoc short and medium query workloads during low-load periods, creating pre-processed tables or materialized views for queries that repeatedly performed joins, aggregations, filters, and projections. RA3’s separated compute and storage architecture, with cost-effective large-scale storage, supported this approach.

Converting regular views to materialized views

Analysis of slow queries revealed the use of joins in views, and frequently referenced tables were being accessed multiple times through these views. As a countermeasure, the customer replaced frequently used regular views with materialized views, removing unnecessary data ranges and redundant columns.

Amazon Redshift supports incremental updates of materialized view contents via the REFRESH MATERIALIZED VIEW command, enabling efficient data updates.

Materialized views and query rewrite

By converting regular views to materialized views, existing queries may be automatically optimized through the “query rewrite” feature provided by the query planner. For more details, refer to “Automatic query rewriting to use materialized views“.

Automatic tuning with AutoMV

On the DC2 cluster, disk utilization consistently exceeded 80%, which disabled the AutoMV feature due to insufficient disk space. With RA3’s expanded storage, automatic tuning through AutoMV became possible, leading to further performance improvements. For more details about AutoMV, refer to Automated materialized views.

Performance tuning results

After applying these optimizations, the customer achieved the following results:

  • Maintained existing performance while controlling cost increases
  • Achieved higher CPU utilization while maintaining throughput
  • Enhanced dynamic throughput during peak load periods using concurrency scaling’s automatic scaling

Conclusion

In this post, you learned how a large retail enterprise successfully migrated from Amazon Redshift DC2 to RA3 instances. The Blue-Green deployment approach enabled a safe migration with quick rollback capability, while the separated compute and storage architecture of RA3 provided flexibility to handle growing data volumes. Although RA3 showed different performance characteristics for short BI queries compared to DC2, the customer achieved significant improvements in long-running ETL query performance (up to 28% faster for data loads and 23% faster for complex transformations). By leveraging RA3-specific features such as materialized views and AutoMV, they optimized overall query performance while maintaining cost efficiency through Reserved Instances and concurrency scaling.

To continue your RA3 migration journey, see Best practices for upgrading from Amazon Redshift DC2 to RA3 and Amazon Redshift Serverless and Resize Amazon Redshift from DC2 to RA3 with minimal or no downtime for additional guidance and best practices.


About the authors

Satoru Ishikawa

Satoru Ishikawa

Satoru specializes in data analytics and AI consulting, focusing on Amazon SageMaker and multi-cloud. He also develops the backend for Classmethod’s “Members,” driving digital transformation through advanced data and AI capabilities.

Junpei Ozono

Junpei Ozono

Junpei drives technical market creation for data and AI solutions, working closely with global teams to build scalable GTM motions. His expertise spans modern data architectures — Data Mesh, Data Lakehouse, and AI — helping customers accelerate their cloud transformation with AWS.

Kinesis On-demand Advantage saves 60%+ on streaming costs

Post Syndicated from Pratik Patel original https://aws.amazon.com/blogs/big-data/kinesis-on-demand-advantage-saves-60-on-streaming-costs/

Amazon Kinesis Data Streams is a serverless streaming data service that helps you capture, process, and store streaming data at any scale. On November 4, 2025, Amazon Kinesis Data Streams introduced On-demand Advantage mode, a capability that enables on-demand streams to handle instant throughput increases at scale and cost optimization for consistent streaming workloads. Historically, you had to choose between provisioned mode, which required managing stream capacity, and on-demand mode, which automatically scaled capacity, but this new offering removes the need to think about stream type at all.

In this post, we show three real-world scenarios comparing different usage patterns and demonstrate how On-demand Advantage mode can optimize your streaming costs while maintaining performance and flexibility. To have a meaningful comparison, we ran simulations in two separate AWS accounts: one with On-demand Standard mode and another with On-demand Advantage mode enabled at the account level. Both deployments maintained identical stream configurations, shard allocations, and ingest patterns, providing a comparison of the billing impact for all of the following scenarios.

All prices displayed in this post are from the us-east-1 Region.

Breaking down On-demand Advantage savings

Now let’s go into more details. In the first post, we talked about the warm throughput feature and how you can use it to warm streams to handle gigabytes or millions of records per second with On-demand Advantage mode. Next, we illustrate how different streaming use cases operate most cost efficiently with the On-demand Advantage mode, while maintaining performance and flexibility.

Enabling On-demand Advantage mode in the account level gives you cost savings across many dimensions compared to On-demand Standard, here are some notable ones:

  • Provides at least 60% savings by committing to an account to stream at least 25 MiBps usage in the AWS Region. The minimum commitment is about $100 a day based on AWS N. Virginia Regions’ public pricing.
  • Enhanced Fan Out consumers usage is priced 68% lower and you can have up to 50 per stream, compared to 20 per stream without On-demand Advantage.
  • Extended retention usage is priced 77% lower when using data storage beyond 24 hours.
  • No minimum per-stream fixed charge, so you can use as many streams as you need without incurring a higher cost.

Why choose Kinesis On-demand Advantage: 60%+ savings

We evaluated Amazon Kinesis Data Streams On-demand using both Standard and Advantage modes by deploying 10 streams and generating a sustained ingest throughput of 100 MiBps across all streams. This scenario models an ecommerce company streaming user clickstream data to generate real-time insights. The simulation was run over two days in two separate AWS accounts with the different modes. On day one, we maintained a steady ingestion rate of 100 MiBps. On day two, in anticipation of a holiday sales event, we increased the warm throughput capacity by 10X across all 10 streams while keeping the actual ingest rate constant at 100 MiBps. Each stream is ingesting 10 MiBps with a total of 100 MiBps across all On-demand streams.

scenario-1-kds-od-metrics-1

On the second day, we enabled warm throughput of 100 MiBps in an account with On-demand Advantage mode configured. With warm throughput, you can proactively pre-scale your streams when expecting a traffic surge.

scenario-1-kds-od-streams

On-demand Standard Cost Explorer:

scenario-1-kds-od-cost

On-demand Advantage Mode Cost Explorer:

scenario-1-kds-od-advantage-cost

This use case is cost effective in On-demand Advantage because of the consistent data throughput traffic and the need to use multiple streams. You can see zero stream hour charges in Advantage mode but in On-demand Standard mode there is a $0.04 charge per stream hour. Additionally, the incoming bytes charge for Advantage mode is $0.032 per GB and On-demand Standard is $0.08 per GB. The On-demand Standard configuration generated total costs of $2,071.75 over the 48-hour period, which comes out to $1,037 per day, which is $378,505 annually. The identical workload running with On-demand Advantage mode costs $823.44 for the same 48-hour period, approximately $412 per day resulting in an annual cost of $150,380. There is also no additional cost for warm throughput in On-demand Advantage. This translates to a 60% cost reduction, which means annual savings of $228,125 for this single workload.

Scenario 2: 15 MiBps throughput and extended retention across 10 streams in one account

A healthcare company requires a two-day data retention period to ensure continuous replay-ability, and regulatory requirements mandate that different types of data be stored in separate streams. To emulate this scenario, we deployed 10 Amazon Kinesis Data Streams configured in both On-demand Standard and Advantage modes, each with an extended 48-hour retention period. The extended retention allows downstream systems to reprocess data and recover from transient failures. This use case is also expected to be cost-efficient in On-demand Advantage mode due to the need for multiple streams and retention, a point we explore further in the following cost breakdown image.

scenario-2-kds-od-streams

We generated consistent ingest traffic of 15 MiBps distributed across 10 streams for 48 hours to evaluate costs. On-demand Standard Mode:

scenario-2-kds-od-cost

On-demand Advantage Mode enabled:

scenario-2-kds-od-advantage-cost

The Cost Explorer screenshot gives us a side-by-side view of the pricing between On-demand Standard and Advantage modes. On-demand Standard came out to a daily cost of $176, annually becoming $64,240. For On-demand Advantage, the daily cost was $104 with an annual cost of $37,960. As a result, we achieve a 41% savings with Advantage mode despite operating with 15 MiBps throughput and implementing extended retention. Standard mode has an additional cost of $0.10 per GB of data stored beyond 24 hours, up to 7 days. Advantage mode costs an additional $0.023 per GB month (beyond 24 hours, up to 365 days) resulting in cost optimization for data storage. This scenario shows how Advantage mode delivers cost benefits across a broader range of workloads. When enabling Advantage mode, you commit to paying for usage at a minimum of 25 MiBps. However, in our simulation with 15 MiBps throughput, we found that customers still achieved significant cost savings. With Advantage mode, as long as you’re ingesting 10 MiBps or more, you will experience lower costs compared to Standard mode, even when committing to the 25 MiBps threshold.

Scenario 3: 1 Kinesis Data Stream with 10 enhanced fan-out consumers

In a microservices architecture, multiple services might need to read from the same data stream concurrently and with low latency. As the system evolves, additional Enhanced Fan-Out (EFO) consumers might be added to support new analytics use cases and derive deeper insights from the streaming data pipeline.

Next, we evaluate the cost comparison of using EFO consumers with On-demand Advantage mode in a microservices architecture. Over a 24-hour period, we tested a single Amazon Kinesis Data Stream with standard 24-hour retention, connected to 10 AWS Lambda functions configured as EFO consumers.

To assess the cost impact of multiple EFO consumers accessing the same stream simultaneously, we generated a consistent ingest rate of 25 MiBps throughout the evaluation period. The following chart shows a Kinesis Data Stream with Enhanced Fan-Out Consumers with 25MiBps payload.

scenario-3-kds-od-metrics-1

The following charts show the cost difference between On-demand Standard mode vs On-demand Advantage mode enabled:

scenario-3-kds-od-cost-comparision

Our Cost Explorer analysis demonstrates that even with multiple EFO consumers, the On-demand Advantage mode resulted in lower overall costs compared to the On-demand Standard mode. The On-demand Standard cost was $1,266 while the Advantage mode cost was $419. For the same workload characteristics, we observed approximately 67% savings with annual savings of $309,155. This is especially important for organizations building event-driven architectures where multiple services need independent, real-time access to streaming data. Enhanced fanout data retrievals are $0.016 per GB per consumer in Advantage mode compared to the $0.05 of Standard mode. Now that we’ve discussed which workloads we recommend for Kinesis On-demand Advantage, let’s turn to workloads that we recommend for On-demand Standard mode. Highly spiky, unpredictable workloads with low sustained throughput (under 10 MiBps) are recommended candidates for On-demand Standard. With these workloads, customers can let the stream automatically scale throughput capacity without committing to a consistent throughput usage.

On-demand Advantage compared to Provisioned mode:

With both On-demand Standard and On-demand Advantage modes available, customers no longer need to rely on Provisioned mode. Instead of continuously managing capacity to balance performance and cost, on-demand streams offer a streamlined pricing model and greater ease of use. Additionally, customers running provisioned workloads that operate many streams or use features such as Extended Retention and Enhanced Fan-Out should strongly consider migrating to On-demand Advantage. Data streams use the same underlying infrastructure, regardless of which mode you choose, so there is no difference in availability or reliability between On-demand Advantage and provisioned.

  1. For ensuring streams can support instant scaling at no extra cost and can reactively scale when needed, On-demand Advantage is a better fit than Provisioned mode, because warm throughput doesn’t incur an additional cost.
  2. Even at a large scale (Gbps), On-demand Advantage’s pay-for-actual-data-usage billing is competitive.
  3. On-demand Advantage has a much lower cost to use EFO (for high fanout needs like microservices) and Extended Retention.

To help you compare, you can use the Kinesis console to check if your account’s top 200 provisioned streams are a good fit to use On-demand Advantage instead.

kds-od-analyze-account-usage

Conclusion

In this post, we explored three real-world scenarios demonstrating how Amazon Kinesis Data Streams On-demand Advantage mode delivers significant cost savings while maintaining performance and flexibility. On-demand Advantage provides you with the best performance at scale, and together with On-demand Standard, Kinesis Data Streams offers you a streamlined way to use and most cost-effective streaming solution for any streaming use case. If your workloads consistently stream at least 10 MBps, fan out to two or more consumers, retain data for more than 24 hours, or operate hundreds of streams, On-demand Advantage is the most cost-effective mode. For all your other workloads, On-demand Standard mode is a great fit. Whether you’re streaming millions of records or gigabytes of data per second from diverse producers to consumers, Kinesis Data Streams has you covered. We look forward to hearing how you and your teams take advantage of Kinesis Data Streams On-demand Advantage to bring real-time insights to your organizations and processes.


About the authors

Sandhya Khanderia

Sandhya Khanderia

Sandhya is a Sr. Technical Account Manager and Data Analytics Specialist at AWS. With deep expertise in analytics and data services, Sandhya specializes in helping organizations optimize their cloud architectures for performance, scalability, and cost-efficiency.

Pratik Patel

Pratik Patel

Pratik is a Sr. Technical Account Manager and streaming analytics specialist. He works with AWS customers and provides ongoing support and technical guidance to help plan and build solutions using best practices and proactively keep customers’ AWS environments operationally healthy.

Varsha Palepu

Varsha Palepu

Varsha is a Solutions Architect at AWS, where she helps small and medium businesses innovate and build on AWS. She serves on the streaming team to create technical content and empower customers to achieve their goals on the cloud.

Kalyan Janaki

Kalyan Janaki

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

Adding a voice layer to WhatsApp conversations with AWS End User Messaging

Post Syndicated from Pavlos Ioannou Katidis original https://aws.amazon.com/blogs/messaging-and-targeting/adding-a-voice-layer-to-whatsapp-conversations-with-aws-end-user-messaging/

Businesses around the world use WhatsApp as a primary channel to connect with customers. It’s familiar, trusted, and effective for everything from booking confirmations to customer support. But most of these conversations are still text-only. For many customers, text is fast and efficient. Yet there are times when typing is inconvenient, slow, or less effective at conveying nuance. In those moments, voice messages can transform the interaction — making it faster, more inclusive, and more human.

With AWS End User Messaging, businesses can now enable both voice note input and voice note responses on WhatsApp. Customers send a voice note, and a bot can respond with a natural-sounding voice note reply. Note: This solution processes asynchronous voice notes (recorded audio messages), not real-time voice calls. In this blog post, we explore why voice notes matter, where they make a difference, and how AWS helps you enable them through a sample voice note messaging solution.

Watch an end to end demo here.

Why voice notes matter in customer messaging

Text remains essential, but research shows that voice notes adds unique advantages:

  • Richer communication: Voice carries tone, urgency, and emotion — reducing misunderstandings and helping businesses respond more appropriately (Preply survey).
  • Natural and fast: Speaking is up to three times faster than typing on mobile devices, especially when users are on the go (Sherry Ruan, Jacob O. Wobbrock, Kenny Liou, Andrew Ng, and James A. Landay. 2018. Comparing Speech and Keyboard Text Entry for Short Messages in Two Languages on Touchscreen Phones. Proc. ACM Interact. Mob. Wearable Ubiquitous Technol. 1, 4, Article 159 (December 2017), 23 pages. https://doi.org/10.1145/3161187).
  • Accessibility and inclusivity: Voice lowers barriers for people with limited literacy or visual impairments. Elderly customers or those with difficulty reading long text messages benefit significantly.
  • Context-driven preference: A YouGov study across 17 markets found that while text is still preferred overall, a notable share of users choose both text and audio depending on situation (YouGov survey).

Where voice notes make a difference

Voice messaging is especially useful when speaking feels more natural than typing—helping customers communicate in ways that fit their situation and needs.

  • Elderly customers – easier to listen than to read.
  • Field workers or drivers – easier to speak than to type while working.
  • Healthcare – patients can describe symptoms naturally by voice.
  • Hospitality and reservations – “Book a table for 7 pm” is faster to say than to navigate online calendar.
  • Customer support escalation – complex issues are often resolved more quickly with a voice exchange.

Voice notes don’t replace text. It complements it — giving customers the flexibility to communicate in the way that best suits their context.

AWS End User Messaging and WhatsApp

AWS End User Messaging is a managed AWS service that enables businesses to send and receive messages across multiple channels, including WhatsApp, SMS, MMS (US only), outbound voice, and push notifications.

When you use AWS End User Messaging for WhatsApp, you benefit from AWS’s global scale, resilience, and security. Inbound WhatsApp messages are automatically published to an Amazon SNS topic, enabling the integration with other AWS services such as Amazon SQS queues, AWS Lambda functions or Amazon Bedrock for downstream processing.

This flexibility is also what makes voice-to-voice messaging possible. Businesses can process inbound voice messages with Lambda, apply speech-to-text and text-to-speech services like Amazon Transcribe and Amazon Polly, or integrate third-party models such as Whisper through the AWS Marketplace for Amazon Bedrock.

Voice notes messaging solution

To demonstrate how voice can be enabled on WhatsApp, check out the AWS CDK sample project:  GitHub – WhatsApp Voice Notes Messaging

The solution shows how to:

  • Receive a WhatsApp voice note through AWS End User Messaging.
  • Transcribe the voice input to text.
  • Process it with conversational bot logic.
  • Convert the response back into a natural-sounding voice note.
  • Send the reply to the user on WhatsApp.

You can enable inbound only, outbound only, or a full voice-to-voice  notes loop depending on your requirements.

Getting started

The complete solution is available as an open-source AWS CDK project. To get started, you’ll need:

Implementation

Clone the repository and deploy the solution:

git clone https://github.com/aws-samples/sample-whatsapp-voice-to-voice-messaging
cd sample-whatsapp-voice-to-voice-messaging
npm install

Before deploying, you’ll need to configure your WhatsApp phone number ID in the CDK context or parameters. The deployment will prompt you for this configuration, or you can set it in the cdk.json file. Once configured, deploy with:

cdk deploy

The CDK stack automatically provisions all required AWS resources including Lambda functions, SNS topics, S3 buckets, and IAM roles.

Clean up

To remove all resources and avoid ongoing charges:

cdk destroy

For detailed architecture diagrams, configuration options, and step-by-step setup instructions, visit the GitHub repository.

Conclusion

Customers are already using voice notes in their personal WhatsApp conversations. Bringing that same option into business communication makes customer interactions more natural, inclusive, and efficient.

With AWS End User Messaging and its WhatsApp channel, you can add voice alongside text without changing how customers connect to you. And with the sample CDK project, you can try it out today, experiment, and extend it for your own business needs.

Explore the project here: AWS Sample – WhatsApp Voice Notes Messaging


About the authors

Zero-ETL integrations with Amazon OpenSearch Service

Post Syndicated from Omama Khurshid original https://aws.amazon.com/blogs/big-data/zero-etl-integrations-with-amazon-opensearch-service/

Amazon OpenSearch Service is a fully managed service that reduces operational overhead, provides enterprise-grade security, high availability, and scalability, and enables you to quickly deploy real-time search, analytics, and generative AI applications. OpenSearch itself is an open-source, distributed search and analytics suite that supports a wide range of use cases, including real-time monitoring, log analytics, and full-text search. OpenSearch Service offers zero-ETL integrations with other Amazon Web Service (AWS) services, enabling seamless data access and analysis without the need for maintaining complex data pipelines.

Zero-ETL refers to a set of integrations designed to minimize or eliminate the need to build traditional extract, transform, load (ETL) pipelines. Traditional ETL processes can be time-consuming and difficult to develop, maintain, and scale. In contrast, zero-ETL integrations allow direct, point-to-point data movement and can also support querying across data silos without physically moving the data.

In this post, we explore various zero-ETL integrations available with OpenSearch Service that can help you accelerate innovation and improve operational efficiency. We cover following types of integrations, their key features, architecture, benefits, pricing, limitation and some general best practices.

  1. Log and storage integrations
  2. Database integrations

The following diagram illustrates the zero-ETL integration architecture in AWS, showing how various AWS services feed data into OpenSearch Service and its associated dashboards:

Zero ETL with Amazon OpenSearch Service

Zero-ETL integration with Amazon S3

Amazon OpenSearch Service direct queries with Amazon S3 provides a zero-ETL integration to reduce the operational complexity of duplicating data or managing multiple analytics tools by enabling you to directly query their operational data, reducing costs and time to action.

Key features of this integration include:

  1. In-place querying: You can use rich analytics capabilities of OpenSearch Service SQL and PPL directly on infrequently-queried data stored outside of OpenSearch Service in Amazon S3.
  2. Selective data ingestion: You can choose which data to bring into OpenSearch Service for detailed analysis, optimizing costs and speeding up queries with indexes like skipping or covering indexes.

The zero-ETL integration with Amazon S3 supports OpenSearch Service. For more information on architecture and feature see the post Modernize your data observability with Amazon OpenSearch Service zero-ETL integration with Amazon S3.

In log analytics use cases, we categorize operational log data into two types:

  • Primary data includes the most recent and frequently accessed logs used for real-time monitoring and analysis.
  • Secondary data consists of historical logs that are accessed less frequently but retained for compliance or trend analysis.

You can offload infrequently queried data, such as archival or compliance data, to Amazon S3. With direct query, you can analyze analytics from Amazon S3 without data movement or duplication. However, query performance in OpenSearch Service might slow down when you’re accessing external data sources due to factors like network latency, data transformation, or large data volumes. You can optimize your query performance by using OpenSearch indexes, such as a skipping index, covering index, or materialized view.

While Amazon S3 direct query integration with OpenSearch Service provides on-demand access to data stored in Amazon S3, it is important to remember that OpenSearch’s alerting, monitoring, anomaly detection, and security analytics capabilities can only operate on data that has been explicitly ingested into OpenSearch Service indices. These capabilities would not work with direct query with Amazon S3. However, it will work if the data is indexed with covering or materialized index.

Benefits

With direct queries with Amazon S3, you no longer need to build complex ETL pipelines or incur the expense of duplicating data in both OpenSearch Service and Amazon S3 storage. You also save time and effort by not having to move back and forth between different tools during your analysis.

Pricing

OpenSearch Service separately charges for the compute needed to query your external data in addition to maintaining indexes in OpenSearch Service. Costs for Direct Query is based on the data volume scanned, query execution time, query frequency and frequency with which the indexed data in OpenSearch is kept updated. For more information, see Amazon OpenSearch Service Pricing.

Considerations

In case you are using OpenSearch service to query directly data on Amazon S3, consider the limitations with Direct Query.

Best practices

These are some general and Amazon S3 recommendations for using direct queries in OpenSearch Service. For more information, see Recommendations for using direct queries in Amazon OpenSearch Service.

  • Use the COALESCE SQL function to handle missing columns and ensure results are returned.
  • Use limits on your queries to ensure you aren’t pulling too much data back.
  • If you plan to analyze the same dataset many times, create an indexed view to fully ingest and index the data into OpenSearch Service and drop it when you have completed the analysis.
  • Drop acceleration jobs and indexes when they’re no longer needed.
  • Ingest data into Amazon S3 using partition formats of year, month, day, hour to speed up queries.
  • When you build skipping indexes, use Bloom filters for fields with high cardinality and min/max indexes for fields with large value ranges. Bloom filters are a space efficient probabilistic data structure that lets you quickly check whether an item is possibly in a set. For high-cardinality fields, consider using a value-based approach to improve query efficiency.
  • Use Index State Management to maintain storage for materialized views and covering indexes.

Zero-ETL integration with Amazon CloudWatch Logs

Amazon CloudWatch Logs serves as a centralized monitoring and storage solution for log files generated across various AWS services. This unified logging service offers a highly scalable platform where all your logging data converges into one manageable system. It provides comprehensive functionality for log management, including real-time viewing, pattern searching, field-based filtering, and secure archival capabilities. By presenting all logs chronologically in a unified stream, CloudWatch Logs eliminates the complexity of managing multiple log sources, transforming diverse logging data into a coherent, time-ordered sequence of events.

The zero-ETL integration between Amazon CloudWatch and Amazon OpenSearch Service enables direct log analysis and visualization while avoiding data redundancy, thereby reducing both technical complexity and costs. You can now leverage two additional query languages alongside the existing CloudWatch Logs Insights QL when using CloudWatch Logs, while as an OpenSearch user, you gain the ability to query CloudWatch logs directly.

Review New Amazon CloudWatch and Amazon OpenSearch Service launch an integrated analytics experience, to explore how the integration works between OpenSearch Service and Amazon CloudWatch Logs.

Benefits

  • The enhanced CloudWatch Logs Insights console now incorporates OpenSearch PPL and SQL functionality. Users can perform complex log analysis using SQL JOIN operations and various functions (including JSON, mathematical, datetime, and string operations). The PPL option provides additional data filtering and analysis capabilities.
  • The integration offers ready-to-use dashboards for various AWS services like Amazon Virtual Private Cloud (VPC), AWS CloudTrail, and AWS Web Application Firewall (WAF). These pre-configured visualizations enable quick insights into metrics such as flow patterns, top users, data transfer volumes, and temporal analysis, without requiring manual dashboard configuration.
  • You can now analyze CloudWatch logs through OpenSearch UI Discover and execute SQL and PPL queries. At the writing of this post, the query execution is limited to 50 log groups.
  • The direct access and analysis of CloudWatch data within OpenSearch Service removes the need for traditional ETL processes, eliminates separate data ingestion pipelines and avoids data duplication. This streamlined approach significantly reduces both storage expenses and operational complexity. It delivers a more efficient data management solution that simplifies the entire workflow while maintaining cost-effectiveness.

Pricing

When you use OpenSearch Service direct queries, you incur separate charges for OpenSearch Service and the resource used to process and store your data on Amazon CloudWatch Logs. As you run direct queries, you see charges for OpenSearch Compute Units (OCUs) per hour, listed as DirectQuery OCU usage type on your bill.

  • For interactive queries, OpenSearch Service handles each query with a separate pre-warmed job, without maintaining an extended session.
  • For indexed view queries, the indexed data is stored in an OpenSearch Serverless collection where you are charged for data indexed (IndexingOCU), data searched (SearchOCU), and data stored in GB.

You can find a pricing example on running an OpenSearch dashboard from either OpenSearch UI or CloudWatch Logs (pricing example n°7).

For more pricing information, see Amazon OpenSearch Service Direct Query pricing.

Considerations

In addition to the OpenSearch Service “direct queries” general limitations, if you are direct querying data in CloudWatch Logs, the following limitations apply:

  • The direct query integration with CloudWatch Logs is only available on OpenSearch Service collections and the OpenSearch user interface.
  • OpenSearch Serverless collections have networked payload limitations of 100 MiB.
  • CloudWatch Logs supports VPC Flow Logs, CloudTrail, and AWS WAF dashboard integrations installed from the console.

Best practices

Besides the general recommendations of OpenSearch Service direct querying, when using OpenSearch Service to direct query data in CloudWatch Logs, the following is recommended:

  • Specify the log group names within logGroupIdentifier in logGroups command to query multiple log groups in one query, see Multi-log group functions.
  • Enclose certain fields in backticks to successfully query them when using SQL or PPL commands. Backticks are needed for fields with special characters, such as `@SessionToken` or `LogGroup-A` (non-alphabetic and non-numeric). Refer to CloudWatch Logs Recommendations to see an example.

Zero-ETL integration with Amazon DynamoDB

Amazon DynamoDB zero-ETL integration with OpenSearch Service lets you perform a search on your DynamoDB data by automatically replicating and transforming it without custom code or infrastructure. This zero-ETL integration uses Amazon OpenSearch Ingestion to synchronize data between Amazon DynamoDB and OpenSearch Service cluster or OpenSearch Serverless collection within seconds of it being available.

It uses DynamoDB export to Amazon S3 to create an initial snapshot to load into OpenSearch Service. After the snapshot has been loaded, the plugin uses DynamoDB Streams to replicate any further changes in near real time. Turn on point-in-time recovery (PITR) for export and the DynamoDB Streams feature for ongoing replication.

This feature allows you to capture item-level changes in your table and push the changes to a stream. Every item in tables is processed as an event in OpenSearch Ingestion and can be modified with processors. You can also specify index mapping templates within ingestion pipelines to ensure that your Amazon DynamoDB fields are mapped to the correct fields in your OpenSearch indices.

To learn more, see DynamoDB zero-ETL integration with Amazon OpenSearch Service in the AWS documentation.

When configuring zero-ETL between DynamoDB and OpenSearch Service, consider the differences between the data models. You have the following options with data layout:

  1. Passthrough: Each item in DynamoDB table is directly mapped to one document in OpenSearch Index.
  2. Routing: A single DynamoDB table mapped to multiple OpenSearch Service indices. In DynamoDB, it is common to store denormalized data in one table to optimize for access patterns. For example, a single DynamoDB table containing both customer profiles and order information can be routed to separate OpenSearch Service indices:
    • Customer attributes → ‘customers’ index
    • Order attributes → ‘orders’ index

    You can achieve this by using the conditional routing feature in the OpenSearch ingestion pipeline.

  3. Merge: In some use cases, you need to combine data from multiple DynamoDB tables into a single OpenSearch index. You can use AWS Lambda integration with OpenSearch Ingestion to perform lookups on other DynamoDB tables and merge data from multiple DynamoDB tables.

Pricing

There is no additional cost to use this feature apart from the cost of the existing underlying components, including OpenSearch Ingestion charges OpenSearch Compute Units (OCUs) which is used to replicate data between Amazon DynamoDB and OpenSearch Service. Furthermore, this feature uses Amazon DynamoDB Streams for the change data capture (CDC), and you incur the standard costs for Amazon DynamoDB Streams.

Considerations

Consider the following limitations when you set up an OpenSearch Ingestion pipeline for DynamoDB:

  • At the writing of this post, the OpenSearch Ingestion integration with DynamoDB doesn’t support cross-Region and cross-account ingestion.
  • An OpenSearch Ingestion pipeline supports only one DynamoDB table as its source.

Best practices

For complete information, see Best practices for working with DynamoDB zero-ETL integration and OpenSearch Service

Integration with Amazon Aurora and Amazon RDS

Amazon RDS and Amazon Aurora integration with OpenSearch Service eliminates complex data pipelines and enables near real-time data synchronization between Amazon Aurora and Amazon RDS databases (including RDS for MySQL and RDS for PostgreSQL) with advanced search capabilities on transactional databases. You can use an OpenSearch Ingestion pipeline with Amazon RDS or Amazon Aurora to export existing data and stream changes (such as create, update, and delete) to OpenSearch Service domains and collections. The OpenSearch Ingestion pipeline incorporates change data capture (CDC) infrastructure to provide a high-scale, low-latency way to continuously stream data from Amazon RDS or Amazon Aurora.

This automated process keeps your data consistently up to date in OpenSearch Service, making it readily available for search and analysis purpose. The pipeline ensures data consistency by continuously polling or receiving changes from the Amazon Aurora cluster or Amazon RDS and updating the corresponding documents in the OpenSearch index. OpenSearch Ingestion supports end-to-end acknowledgement to ensure data durability. An OpenSearch Ingestion pipeline also maps incoming event actions into corresponding bulk indexing actions to help ingest documents. This keeps data consistent, so that every data change in Amazon RDS is reconciled with the corresponding document changes in OpenSearch.

For details on the architecture, refer to Integrating Amazon OpenSearch Ingestion with Amazon RDS and Amazon Aurora. To get started, refer to OpenSearch Ingestion pipeline with Amazon RDS or Using an OpenSearch Ingestion pipeline with Amazon Aurora.

Pricing

There is no additional charge for using this feature beyond the cost of your existing underlying resources, such as OpenSearch Service, OpenSearch Ingestion pipelines (OCUs), and Amazon RDS or Amazon Aurora. Additional costs may include storage used for enabling enhanced binlogs for MySQL and WAL logs for PostgreSQL for change data capture. You also incur storage costs for snapshot exports from your database to Amazon S3 used for the initial data.

Considerations

Consider the following limitations when you set up the integration for Amazon RDS or Amazon Aurora:

  • Support both Aurora MySQL or RDS for MySQL (8.0 and above) and Aurora PostgreSQL or RDS for PostgreSQL (16 and above).
  • Requires same-Region and same-account deployment, primary keys for optimal synchronization, and currently has no data definition language (DDL) statement support.
  • The integration only supports one Aurora PostgreSQL database per pipeline.
  • The existing pipeline configuration can’t be updated to ingest data from a different database and/or a different table. To update the database and/or table name of a pipeline, stop the pipeline and restart it with an updated configuration or create a new pipeline.
  • Ensure that the Amazon Aurora or Amazon RDS cluster has authentication enabled using AWS Secrets Manager, which is the only supported authentication mechanism.

Best practices

The following are some best practices to follow while setting up the integration with OpenSearch Service:

  • If a mapping template is not specified in OpenSearch, it automatically assigns field types using dynamic mapping based on the first document received. However, it is always recommended to define field types explicitly by creating a mapping template that suits your requirements.
  • To maintain data consistency, the primary and foreign keys of tables remain unchanged.
  • You can configure the dead-letter queues (DLQ) in your OpenSearch Ingestion pipeline. If you’ve configured the queue, OpenSearch Service sends all failed documents that can’t be ingested due to dynamic mapping failures to the queue.
  • Monitor recommended CloudWatch metrics to measure the performance of your ingestion pipeline.

Zero-ETL integration with Amazon DocumentDB

Amazon Document DB is a fully managed database service built for JSON data management at scale. It offers built-in text and vector search functionalities. By leveraging OpenSearch Service, you can execute search analytics, including features like fuzzy matching, synonym detection, cross-collection queries, and multilingual search capabilities on DocumentDB data.

The zero-ETL integration initiates the process with a full historical data extraction to OpenSearch using an ingestion pipeline. After the initial data load is completed, the pipelines read from Amazon DocumentDB change streams ensuring near real-time data consistency between the two systems. OpenSearch organizes the incoming data into indexes, with flexibility to either consolidate data from a DocumentDB collection into a single index or partition data across multiple indices. The ingestion pipelines synchronize all create, update, and delete operations from the DocumentDB collection, maintaining corresponding document modifications in OpenSearch. This ensures both data systems remain synchronised.

The pipelines offer configurable routing options, allowing data from a single collection to be written to one index or conditionally route to multiple indexes. Users can configure ingestion pipelines to stream data from Amazon DocumentDB to OpenSearch Service through three primary modes namely full load only, streaming change events without initial full load and full load followed by change streams. You can also monitor the state of ingestion pipelines in the OpenSearch service console. Additionally, you can use Amazon Cloudwatch to provide real-time metrics and logs and setting up alerts.

Pricing

There is no additional charge for using this feature apart from the cost of your existing underlying resources, including OpenSearch Service, OpenSearch Ingestion pipelines (OCUs), and Amazon DocumentDB. The integration performs an initial full load of Amazon DocumentDB data and continuously streams ongoing changes to OpenSearch Service using change streams. The change streams feature is disabled by default and does not incur any additional charges until the feature is enabled. Using change streams on a DocumentDB cluster incurs additional read and write input/output (I/O), as well as storage costs.

To learn more on pricing see the DocumentDB pricing page.

Considerations

The following are the limitations for the DocumentDB to OpenSearch Service integration:

  • Only one Amazon DocumentDB collection as the source per pipeline is supported.
  • Cross-region and cross-account data ingestion is not supported.
  • Amazon DocumentDB elastic clusters are not supported, only instance-based clusters are supported.
  • AWS Secrets Manager is the only supported authentication mechanism.
  • You can’t update an existing pipeline configuration to ingest data from a different database and/or a different collection. To update the database and/or collection name of a pipeline, create a new pipeline.

Best practices

The following are some best practices to follow while setting up the DocumentDB zero-ETL with OpenSearch Service:

  • Configure dead-letter queues (DLQ) to handle any failed document ingestion.
  • Configure AWS Secrets Manager and enable secrets rotation to provide the pipeline secure access.
  • If you’re using change streams in DocumentDB, it’s important to extend the retention period to up to 7 days. This ensures you don’t lose any data changes during the ingestion process.

To get started, see zero-ETL integration of Amazon DocumentDB with OpenSearch Service.

Benefits for Database Integrations

With zero-ETL integrations, you can use the powerful search and analytics features of OpenSearch Service directly on your latest database data. These include full-text search, fuzzy search, auto-complete, and vector search for machine learning (ML) workloads—enabling intelligent, real-time experiences that enhance your applications and improve user satisfaction. This integration uses change streams to automate the synchronisation of transactional data from Amazon Aurora, Amazon RDS, Amazon DynamoDB and Amazon DocumentDB to OpenSearch Service without manual intervention. Once the data is available in OpenSearch Service, you can perform real-time searches to quickly retrieve relevant results for your applications.This eliminates the need for manual Extract-Transform-Load (ETL) processes, reduces operational complexity, and accelerates time-to-insight for real-time dashboards, search, and analytics.

Conclusion

In this post, you learned that zero-ETL integrations represent a significant advancement in simplifying data analytics workflows and reducing operational complexity. As you’ve explored throughout this post, these integrations offer several advantages such as elimination of complex ETL pipelines and reduced infrastructure and operational costs by removing the need for intermediate storage and processing that enhance developer productivity.

It is time to accelerate your analytics journey with OpenSearch Service zero ETL – where your data flows seamlessly, eliminating complex pipelines and delivering real-time insights. Get started with Amazon OpenSearch Service or learn more about integrations with other services and applications in the AWS documentation.


About the authors

Omama Khurshid

Omama Khurshid

Omama is GTM Specialist Solutions Architect Analytics at Amazon Web Services. She focuses on helping customers across various industries build reliable, scalable, and efficient solutions. Outside of work, she enjoys spending time with her family, listening to music, and learning new technologies.

Canberk Keles

Canberk Keles

Canberk is an Associate Solutions Architect at Amazon Web Services, helping software companies achieve their business goals by leveraging AWS technologies. He is part of OpenSearch specialist community within AWS and has been guiding customers harness the power of OpenSearch. Outside of work, he enjoys sports, reading, traveling and playing video games.

Building a modern lakehouse architecture: Yggdrasil Gaming’s journey from BigQuery to AWS

Post Syndicated from Edijs Drezovs, Viesturs Kols, Krisjanis Beitans original https://aws.amazon.com/blogs/big-data/building-a-modern-lakehouse-architecture-yggdrasil-gamings-journey-from-bigquery-to-aws/

This is a guest post by Edijs Drezovs, CEO and Founder of GOStack, Viesturs Kols, Data Architect at GOStack, and Krisjanis Beitans, Senior Data Engineer at GOStack, in partnership with AWS.

Yggdrasil Gaming develops and publishes casino games globally, processing massive amounts of real-time gaming data for game performance analytics, player behavior insights, and industry intelligence. As Yggdrasil’s system grew, managing dual-cloud environments created operational overhead and limited their ability to implement advanced analytics initiatives. This challenge became critical ahead of the launch of the Game in a Box solution on AWS Marketplace, which generates increases in data volume and complexity.

Yggdrasil Gaming reduced multi-cloud complexity and built a scalable analytics foundation by migrating from Google BigQuery to AWS analytics services. In this post, you’ll discover how Yggdrasil Gaming transformed their data architecture to meet growing business demands. You will learn practical strategies for migrating from proprietary systems to open table formats such as Apache Iceberg while maintaining business continuity.

Yggdrasil worked with GOStack, an AWS Partner, to migrate to an Apache Iceberg-based lakehouse architecture. The migration helped reduce operational complexity and enabled real-time gaming analytics and machine learning.

Challenges

Yggdrasil faced several critical challenges that prompted their migration to AWS:

  • Multi-cloud operational complexity: Managing infrastructure across AWS and Google Cloud created significant operational overhead, reducing agility and increasing maintenance costs. The data team had to maintain expertise in both environments and coordinate data movement between clouds.
  • Architecture limitations: The existing setup couldn’t effectively support advanced analytics and AI initiatives. More critically, the launch of Yggdrasil’s Game in a Box solution required a modernized, scalable data environment capable of handling increased data volumes and enabling advanced analytics.
  • Scalability constraints: The architecture lacked the unified data foundation with open standards and automation required to scale efficiently. As data volumes grew, costs increased proportionally, and the team needed an environment designed for modern analytics at scale.

Solution overview

Yggdrasil worked with GOStack, an AWS APN partner, to design their new lakehouse architecture. The following diagram shows the high level overview of this architecture.

Figure 1: High-level architecture diagram of Yggdrasil's modern lakehouse on AWS

Figure 1: High-level architecture diagram

Yggdrasil successfully migrated from Google BigQuery to a data lakehouse architecture using Amazon Athena, Amazon EMR, Amazon Simple Storage Service (Amazon S3), AWS Glue Data Catalog, AWS Lake Formation, Amazon Elastic Kubernetes Service (Amazon EKS) and AWS Lambda. Their strategic approach aims to reduce multi-cloud complexity while building a scalable foundation for their Game in a Box solution and specific AI/ML initiatives like personalized game recommendations and fraud detection.

The combination of Amazon S3, Apache Iceberg, and Amazon Athena allowed Yggdrasil to move away from provisioned, always-on compute models. The Amazon Athena pay-per-query pricing charges only for data scanned, removing idle compute costs during off-peak periods. Internal cost modeling performed during the evaluation phase indicated that this architecture could reduce analytics system costs by 30–50% compared to compute-based warehouse pricing models of other solutions, particularly for bursty workloads driven by game launches, tournaments, and seasonal traffic. By adopting AWS-native analytics services, Yggdrasil reduced operational complexity through native integration with AWS Identity and Access Management (AWS IAM), Amazon EKS, and AWS Lambda, helping simplify security, governance, and automation across the analytics system.

The solution centers on a modern lakehouse architecture built on Amazon S3, which provides durable and cost-efficient storage for Iceberg tables in Apache Parquet format. Apache Iceberg table format provides ACID transactions, schema evolution, and time travel capabilities while maintaining an open standard. AWS Glue Data Catalog serves as the central technical metadata repository, while Amazon Athena acts as the serverless query engine used by dbt-athena and for ad-hoc data exploration. Amazon EMR runs Yggdrasil’s legacy Apache Spark application in a fully managed environment, and AWS Lake Formation provides centralized security and governance for data lakes, allowing fine-grained access control at database, table, column, and row levels.

The migration followed a phased approach:

  1. Establish lakehouse foundation – Set up Apache Iceberg-based architecture with Amazon S3 with AWS Glue Data Catalog
  2. Implement real-time data ingestion – Deploy Debezium connectors for real-time change data capture from EKS and Google Kubernetes Engine (GKE) clusters
  3. Migrate processing pipelines – Re-system ETL pipelines using AWS Lambda, and legacy data applications re-systemed on Amazon EMR
  4. Modernizing the transformation layer – Implement dbt with Amazon Athena for modular, reusable models
  5. Enable governance – Configure AWS Lake Formation for comprehensive data governance

Establish lakehouse foundation

The first phase of the migration focused on building a solid foundation for the new data lakehouse architecture on AWS. The goal was to create a scalable, secure, and cost-efficient environment that could support analytical workloads with open data formats and serverless query capabilities.

GOStack provisioned an Amazon S3-based data lake as the central storage layer, providing virtually unlimited scalability and fine-grained cost control. This storage-compute separation enables teams to decouple ingestion, transformation, and analytics processes, with each component scaling independently using the most appropriate compute engine.

To establish dataset interoperability and discoverability, the team adopted AWS Glue Data Catalog as the unified metadata repository. The catalog stores Iceberg table definitions and makes schemas accessible across services such as Amazon Athena and Apache Spark workloads on Amazon EMR. Most datasets, both batch and streaming, are registered here, enabling consistent metadata visibility across the lakehouse.

The data is stored in Apache Iceberg tables on Amazon S3, selected for its open table format, ACID transaction support, and powerful schema evolution features. Yggdrasil required ACID transactions for consistent financial reporting and fraud detection, schema evolution to accommodate rapidly changing gaming data models, and time travel queries to align with regulatory audit requirements.

GOStack built a custom schema conversion and table registration service. This internal tool converts source-system Avro schemas into Iceberg table definitions and manages the creation and evolution of raw-layer tables. By controlling schema translation and table registration directly, the team makes sure that metadata stays consistent with the source systems and provides predictable, versioned schema evolution aligned with ingestion needs.

The initial setup made the following components:

  • Amazon S3 bucket structure design: Implemented a multi-layer layout (raw, curated, and analytics zones) aligned with data lifecycle best practices.
  • AWS Glue Data Catalog integration: Defined database and table schemas with partitioning strategies optimized for Athena performance.
  • Iceberg configuration: Enabled versioning and metadata retention policies to balance storage efficiency and query flexibility.
  • Security and compliance: Configured encryption at rest using AWS Key Management Service (AWS KMS), helped enforce access controls via AWS IAM and Lake Formation, and implemented Amazon S3 bucket policies following the principle of least privilege.

The redesign of the previous GCP setup helped deliver price-performance improvements. Yggdrasil reduced ingestion and processing costs by approximately 60% while also lowering operational overhead through a more direct, event-driven pipeline.

Implement real-time data ingestion

After establishing the lakehouse architecture, the next step focused on enabling real-time data ingestion from Yggdrasil’s operational databases into the raw data layer of the lakehouse. The objective was to capture and deliver transactional changes as they occur, making sure that downstream analytics and reporting reflect the most up-to-date information.

To achieve this, GOStack deployed Debezium Server Iceberg, an open-source project that integrates change data capture (CDC) directly with Apache Iceberg tables. It was deployed as Argo CD applications on Amazon EKS and used Argo’s GitOps-based model for reproducibility, scalability, and seamless rollouts.

This architecture provides an efficient ingestion pathway – streaming data changes directly from the source system’s outbox tables into the Apache Iceberg tables registered in the AWS Glue Data Catalog and physically stored on Amazon S3, bypassing the need for intermediate brokers or staging services. By writing data in the Iceberg table format, the ingestion layer maintained transactional guarantees and immediate query availability through Amazon Athena.

Figure 2: Streaming ingestion pipeline using Debezium in Amazon EKS

Because Yggdrasil’s source systems emitted outbox events containing Avro records, the team implemented a custom outbox-to-Avro transformation within Debezium. The outbox table stored two key components:

  • The Avro schema definition
  • The JSON-encoded payload of each record

The custom transformation module combined these elements into valid Avro records before persisting them into the target Iceberg tables. This approach preserved schema fidelity and verified compatibility with downstream processing tools.

To dynamically route incoming change events, the team leveraged Debezium’s event router configuration. Each record was routed to the appropriate Apache Iceberg table (backed by Amazon S3) based on topic and metadata rules, while table schemas and partitioning were governed on the AWS Glue side to maintain stability and alignment with the lakehouse’s data organization standards.

This setup helped deliver low-latency ingestion with end-to-end streaming from database outbox to S3-based Iceberg tables in near real time. The team managed operations end to end on Amazon EKS using Helm charts deployed via Argo CD in a GitOps model for fully declarative, version-controlled operations. ACID-compliant Iceberg writes verified that partially written data could not corrupt downstream analytics. The modular transformation logic allowed future expansion to new source systems or event formats without rearchitecting the ingestion pipeline.

This Debezium Server solution provides fast, real-time data ingestion. GOStack considers it an interim architecture. In the long term, the ingestion pipeline will evolve to use Amazon Managed Streaming for Apache Kafka (Amazon MSK) as the central event backbone. Debezium connectors will act as producers, publishing change events to Apache Kafka topics, while Apache Flink applications will consume, process, and write data into Iceberg tables.

This planned evolution toward a Kafka-based streaming architecture verifies Yggdrasil’s lakehouse remains not only scalable and cost-efficient today, but also future-ready – capable of supporting richer streaming analytics and broader data integration scenarios as the organization grows.

Migrate processing pipelines

Once real-time data ingestion was established, GOStack turned its focus to modernizing the data transformation layer. The goal was to simplify the transformation logic, reduce operational overhead, and unify the orchestration of analytical workloads within the new AWS-based lakehouse.

GOStack adopted a lift-and-shift approach for some of Yggdrasil’s data pipelines to support a fast and low-risk transition away from GCP. The lightweight Cloud Run functions that previously handled extraction tasks – pulling data from file shares, SharePoint, Google Sheets, and various third-party APIs – were re-implemented using AWS Lambda. These Lambda functions now integrate with the same external systems and write data directly into Iceberg tables.

For more complex processing, previous Apache Spark applications running on Dataproc were migrated to Amazon EMR with minimal code changes. This allowed it to preserve the existing transformation logic while benefiting from the managed scaling capabilities of EMR and improved cost control on AWS.

Over time, these processes will be gradually refactored and consolidated into containerized workflows on the EKS cluster, fully orchestrated by Argo Workflows. This phased migration allows Yggdrasil to move workloads to AWS quickly and decommission GCP resources sooner, while still leaving room for continuous improvement and modernization of the data system over time.

Finally, a lot of analytical transformations that previously lived as BigQuery stored procedures and scheduled queries, that were now rebuilt as modular dbt models executed with dbt-athena. This shift made transformation logic more transparent, maintainable, and version-controlled, improving both developer experience and long-term governance.

Modernizing the transformation layer

With the ingestion pipelines migrated to AWS, GOStack turned its focus to simplifying and modernizing Yggdrasil’s analytical transformations. Rather than replicating the previous stored-procedure–driven approach, the team rebuilt the transformation layer using dbt to help improve maintainability, lineage visibility, orchestration, and long-term governance.As part of this redesign, several data models were reshaped to fit the new lakehouse architecture. The most significant effort involved rewriting a critical Spark-based financial transformation into a set of SQL-driven dbt models. This shift not only aligned the logic with the lakehouse design but also removed the need for long-running Spark clusters, helping generate operational and cost savings.For the curated data layers, replacing the legacy warehouse, GOStack consolidated numerous scheduled queries and stored procedures into structured dbt models. This provides standardized, version-controlled transformations and clear lineage across the analytical stack.

Orchestration was simplified as well. Previously, coordination was split between Apache Airflow for Spark workloads and scheduled queries analytical transformations, creating operational friction and dependency risks. In the new architecture, Argo Workflows on Amazon EKS orchestrates dbt models centrally, consolidating the transformation logic within a single workflow engine. While most transformations still run on time-based schedules today, the system now supports event-driven execution through Argo Events, giving the opportunity to progressively adopt trigger-based workflows as the transformation layer evolves.

This unified orchestration framework can bring multiple benefits:

  • Consistency: One orchestration layer for data workflows across ingestion and transformation.
  • Automation: Event-driven dbt runs help remove manual scheduling and reduce operational overhead.
  • Scalability: Argo Workflows scales with the EKS cluster, handling concurrent dbt jobs seamlessly.
  • Observability: Centralized logging and workflow visualization help improve visibility into job dependencies and data freshness.

Through this transformation, Yggdrasil successfully unified its data lakes and warehouses into a modern lakehouse architecture, powered by open data formats, serverless query engines, and modular transformation logic. The move to dbt and Athena not only simplified operations but also helped pave the way for faster iteration, simpler governance, and greater developer productivity across the data environment.

Lakehouse performance optimizations

While performance tuning is an ongoing journey, as part of the transformation redesign, GOStack made few performance-oriented tweaks to make sure Athena queries can be fast and cost-efficient. The Apache Iceberg tables were stored in Parquet with ZSTD compression, providing strong read performance and reducing the amount of data scanned by Athena.

Partitioning strategies were also aligned to actual access patterns using Iceberg’s native partitioning. Raw data zones were partitioned by ingestion timestamp, enabling efficient incremental processing. Curated data used business-driven partition keys, such as player or game identifiers and date dimensions, to help optimize analytical queries. These designs made sure Athena could prune unneeded data and consistently scan only the relevant partitions.

Iceberg’s native partitioning features, including transforms such as bucketing and time slicing, replace traditional Hive partitioning patterns. Because Iceberg manages partitions internally in its metadata layer, not all Glue or Athena partition constructs apply. Relying on Iceberg’s native partitioning helps provide predictable pruning and consistent performance across the lakehouse without introducing legacy Hive behaviors.

To handle the high volume of small files produced by real-time ingestion, GOStack enabled AWS Glue Iceberg compaction. This automatically merges small Parquet files into larger segments, helping improve query performance and reduce metadata overhead without manual intervention.

Enable governance

The team adopted AWS Lake Formation as the primary governance layer for the curated zone of the lakehouse, leveraging Lake Formation hybrid access mode to manage fine-grained permissions alongside existing IAM-based access patterns. This hybrid mode provides an incremental and flexible pathway to adopt Lake Formation without forcing a full migration of legacy permissions or internal pipeline roles, making it an ideal fit for Yggdrasil’s phased modernization strategy.

Lake Formation offers centralized authorization, supporting database, table, column, and, critically for Yggdrasil, row-level permissions. These capabilities are essential because of the company’s multi-tenant operating model:

  • Game development partners require access to data and reports pertaining only to their own games, facilitating both security and compliance alignment with partner agreements.
  • iGaming operators integrating with Yggdrasil’s system must receive operational and financial insights exclusively for their own data, enforced automatically through reporting tools backed by curated Iceberg tables.

With Lake Formation hybrid access mode, tenant-specific row-level access policies are consistently enforced across Amazon Athena, AWS Glue, and Amazon EMR, without introducing breaking changes to existing IAM-based workloads. This allowed Yggdrasil to implement strong governance for external consumers while keeping internal operations stable and predictable.

Internally, Lake Formation is also used to grant the Analytics team and BI tools targeted access to curated datasets, straightforward but centrally managed to maintain consistency and reduce administrative overhead.

For ingestion and transformation workloads, the team continues to rely on IAM roles and policies. Services such as Debezium, dbt, and Argo Workflows require broad but controlled access to raw and intermediate storage layers, and IAM provides a straightforward, least-privilege mechanism for granting those permissions without involving Lake Formation in the internal pipeline path.

By adopting Lake Formation in hybrid access mode and combining it with IAM for internal services, Yggdrasil established a governance model that can balance strong security with operational flexibility – enabling the lakehouse to scale securely as the business grows.

Results and business impact

The new lakehouse, built on Amazon Athena, Amazon S3, and AWS Glue Data Catalog, now underpins advanced analytics and AI/ML use cases such as player behavior modeling, predictive game recommendations, and fraud detection.

The optimized lakehouse design allows Yggdrasil to rapidly onboard new analytics workloads and business use cases, helping deliver measurable outcomes:

  • Reduced operational complexity through consolidation on AWS analytics services
  • Cost optimization with a 60% reduction in data processing costs
  • Improved data freshness with 75% lower latency for analytics results (from 2 hours to 30 minutes)
  • Enhanced governance using the AWS Lake Formation fine-grained controls
  • Future-ready architecture leveraging open formats and serverless analytics

Conclusion

Yggdrasil Gaming’s migration journey illustrates how organizations can successfully transition from proprietary analytics systems to an open, flexible lakehouse architecture. By following a phased approach guided by AWS Well-Architected Framework principles, Yggdrasil maintained business continuity while establishing a modern foundation for their data needs.

Based on this experience, several lessons emerged to help guide your own move to an AWS-based lakehouse:

  1. Assess your current state: Identify pain points in your existing data architecture and establish clear objectives for modernization.
  2. Start small: Begin with a pilot project using AWS analytics services to validate the lakehouse approach for your specific use cases.
  3. Design for openness: Leverage open table formats like Apache Iceberg to maintain flexibility and avoid vendor lock-in.
  4. Implement gradually: Follow a phased migration strategy similar to Yggdrasil’s, prioritizing high-value workloads.
  5. Optimize continuously: Use performance tuning techniques for Amazon Athena to help maximize efficiency and minimize costs.

To learn more about building modern lakehouse architectures, refer to “The lakehouse architecture of Amazon SageMaker”.


About the authors

Edijs Drezovs

Edijs Drezovs

Edijs is the CEO and Founder of GOStack an AWS Partner specializing in modernizing cloud-native infrastructures, data systems and analytics architectures. He brings over 12 years of experience driving complex cloud transformations and data engineering initiatives.

Viesturs Kols

Viesturs Kols

Viesturs is a Data Architect at GOStack with deep expertise in lakehouse architectures and real-time analytics. He led the technical implementation of Yggdrasil Gaming’s migration to AWS analytics services and specializes in Apache Iceberg and streaming data systems.

Krisjanis Beitans

Krisjanis Beitans

Krisjanis is Senior Data Engineer at GOStack specializing in lakehouse architectures, Apache Iceberg, Amazon Athena, and dbt-based transformation frameworks. During Yggdrasil Gaming’s migration to AWS, he rebuilt the analytical layer, designing Iceberg table structures, optimizing Athena performance, and implementing the dbt-driven transformation pipeline.

Alvaro Guerrero

Alvaro Guerrero

Alvaro is an AWS Solutions Architect who helps customers build innovative cloud solutions – specialised in AWS analytics services.

Aleksandra Zgnilec

Aleksandra Zgnilec

Aleksandra is an Account Executive at AWS supporting Betting & Gaming customers in their cloud and business transformations.

Zahi Njeim

Zahi Njeim

Zahi is a Business Development Manager at AWS for Betting & Gaming, Media, Entertainment, Games and Sports.

Standardize Amazon Redshift operations using Templates

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

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

The challenge: Managing repetitive data operations at scale

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

customer transactions | product catalogs | inventory updates

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

Their data engineering team faces several pain points:

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

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

Introducing Redshift Templates

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

Template management best practices

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

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

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

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

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

Solution overview

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

Scenario 1: Standardizing client data ingestion

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

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

This template defines their standard approach:

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

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

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

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

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

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

Scenario 2: Handling client-specific variations with parameter overrides

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

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

This demonstrates the Redshift Templates parameter hierarchy:

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

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

Scenario 3: Simplified template maintenance

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

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

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

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

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

DROP TEMPLATE data_ingestion.standard_client_load;

Scenario 4: Environment-specific templates for development and production

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

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

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

Key benefits

The key benefits of using templates include:

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

Additional use cases across industries

Templates can be used across industries.

Financial services: Standardizing regulatory data loads

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

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

Healthcare: Loading patient data with strict standards

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

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

Retail: JSON data loading standardization

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

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

Conclusion

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

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

How Twilio secured their multi-engine query platform with AWS Lake Formation

Post Syndicated from Aakash Pradeep, Venkatram Bondugula original https://aws.amazon.com/blogs/big-data/how-twilio-secured-their-multi-engine-query-platform-with-aws-lake-formation/

This is a guest post by Aakash Pradeep, Principal Software Engineer, and Venkatram Bondugula, Software Engineer at Twilio, in partnership with AWS.

Twilio is a cloud communications platform that provides programmable APIs and tools for developers to easily integrate voice, messaging, email, video, and other communication features into their applications and customer engagement workflows.

In this blog series we discuss how we built a multi-engine query platform at Twilio. The first part introduces the use case that led us to build a new platform and why we selected Amazon Athena alongside our open-source Presto implementation. This second part discusses how Twilio’s query infrastructure platform integrates with AWS Lake Formation to provide fine-grained access control to all their data.

At Twilio, we faced critical challenges in managing our multi-engine query platform across a complex data mesh architecture spanning multiple AWS accounts and Lines of Business. We needed a unified permissions model that could work consistently across different query engines like OSS Presto and Amazon Athena, eliminating the fragmented authentication experiences in our infrastructure. The growing demand for secure cross-account data sharing required moving beyond manual, multi-step provisioning processes that depended heavily on human intervention. Additionally, Twilio’s compliance and data stewardship requirements demanded fine-grained access controls at row, column, and cell levels, necessitating a scalable and flexible approach to permission management. By adopting the AWS Glue Data Catalog as our managed metastore and AWS Lake Formation for governance, we implemented Tag-Based Access Control (LF-TBAC) to simplify access management, enabled data sharing through automated workflows, and established a centralized governance framework that provided uniform permissions management across all AWS services.

Transitioning to a managed metastore and governance solutions

We discussed in part 1, how we were looking to move to managed services to alleviate us of the burden of managing the underlying infrastructure of a query platform. Along with our decision to adopt Amazon Athena, we also began to evaluate the adoption of Amazon EMR Serverless for our Spark workloads, which made us aware of the fact that we needed to migrate to a managed solution for our Apache Hive metastore.

We selected the AWS Glue Data Catalog as our managed metastore repository to support our enterprise-wide data mesh architecture. For managing permissions to the Data Catalog assets, we chose AWS Lake Formation, a service that enables data governance and security at scale using familiar database-like permissions. Lake Formation provides a unified permissions model as well as support for enabling data mesh architecture that we were seeking.

Lake Formation’s support for row, column, and cell-level access controls provides the fine-grained access control (FGAC) capabilities required by our compliance and data stewardship policies. Additionally, Lake Formation’s tag-based access control (LF-TBAC) feature allows us to define FGAC permissions based on tags attached to the Data Catalog resources, enabling flexible and scalable permission management.

Integrating Odin with AWS Lake Formation

Odin, our Presto-based gateway, serves as a central hub for query processing, managing authentication, routing, and the complete workflow throughout a query’s lifecycle. As the primary interface, Odin enables users to connect through JDBC or APIs from various BI tools, SQL IDEs, and other applications.

Beyond its core routing capabilities, Odin utilizes local caches implemented using Google’s Guava caching library to optimize performance across the platform. Guava delivers efficient in-memory caching for Java applications by storing data locally within the application instance, resulting in significantly faster retrieval times. Odin employs multiple Guava caching layers across various modules to ensure optimal response times for frequently accessed data and metadata.

Building on this performance foundation, Odin implements authentication and authorization layers to ensure secure and controlled access to data across multiple query engines. These security components work together to verify user identities and enforce data access policies, providing a unified security framework that abstracts away the complexities of individual engine implementations while maintaining strict governance standards.

The authentication layer

Different query engines like OSS Presto and Amazon Athena each implement their own authentication mechanisms. To create a consistent user experience, Odin provides a unified authentication layer that shields users from these underlying differences. Currently, Odin’s pluggable authentication system supports LDAP integration, with plans to expand this capability to include Okta authentication using IAM Identity center in the future.

The authorization layer

For data consumers using AWS Analytics services such as AWS Glue, Amazon EMR, and Athena through an IAM federated role-based access, AWS Lake Formation provided critical authorization capabilities for data governance through their existing integrations. However, we needed to extend its capabilities to integrate with OSS Presto. Additionally, our users for the query infrastructure platform were not mapped to an IAM user so would need to build a custom authorization layer in Odin to verify permissions and integrate with Lake Formation. Our challenge was creating a consistent way to control data access across all our query engines.

When a user runs a query, Odin’s authorization layer checks three key pieces of information:

  • Table details: which database and table the query is accessing
  • User permissions: what data tags the user has access to
  • Resource tags: what security tags are attached to the requested table

We store user permissions in Amazon DynamoDB, which allows us to quickly look up what each user can access. By matching the user’s tags with the table’s Lake Formation tags, we can determine if the query should be allowed. To keep things fast, we cache this information temporarily, allowing us to expedite authorization for recent requests.

How the authorization works:

  1. Initial check: First, we see if this user recently ran a similar successful query (within the last 5 minutes).
  2. Gather information: We collect the table details, user permissions, and security tags—first checking our cache, then fetching from AWS Glue Data Catalog and Lake Formation if needed.
  3. Match permissions: We compare the user’s access tags stored in a DynamoDB table against the table’s security tags in Lake Formation.
  4. Make decision: If the user’s permissions match what’s required for their query action (like SELECT or INSERT), access is granted.

This approach allows us to make use of Lake Formation tag-based access control while keeping our authorization logic separate from the individual query engines. By using smart caching and efficient lookups, we can verify permissions in just milliseconds.

Building a data mesh

At Twilio, we have multiple line of business (LoBs) each managing their own data platform infrastructure. The individual platforms are spread across multiple AWS accounts, and primarily store data on Amazon S3 in variety of open table formats, such as Apache Hudi, Apache Iceberg, and Delta Lake. Each platform independently supports analytics and machine learning use cases, however, there was a growing need for secure sharing of data across LoBs. Additionally, we needed to enable self-service discovery and provisioning of access to the data with a centralized governance framework.

Data consumers bring their own AWS accounts and choice of tools, which include not only AWS services such as Amazon Athena, AWS Glue ETL jobs (Spark), and Amazon EMR, but also AWS partner solutions. To improve the process of access fulfillment, data auditability and lowering the operational overhead involved, we needed an automated framework in place that had minimal human intervention and oversight.

Implementing a data subscription workflow

Previously, consumers requiring access to specific data sets would need to go through multiple steps to secure access, which involved several dependencies and manual actions. To simplify this process and provide a self-service capability, we decided to build a custom integration solution between ServiceNow and AWS Lake Formation. At Twilio, ServiceNow is used extensively to automate workflows and build custom applications to connect disparate systems and improve operational efficiency.

We automated key parts of the data access process using Twilio’s standard tools: Git for version control, Terraform for infrastructure management, and custom scripts to execute the necessary AWS actions.

We automated three main use cases:

1. Sharing data between accounts

When one team needs to share data with another team or with our central governance account, the process starts with a Git pull request (PR). This triggers our custom Lake Formation automation tool, which:

  • Connects to the source AWS account with admin permissions
  • Sets up data sharing using the security tags (LF-Tags) specified in a YAML configuration file
  • Completes the share using AWS Resource Access Manager (RAM)
  • Creates resource links in the target account so the data appears in their catalog
  • Updates ServiceNow with the newly shared database and table information

2. Granting permissions to user roles

When users request access to data, our automation tool grants tag-based permissions directly to their IAM roles in Lake Formation. This happens after approval of either a Git PR or ServiceNow ticket.

3. Granting access to individual users

For individual user access requests:

  • Users submit a request in ServiceNow for specific tables
  • After approval, ServiceNow calls our internal API that checks relevant Lake Formation tags
  • The request is validated and sent to an Amazon Simple Queue Service (Amazon SQS) queue
  • A consumer service processes the request, updates the user’s permissions in our DynamoDB table (which Odin uses for authorization checks), and includes retry logic for reliability
  • Once complete, the service updates the ServiceNow ticket to notify the user

The overall subscription and authorization flow is as shown in the diagram below:

Diagram of Twilio's AWS data query platform showing user access requests flowing through ServiceNow and LF-Tag validation before queries reach Amazon Athena via Odin EC2 instances.

  1. Users submit a request in ServiceNow for access to a database, table, or LF-Tag
  2. The system retrieves the relevant LF-Tags from Lake Formation through our API integration
  3. Upon approval, the automation procedure adds the user to the User-To-Tag DynamoDB table, grants IAM role permissions in Lake Formation, and sets up cross-account sharing via RAM as needed
  4. Users submit SQL query to the Odin presto gateway
  5. Odin authorizes the user through LDAP
  6. Odin parsers the SQL query to identify the tables involved and the action being performed (SELECT, DDL, and more)
  7. Odin validates permissions using the User to LF-Tag mapping and Lake formation grants to authorize the SQL query based on granted permissions
  8. If authorized, Odin routes the query to Amazon Athena or Presto

Using standardized tools and processes to provide self-service capabilities to the users helped us scale the governance framework and support broader use cases. Important capabilities in Lake Formation, such as Tag-based access control (TBAC) and cross-account sharing of data, simplified developing automations and our overall approach to governance.

Lessons learned- Cache is king

“By adopting AWS Glue Data Catalog as our managed metastore and AWS Lake Formation for Tag-Based Access Control, we simplified access management and enabled data sharing by reducing auth overhead to just 6-10 milliseconds through caching and targeted scaling.”

As Odin began handling queries at scale, we encountered performance bottlenecks in our customized authorization process as we had to retrieve information from multiple services, particularly with complex queries spanning multiple tables. The authorization checks involved in the performance bottleneck frequently caused query timeouts which impacted overall system reliability. The root of the problem lay in our sequential authorization workflow: our system first had to parse each query to identify all tables requiring identity verification, then make separate API calls to the AWS Glue Data Catalog and Lake Formation for each table’s permissions. It became clear that we needed to optimize this authentication process to reduce response times and improve the overall query experience.

We also recognized there were different caching needs between our POST operations and GET/DELETE HTTP calls, so we decided to separate them into two different Application Load Balancer (ALB) target groups. For POST requests, which required Lake Formation authentication, we found that concentrating traffic through just 2-3 target instances distributed across multiple Availability Zones (AZ) was more efficient. This approach allowed authentication information to be effectively cached locally on these dedicated instances, dramatically reducing the volume of API calls to the Lake Formation service.

GET and DELETE requests follow a more simplified workflow. Since users have already completed initial authorization, there is no need to continue to perform authorization checks. Although they follow a simpler workflow, these requests have much higher volume with requests numbering into the 10s of millions per hour. Due to this scale, we opted to implement horizontal scaling to scale the target ALB to 10 Amazon EC2 instances to fetch the query history from the DynamoDB table. These EC2 instances make use of local LRU caching with a 5-minute expiration policy for authentication data.

By implementing authentication caching and adopting specialized approaches for different HTTP request types with targeted scaling groups, we successfully reduced Odin’s overall overhead to a maximum of 6-10 milliseconds for both authentication and authorization.

Conclusion and what’s next

In this post, we explored how we enhanced Odin, our unified multi-engine query platform, with authentication and authorization capabilities using AWS Lake Formation and a custom authorization workflow. By using AWS services including Lake Formation, AWS Glue Data Catalog, and Amazon DynamoDB alongside Twilio’s existing infrastructure, we created a scalable self-service governance framework that streamlines user access management, simplifies auditing, and enables seamless data sharing across our complex cloud environment. With this workflow automation, we eliminated operational overhead while building a secure, robust platform that serves as the foundation for Twilio’s data mesh architecture.

Going forward, we are focusing on strengthening our authentication and authorization framework by enabling trusted federation with an identity provider(IdP) through AWS IAM Identity Center, which integrates directly with Lake Formation. Using Trusted Identity Propagation capabilities supported by IAM IDC will allow us to establish a consistent governance flow based on a user identity and will allow us to unlock the full capabilities of AWS Lake Formation such as fine-grained access control with data filters.

To learn more and get started with building with AWS Lake Formation, see Getting started with Lake Formation, and How to build a data mesh architecture at scale using AWS Lake Formation tag-based access control.


About the authors

Aakash Pradeep

Aakash Pradeep

Aakash is a Principal Software Engineer with over 15 years of experience across ingestion, compute, storage, and query platforms. Aakash is a PrestoCon speaker, holds multiple patents in real-time analytics, and is passionate about building high-performance distributed systems.

Venkatram Bondugula

Venkatram Bondugula

Venkatram is a seasoned backend engineer with over a decade of experience specializing in the design and development of scalable data platforms for big data and distributed systems. With a strong background in backend architecture and data engineering, he has built and optimized high-performance systems that power data-driven decision-making at scale.

Aneesh Chandra PN

Aneesh Chandra PN

Aneesh is a Principal Analytics Solutions Architect at AWS working with Strategic customers. He is passionate about using technology advancements to solve customers’ data challenges. He uses his strong expertise on analytics, distributed systems and open source frameworks to be a trusted technical advisor for AWS customers.

Amber Runnels

Amber Runnels

Amber is a Senior Analytics Specialist Solutions Architect at AWS specializing in big data and distributed systems. She helps customers optimize workloads in the AWS data ecosystem to achieve a scalable, performant, and cost-effective architecture. Aside from technology, she is passionate about exploring the many places and cultures this world has to offer, reading novels, and building terrariums.

Amazon OpenSearch Serverless introduces collection groups to optimize cost for multi-tenant workloads

Post Syndicated from Madhusudhan Narayana original https://aws.amazon.com/blogs/big-data/amazon-opensearch-serverless-introduces-collection-groups-to-optimize-cost-for-multi-tenant-workloads/

Today, we’re excited to announce the general availability of the collection groups feature for Amazon OpenSearch Serverless. With this feature you can reduce compute costs for multi-tenant workloads while creating secure tenant boundaries through per-tenant encryption, giving you the flexibility to balance cost efficiency with the exact level of isolation and security your applications requires.

Amazon OpenSearch Serverless is a serverless deployment option for Amazon OpenSearch Service, that eliminates the complexity of infrastructure management for running search and analytics workloads at scale. It automatically provisions and scales resources to deliver fast data ingestion rates and millisecond response times, even as usage patterns change. For organizations that are managing multi-tenant environments, data isolation, where the tenant’s data must be encrypted and protected (often with their own encryption keys), is a compliance requirement.

Previously, OpenSearch Serverless provided maximum security through physical isolation: each AWS Key Management Service key (KMS key) required dedicated OpenSearch Compute Units (OCUs) to maintain complete physical data separation. While this architecture provided the highest level of protection, it created challenges for multi-tenant deployments at scale. For customers managing multiple tenants with shared encryption keys, OCU resources are efficiently pooled, making the economics favorable. However, customers managing large numbers of smaller tenants, each requiring their own KMS key for data isolation, faced a challenge with higher cost. With dedicated OCU resources needed per unique key, the infrastructure costs could become prohibitive when individual tenants required only a fraction of an OCU’s capacity. This particularly impacted service providers wanting to offer bring your own key (BYOK) capabilities to their customers, forcing them to either absorb unsustainable costs or limit their service offerings.

OpenSearch Serverless has always provided flexible capacity management with maximum OCU settings to help you control costs. For most workloads, this model works seamlessly capacity scales up and down in response to demand, so you only pay for what you use. However, some workload patterns are simply better served by having a guaranteed baseline of compute ready to go from the start. Workloads with sudden traffic spikes, high-speed data ingestion pipelines, or load testing scenarios benefit from having capacity pre-allocated, so that the first requests are handled with the same responsiveness as any other. Similarly, multi-tenant architectures and time-sensitive operations often require predictable, consistent performance from the moment a collection becomes active.

Flexible controls with collection groups

Collection groups give you flexible control over security boundaries and resource allocation. Instead of forcing a one-size-fits-all approach, you can now tailor your architecture to match your specific security and cost requirements. Here’s how it works:

  1. Define your security boundary that matches your need: Collection groups is a logical security construct for related collections. Each collection groups maintains strong isolation with physically separated memory, CPU and disk from other collection groups, ensuring robust security boundaries between different security constructs.
  2. Share resources across encryption keys: Allocate collections to your collection groups regardless of whether they share KMS keys or use separate ones. Collections with different encryption keys can now share OCU resources within the same security boundary, dramatically reducing costs while maintaining full encryption protection and logical separation for each tenant.
  3. Deploy with flexible network access: Collection groups support collections with different network access types, allowing you to combine collections with public endpoints and VPC endpoints within the same group. This flexibility lets you match your security and connectivity requirements while benefiting from shared resource management across all collections in the group.
  4. Control cost and performance: Set maximum OCUs to cap spending and minimum OCUs to guarantee baseline performance. This dual control gives you a defined resource envelope for each collection groups, eliminating cost surprises while ensuring consistent performance.
  5. Optimize with insights: Access detailed CloudWatch metrics showing resource consumption, relative usage patterns, and latency across collection groups. These insights help you right-size allocations, identify optimization opportunities, and tune performance based on actual workload behavior.

With collection groups, you now have full control over resource allocation through both minimum and maximum OCU settings

Maximum OCUs: Cost control

Set an upper limit on resources to prevent runaway scaling and control costs per collection groups. This helps ensure you never exceed your budget, even during unexpected traffic spikes. Collection groups capacity limits operate independently from account-level limits. Account-level maximum OCU settings apply only to collections not associated with any collection groups, while collection groups maximum OCU settings apply to collections within that specific group. The sum of (Max OCUs across all your collection groups + Max OCU setting at the account level) should be less than your Service Quota Max OCUs allowed for your account. This separation gives you granular cost control across different security contexts.

Minimum OCUs: Performance guarantees

Define the baseline compute resources that will always be allocated to your collection groups, for consistent performance and resource availability. These OCUs are reserved exclusively for your collection groups and provide:

  • Instant availability with no cold starts: Your collections benefit from instant availability without scaling delays. Resources are always warm and ready, eliminating scaling delays when traffic arrives.
  • Guaranteed capacity: Resources are always available, even during periods of low activity or when competing with other collection groups, ensuring predictable performance even during low-traffic periods.
  • Predictable costs: Minimum OCUs are charged continuously, providing you with reserved capacity in exchange for predictable billing giving you cost certainty in exchange for guaranteed performance. This reserved baseline serves as the foundation for auto-scaling, which expands capacity up to your maximum limit as demand increases.

This combination gives you the flexibility to balance cost optimization with performance guarantees based on your specific requirements.

Multi-tenant cost economics with collection groups

Managing costs in multi-tenant architectures has always required balancing isolation, performance, and efficiency often at the expense of one another. Collection groups change that equation by enabling shared capacity across collections without sacrificing security boundaries. The following details how this plays out when you work with collection groups or without.

Before collection groups: Consider a customer with 10 tenants, each requiring their own KMS key for data isolation. Most of these tenants have modest data requirements typically 10-100GB, with the majority on the smaller end of that range. Managing dedicated resources for each tenant’s encryption key, regardless of their actual capacity needs, created operational complexity and cost challenges at scale.

With collection groups: The same customer can now group their tenants with similar security requirements into the collection groups, sharing OCU resources across collections. Tenants requiring only a small portion of OCU capacity no longer force the allocation of dedicated resources, reducing costs by up to 90% for large number of smaller tenant workloads.

With minimum OCU configuration: Premium tenants can be placed in collection groups with minimum OCUs set to guarantee performance, while standard tenants use collection groups with lower minimum thresholds for cost efficiency.

The following table illustrates how these cost savings play out across different tenant configurations, comparing infrastructure costs with and without collection groups across varying data sizes and query loads.

Number of tenants with unique KMS keys

Data size and query parameters

Cost with complete data isolation (without collection groups)

Cost with collection groups

Additional comments

10

Data size: 60GB or less

Query: Not needing more than base OCU (1 for redundant collection) compute

$3,500 $350 10x Savings in cost.
10

Data size: 60GB or less

Query: More than base OCU (1 for redundant collection) compute during peak times (For example – 5 additional OCUs per tenant without collection groups & 40 OCUs across all tenants based with collection groups due to benefit of shared infra).

$3500 + Peak time scale out per tenant ($8650) $350+ Peak time scale out ($6912). The system will scale up when there is additional query load, additional OCUs are deployed during this time. However when the load scales back, the system will scale-in to base OCU’s.
10 Data size: Sample data size in GB per tenant [3, 5, 7, 8, 10, 15, 18, 25, 28, 150]

Query: Can handle queries upto certain level with minimum OCU for the data size and then scales out on load.

For the sample data sizes, minimum OCU requirement will be [2, 2, 2, 2, 2, 2, 2, 2, 2, 8] = 26 OCUs [$4492] + Peak time scale out per tenant Minimum cost is determine by the number of OCUs required to hold the data across all tenants (120GB per OCU *2) + Peak time scale out.For the sample data sizes, 8 OCUs [$1382] + Peak time scale out per tenant The system will scale up when there is additional query load, additional OCUs are deployed during this time. However when the load scales back, the system will scale-in to minimum number of OCU required to hold the data.

Note: Above calculations are made with assumption for redundant enabled collections. For non-redundant mode it will be half the above calculations.

Getting started with collection groups

Collection groups and minimum OCU configuration are available in all AWS Regions where OpenSearch Serverless is offered, at no additional charge. Collection groups offers a new organizational feature to create collection groups and add new collections directly to these groups for enhanced management capabilities. While your existing collections will continue to operate unchanged and remain independent of any collection groups, you can immediately start using collection groups for new collections to benefit from improved organization and workflow management.

Currently, only newly created collections can be associated with collection groups, and all collections within a group must be of the same type (search, time series, or vector search). Existing collections continue to operate independently with their current capacity management settings, and you cannot mix different collection types within a single collection groups. You can use the AWS Management Console, AWS CLI, AWS CloudFormation, or AWS CDK to create the collection groups. In the following section we will show you how you can create the collection groups using the OpenSearch Service console.

To create your first collection groups:

  1. Open the OpenSearch Service console.
  2. In the left navigation pane, choose Serverless, then choose Collection groups.
  3. Choose Create collection groups.
  4. For collection groups name, enter a name for your collection groups. The name must be 3-32 characters long, start with a lowercase letter, and contain only lowercase letters, numbers, and hyphens.
  5. (Optional) For Description, enter a description for your collection groups.
  6. In the Capacity management section, configure the OCU limits:
    1. Maximum indexing capacity – The maximum number of indexing OCUs that collections in this group can scale up to.
    2. Maximum search capacity – The maximum number of search OCUs that collections in this group can scale up to.
    3. Minimum indexing capacity – The minimum number of indexing OCUs to maintain for consistent performance.
    4. Minimum search capacity – The minimum number of search OCUs to maintain for consistent performance.
  7. (Optional) In the Tags section, add tags to help organize and identify your collection groups.
  8. Choose Create collection groups.

To assign collection to the collection groups

  1. Open the Amazon OpenSearch Service console.
  2. In the left navigation pane, choose Serverless, then choose Collections.
  3. Choose Create collection.
  4. For Collection name, enter a name for your collection. The name must be 3-28 characters long, start with a lowercase letter, and contain only lowercase letters, numbers, and hyphens.
  5. (Optional) For Description, enter a description for your collection.
  6. In the Collection groups section, select the collection groups you want the collection to be assigned to. A collection can only belong to one collection groups at a time.
    (Optional) You can also choose to Create a new group. This will navigate you to the Create collection groups workflow. After you finish creating the collection groups, return to the step 1 of this procedure to begin creating your new collection.
  7. Continue through the workflow to create the collection.

Managing collection groups

Once you’ve created your collection groups, you can update their settings as your architecture evolves. The Amazon OpenSearch Serverless documentation provides step-by-step guidance on how to edit and delete collection groups, including updating OCU limits and modifying group configurations using the AWS Management Console, CLI, and CloudFormation.

Conclusion

OpenSearch Serverless collection groups transform how you can architect multi-tenant deployments by offering flexible deployment modes that balance security requirements with operational efficiency. You can now choose the collection groups where you define logical security boundaries that allow collections, regardless of whether they share the same KMS key or use different KMS keys to share OCU resources.

This flexibility directly addresses the cost challenges that previously made multi-tenant deployments prohibitive. By consolidating collections within collection groups, you can reduce infrastructure costs while maintaining robust encryption and tenant isolation. Configuring both minimum and maximum OCUs for each collection groups solves the cold-start and capacity guarantee challenges: minimum OCUs ensure your collections maintain ready compute resources to handle high-speed ingestion, sudden traffic spikes, and load testing without performance degradation. Maximum OCUs provide cost predictability and spending controls. This dual configuration gives you a defined resource envelope that eliminates both the uncertainty of cold starts and the risk of runaway costs.

To dive deeper into the collection groups and minimum OCU configuration, visit the Amazon OpenSearch Serverless documentation.

About the authors

Madhusudhan Narayana

Madhusudhan Narayana

Madhusudhan is Senior Software Engineer with Amazon Web Services. He is focused on OpenSearch Service and has years of experience in software engineering, distributed and autonomous systems. He holds a MS in Computer Science.

Prashant Agrawal

Prashant Agrawal

Prashant is a Sr. Search Specialist Solutions Architect with Amazon OpenSearch Service. When not working, you can find him traveling and exploring new places. In short, he likes doing Eat → Travel → Repeat.

Xian Huang

Xian Huang

Xian is a Product Marketing Manager at AWS.

Improving order history search using semantic search with Amazon OpenSearch Service

Post Syndicated from Shwetabh . original https://aws.amazon.com/blogs/big-data/improving-order-history-search-using-semantic-search-with-amazon-opensearch-service/

If you’ve ever shopped on Amazon, you’ve used Your Orders. This feature maintains your complete order history dating back to 1995, so you can track and manage every purchase you’ve made. The order history search feature lets you find your past purchases by entering keywords in the search bar. Beyond just finding items, it provides a straightforward way to repurchase the same or similar items, saving you time and effort.

Various features across Amazon’s shopping experience, such as Rufus and Alexa, use order history search to help you find your past purchases. Therefore, it’s important that order history search can locate your past purchased items as accurately and quickly as possible.

In this post, we show you how the Your Orders team improved order history search by introducing semantic search capabilities on top of our existing lexical search system, using Amazon OpenSearch Service and Amazon SageMaker.

Limitations of lexical search

Order history search uses lexical matching to find items from the entire order history of a customer that match at least one word of the search keywords. For example, if a customer searches for “orange juice,” the system retrieves all orange juice items as well as fresh oranges and other fruit juices the customer had previously ordered. Although lexical matching can provide a high recall of items with terms matching the search keywords precisely, it doesn’t work well for related or generic search keywords, like “health drinks” in this example.

Since the launch of Rufus, Amazon’s AI-enabled shopping assistant, a growing number of customers are experiencing a streamlined and richer shopping journey, including searching for their previous purchases with Rufus. Customers can now ask “Show me healthy drinks” without worrying about using lengthy, more precise terms like “kombucha”, “green tea”, and “protein shakes”. This makes the search experience more conversational and intent-based, presenting an opportunity to make item discovery more intuitive. For Rufus to answer order history searches with the same intuitive experience such as “Show me the healthy drinks I bought last year”, the underlying order history data store (“Your Orders”) needs semantic search capability to understand the underlying semantics of search keywords beyond the conventional lexical matching.

Challenges implementing semantic search

Implementing semantic search at our scale presented several technical challenges:

  • Scale – We needed to enable semantic search across billions of records corresponding to customers’ order history globally.
  • Zero downtime – We needed to keep the system 100% available while making changes on the backend to introduce semantic search.
  • Preventing search quality degradation – Semantic search is intended to improve the quality of search results. However, in some cases, it can reduce search quality. For example, if a customer remembers their item name exactly and wants to find only items matching that name, surfacing similar items in addition to the exactly matching items will increase crowding in results and make it harder to find the relevant item. Similarly, semantic search will not work for cases where the customer intends to search by identifier values, like order ID, which lack an inherent semantic meaning. For these scenarios, we use lexical search only.

Solution overview

Semantic search is powered by large language models (LLMs), which are mostly trained on human languages. These models can be adapted to take a piece of text in any language they were trained in and emit an embedding vector of a fixed length, irrespective of the input text length. By design, embedding vectors capture the semantic meaning of input text such that two semantically similar text strings have high cosine similarity computed on their respective embedding vectors. For semantic search on order history, the input text subject to embedding generation and similarity computation are the customer search phrases and the product text of purchased items.

We divide our solution into two parts:

  • Improving system scalability and resiliency for handling requests at scale – Before implementing semantic search, we needed to ensure our infrastructure could handle the increased computational load, leading us to adopt a cell-based architecture. This step is not needed for every use case, but systems with very high scale in terms of request or data volume can benefit a lot from its use before implementing a resource-intensive use case like semantic search.
  • Implementing semantic search – We began by evaluating the available embedding models, using the offline evaluation capabilities of Amazon Bedrock to test different models. After we selected our model, we could establish the infrastructure for generating embedding vectors.

Improving system scalability and resiliency

We used the cell-based architecture design pattern for improving our scalability and resiliency. A cell-based design entails partitioning the system into identical, smaller, self-contained chunks, or cells, which handle only a part of the overall traffic received by the system. The following diagram shows a high-level representation of a cell-based design for order history search.

Cell-based architecture diagram showing customer request routing to Amazon OpenSearch Service domains via hash-based partitioning

Each cell serves a defined subset of our customers. Cells don’t need to communicate with one another to serve a customer request. Each customer is assigned to a cell and each request from that customer is routed to that cell. The OpenSearch Service domain in each cell holds data only for the subset customers that it is supposed to serve. The number of cells (N) and distribution of data among those cells depends on the business use case, but the goal is to achieve as even a distribution of data and traffic as possible.

The routing logic can be kept as simple or as sophisticated as the use case requires it to be. The cell assignment values can either be computed at runtime for each request, or they can be computed one time and written to a cache or persistent data store like Amazon DynamoDB, from where cell assignment values can be fetched for subsequent requests. For order history search, the logic was simple and quick enough to be executed at runtime for each request. Looking up cell assignment from a persistent data store is especially useful for cases where there is a risk of some cells becoming “heavier” than others over time. In such cases, it becomes easier to redistribute the heavy cell’s data by simply overriding cell assignment values for specific keys in the data store, instead of having to change the partitioning logic immediately, which might have an impact on data distribution across all the cells.

As the system’s load grows, the number of cells in the system can be increased to handle the additional traffic. Even without increasing the number of cells in the system, we can redistribute current data among the existing N cells by reassigning some keys from one or more heavily populated cells to different lightly populated cells to spread out the load more evenly across all the cells and make more efficient use of the infrastructure.

A cell-based architecture also helps make the system more resilient. For example, if we lose one cell, our capacity is diminished only by 1/N, instead of 100%. This arrangement can also be improved to reduce the capacity loss even further by assigning partitioning keys to two or more cells such that they get written to two or more cells. In such cases, loss of a single cell does not result in data loss.

Implementing semantic search

Implementing semantic search for our order history search required several key decisions and technical steps. We began by evaluating the available embedding models, using the offline evaluation capabilities of Amazon Bedrock to test different models against our specific business domain requirements. This evaluation process helped us identify which model would deliver the best performance for our use case. After we selected our model, we needed to establish the infrastructure for generating embedding vectors. We containerized our embedding model and registered it in Amazon Elastic Container Registry (Amazon ECR), then deployed it using SageMaker inference endpoints to handle the actual vector computation at scale.

For the search infrastructure itself, we chose OpenSearch Service to implement our semantic search capabilities. OpenSearch Service provided both the vector storage we needed and the search algorithms required to deliver relevant results to our users.

One of our biggest challenges was updating our historical data to support semantic search on existing orders. We built a data processing pipeline using AWS Step Functions to orchestrate the workflow and AWS Lambda functions to handle the actual vector generation for our legacy data, so we could provide semantic search for all the records we wanted to.

The following diagram illustrates the high-level architecture.

Architecture diagram showing read-flow and write-flow for semantic search using Amazon OpenSearch Service and Amazon SageMaker embedding vectors

Model evaluation and selection

Order history search uses an embedding model trained on Amazon-specific data. Domain-specific training is critical because the generated embedding vectors must work well for the business context to return quality results.

We used an LLM-as-a-judge methodology with Anthropic’s Claude on Amazon Bedrock to evaluate candidate models. Anthropic’s Claude received prompts containing anonymized item text and search phrases from customer order history, then filtered and ranked items by relevance. These results served as ground truth for comparison.

We evaluated models using standard ranking metrics:

  • Normalized Discounted Cumulative Gain (NDCG) – Measures ranking quality against ideal order
  • Mean Reciprocal Rank (MRR) – Considers position of first relevant item
  • Precision – Rates accuracy of retrieved results
  • Recall – Rates ability to retrieve all relevant items

This process helped us determine the best model.

Retrieval strategy: Customer-scoped comprehensive search

Order history search has two key requirements:

  • Search only through the requesting customer’s order history – We don’t want items from one customer’s order history showing up in search results for another customer
  • Search all of that customer’s history – We don’t want to miss showing an item that would have been relevant for the customer’s search phrase just because the search algorithm missed evaluating it for some reason

Our approach involves using OpenSearch Service to retrieve all items for the customer who issued the search query, calculating relevance scores for each of them against the search phrase, sorting by score, and returning top K results. This provides comprehensive results coverage for each customer.

Vector storage with OpenSearch Service

We used two OpenSearch Service features for efficient vector storage and search:

  • knn_vector datatype – Built-in support for storing embedding vectors. Existing domains can add this field type without reindexing, enabling exact kNN search across all records. We didn’t need approximate kNN because the number of records for most customers was small enough for exact kNN to scale.
  • Scripted scoring – Painless scripts compute vector similarity server-side, reducing client complexity and maintaining low latency.

Hybrid search

Hybrid search refers to combining the results of lexical and semantic search to benefit from the strengths of each. The hybrid query capabilities of OpenSearch Service simplify implementing hybrid search by letting clients specify both types of queries in a single request. OpenSearch Service runs both queries in parallel, merges their results, normalizes the relevance scores of the sub-queries, and sorts results by the provided sort order (relevance score by default) before returning them to clients.

This gives clients the best of both types of searches. For example, there are certain scenarios where the search phrase doesn’t make much sense semantically, like when customers search by their orderId values. Semantic search is not designed for such cases; these are best served using keyword matching.

The hybrid search functionality helped save implementation effort and potential latency increase for order history search.

Updating historical data

After the infrastructure has been set up, newly ingested records are persisted with the relevant embedding vectors and support semantic search on those records. However, when customers search, they typically search for products they had purchased earlier. Therefore, the system might not help improve customer experience much unless the older records are updated to include the relevant embeddings. The approach to populate this data depends on the scale of the problem at hand.

Releasing the change to minimize potential customer impact

Our final step was to release the change to clients in a manner such that the impact of any potential problems is as small as possible. There are multiple ways to do that, including:

  • Implementing semantic search in a manner such that any transient issues in the semantic search flow make the logic fall back to lexical-only search, instead of failing the request completely. Even if semantic search doesn’t execute, the system should still be able to return results of lexical search to the client, instead of empty results.
  • Gating the change such that the default behavior remains lexical-only search and clients who need the semantic search feature must pass an additional flag in the request, for example, which executes the semantic or hybrid flow only for those requests.
  • Keeping the new flow behind a feature flag during the initial period such that it could be turned off completely if some critical problem is detected.

Examples of improved customer experience

The following are some examples of customer interactions with Rufus that required Rufus to query the respective customer’s order history to answer their question and give them the required pieces of information.

The following screenshots show how semantic search picks up wooden spoons for a “sustainable utensils” query and different kinds of chargers despite not having the keyword “charger” in the title description, in the case of the wall connector.

Two side-by-side screenshots demonstrating semantic search results for sustainable utensils and chargers in an e-commerce interface.

The following screenshots show how semantic search picks up relevant results even though the title description doesn’t include the queried keywords.

Two side-by-side screenshots demonstrating semantic search results for healthy snacks and kids educational items in an e-commerce order interface.

The semantic search feature of order history search helped Rufus fetch them and show to the customers. Before semantic search, Rufus wasn’t able to show any results to customers for such queries.

Business impact

Our solution resulted in the following key business impacts:

  • Customer experience improvements – The solution achieved 10% improvement in query recall, increasing the percentage of searches that return relevant results. It also reduced customer service contacts for issues related to locating past orders.
  • Partner integration success – The solution strengthened natural language processing capabilities for Alexa and Rufus, enhancing their ability to interpret order history queries. It also reduced the need for reranking and postprocessing by partner teams. We improved query success rate by 20%, meaning more customer searches now return at least one relevant item. We also observed enhanced result coverage by 48%, with semantic search consistently surfacing additional relevant matches that lexical search would have missed.

Conclusion

In this post, we showed you how we evolved Amazon order history search to support semantic search capabilities. This transition involved using cutting-edge AI technology while working within existing infrastructure limitations to develop solutions that avoided disruption and maintained SLAs during the feature upgrade. The implementation also involved backfilling, where billions of documents were processed at rates multiple times higher than normal ingestion to compute embedding vectors for previously purchased items. This operation required careful engineering and took advantage of the resilience OpenSearch Service offers even under extreme load.

Beyond the immediate implementation, this foundation enables continued innovation in search technology. The embedding vectors framework can incorporate improved models as they become available, and the architecture supports expansion into new capabilities such as personalization and multi-modal search.

You can get started with exact k-NN search today following the instructions in Exact k-NN search. If you’re looking for a managed solution for your OpenSearch cluster, check out Amazon OpenSearch Service.


About the authors

Shwetabh

Shwetabh

Shwetabh is a Senior Software Engineer at Amazon with interests in distributed systems and machine learning. Outside of work, he’s an avid reader with a particular love for technical deep-dives and thought-provoking non-fiction.

Harshavardhan Miryala

Harshavardhan Miryala

Harshavardhan is a Software Engineer at Amazon. He is passionate about machine learning, with particular interest in information retrieval and distributed computing. Outside of work, he enjoys playing racquet sports and watching football.

Ayush Kumar

Ayush Kumar

Ayush is a Tech Leader at Amazon. He is a passionate builder with an experience of over 14 years and leads the Your Orders Search product. In his spare time, he enjoys watching cricket and playing with his toddler.

Verisk cuts processing time and storage costs with Amazon Redshift and lakehouse

Post Syndicated from Karthick Shanmugam, Srinivasa Are original https://aws.amazon.com/blogs/big-data/verisk-cuts-processing-time-and-storage-costs-with-amazon-redshift-and-lakehouse/

This post is co-written with Srinivasa Are, Principal Cloud Architect, and Karthick Shanmugam, Head of Architecture Verisk EES (Extreme Event Solutions).

Verisk, a catastrophe modeling SaaS provider serving insurance and reinsurance companies worldwide, cut processing time from hours to minutes-level aggregations while reducing storage costs by implementing a lakehouse architecture with Amazon Redshift and Apache Iceberg. If you’re managing billions of catastrophe modeling records across hurricanes, earthquakes, and wildfires, this approach eliminates the traditional compute-versus-cost trade-off by separating storage from processing power.

In this post, we examine Verisk’s lakehouse implementation, focusing on four architectural decisions that delivered measurable improvements:

  • Execution performance: Sub-hour aggregations across billions of records replaced long batch process
  • Storage efficiency: Columnar Parquet compression reduced costs without sacrificing response time
  • Multi-tenant security: Schema-level isolation enforced complete data separation between insurance clients
  • Schema flexibility: Apache Iceberg support column additions and historical data access without downtime

The architecture separates compute (Amazon Redshift) from storage (Amazon S3), demonstrating how to scale from billions to trillions of records without proportional cost increases.

Current state and challenges

In Verisk’s world of risk analytics, data volumes grow at exponential rates. Every day, risk modeling systems generate billions of rows of structured and semi-structured data. Each record captures a micro-slice of exposure, event probability, or loss correlation. To convert this raw information into actionable insights at scale, experts need a data engine designed for high-volume analytical workloads.

Each Verisk model run produces detailed, high-granularity outputs that include billions of simulated risk factors and event-level results, multi-year loss projections across thousands of perils, and deep relational joins across exposure, policy, and claims datasets.

Running meaningful aggregations (such as, loss by region, peril, or occupancy type) over such high volumes created performance challenges.

Verisk needed to build a SQL service that could aggregate at scale in the fastest time possible and integrate into their broader AWS solutions, requiring a serverless, open, and performant SQL engine capable of handling billions of records efficiently.

Prior to this cloud-based release, Verisk’s risk analytics infrastructure operated on an on-premises architecture centered around relational database clusters. Processing nodes shared access to centralized storage volumes through dedicated interconnect networks. This architecture required capital investment in server hardware, storage arrays, and networking equipment. The deployment model required manual capacity planning and provisioning cycles, limiting the organization’s ability to respond to fluctuating workload demands. Database operations depended on batch-oriented processing windows, with analytical queries competing for shared compute resources.

Amazon Redshift and lakehouse architecture

Lakehouse architecture on AWS combines data lake storage scalability with data warehouse analytical performance in a unified architecture. This architecture stores vast amounts of structured and semi-structured data in cost-effective Amazon S3 storage while maintaining Amazon Redshift’s massively parallel SQL analytics.

Amazon Redshift is a fully managed, petabyte-scale cloud data warehouse service that delivers fast query performance using massively parallel processing (MPP) and columnar storage. Amazon Redshift eliminates the complexity of provisioning hardware, installing software, and managing infrastructure, keeping focus on deriving insights from their data rather than maintaining systems.

To meet their challenge, Verisk designed a hybrid data lakehouse architecture that combines the storage scalability of Amazon S3 with the compute power of Amazon Redshift. The following diagram shows the foundational compute and storage architecture that powers Verisk’s analytical solution.

Compute and Storage Layer

Architecture Overview

The architecture processes risk and loss data through three distinct stages within the lakehouse architecture, with comprehensive multi-tenant delivery capabilities to maintain isolation between insurance clients.

Amazon Redshift allows retrieving data directly from S3 using standard SQL for background processing. This solution collects detailed result outputs, join them with internal reference data, and executes aggregations over billions of rows. Concurrency scaling guarantees that hundreds of background analyses using multiple serverless clusters can run simultaneous aggregation queries.

The following diagram shows the architecture designed by Verisk

Architecture Design used by Verisk

Data ingestion and storage foundation

Verisk stores risk model outputs, location level losses, exposure tables, and model data in columnar Parquet format within Amazon S3. An AWS Glue crawler extracts metadata from S3 and feeds it into the lakehouse processing pipeline.

For versioned datasets like exposure tables, Verisk adopted Apache Iceberg, an open table format that addresses schema evolution and historical versioning requirements. Apache Iceberg provides transactional consistency through atomicity, consistency, isolation, durability ACID-compliant operations that maintain consistent snapshots during concurrent updates. Snapshot-based time travel allows data retrieval at previous points in time for regulatory compliance, audit trails, and model comparison with rollback capabilities. Schema evolution supports adding, dropping, or renaming columns without downtime or dataset rewrites. Incremental processing uses metadata tracking to process only changed data, reducing refresh times. Hidden partitioning and file-level statistics reduce I/O operations, improving aggregation performance. Engine interoperability allows accessing the same tables across Amazon Redshift, Amazon Athena, Spark, and other engines without data duplication.

Verisk built a foundation that combines S3’s cost-effectiveness with data management by adopting Apache Iceberg as open table format for this solution.

Three-stage processing pipeline

This pipeline orchestrates data flow from raw inputs to analytical outputs through three sequential stages. Pre-processing prepares and cleanses data, modeling applies risk calculations and analytics, and post-processing aggregates results for delivery.

  • Stage 1: Pre-processing transforms raw data into structured formats using Iceberg Tables and Parquet files, then processes it through Amazon Redshift Serverless for initial data cleaning and transformation.
  • Stage 2: Modeling takes place with a process built on AWS Batch the pre-processed data and applies advanced analytics and feature engineering. Results are stored in Iceberg Tables and Parquet files.
  • Stage 3: Aggregated Results are obtained during post-processing using Amazon Redshift Serverless, it produces the final analytical outputs in Parquet files, ready for consumption by end users.

Multi-tenant delivery system

The architecture delivers results to multiple insurance clients (tenants) through a secure, isolated delivery system that includes:

  • Amazon Quick Sight dashboards for visualization and business intelligence
  • Amazon Redshift as the data warehouse for querying aggregated results
  • AWS Batch for modelling processing.
  • AWS Secrets Manager to manage tenant-specific credentials
  • Tenant Roles implementing role-based access control to provide data isolation between clients

Summarized results are exposed through Amazon Quick Sight dashboards or downstream APIs to underwriting teams.

Multi-tenant security architecture

A critical requirement for Verisk’s SaaS solution was supporting comprehensive data and compute isolation between different insurance and reinsurance clients. Verisk implemented a comprehensive multi-tenant security model that provides isolation while maintaining operational efficiency.

Our solution implements an isolation strategy in two layers combining logical and physical separation. At the logical layer, each client’s data resides in dedicated schemas with access controls that prevent cross-tenant operations. Amazon Redshift Metadata security restricts tenants from discovering or accessing other clients’ schemas, tables, or database objects through system catalogs. At the physical layer, for larger deployments, dedicated Amazon Redshift clusters provide workload separation at the compute level, preventing one tenant’s analytical operations from impacting another’s performance. This dual approach meets regulatory requirements for data isolation in the insurance industry through schema-level isolation within clusters for standard deployments and complete compute separation across dedicated clusters for larger-scale implementations.

The implementation uses stored procedures to automate security configuration, maintaining consistent application of access controls across tenants. This defense-in-depth approach combines schema-level isolation, system catalog lockdown, and selective permission grants to create a security model.

For data architects interested in implementing similar multi-tenant architectures, review Implementing Metadata Security for Multi-Tenant Amazon Redshift Environment.

Implementation considerations

Verisk’s architecture reveals three decision points for companies building similar systems.

When to adopt open table formats

Apache Iceberg proved essential for datasets requiring schema evolution and historical versioning. Data engineers should evaluate open table formats when analytical workloads span multiple engines (Amazon Redshift, Amazon Athena, Spark) or when regulatory requirements demand point-in-time data reconstruction.

Multi-tenant isolation strategy

Schema-level separation combined with metadata security prevented cross-tenant data discovery without performance overhead. This approach scales more efficiently than database-per-tenant architectures while meeting insurance industry compliance requirements. Security experts should implement isolation controls during initial deployment rather than retrofitting them later.

Stored procedures or application logic

Redshift stored procedures standardized aggregation calculations across teams and constructed dynamic SQL queries. This approach works best when business logic changes frequently or when multiple teams need different aggregation dimensions on the same datasets.

Conclusion

Verisk’s implementation of Amazon Redshift Serverless with Apache Iceberg and lakehouse architecture shows how separating compute from storage addresses enterprise analytics challenges at billion-record scale. By combining cost-effective Amazon S3 storage with Redshift’s massively parallel SQL compute, Verisk achieved aggregations across billions of catastrophe modeling records, reduced storage costs through efficient parquet compression, and eliminated ingestion delays. Now underwriting teams can run ad-hoc analyses during business hours rather than waiting for long-running batch jobs. The combination of open standards like Apache Iceberg, serverless compute with Amazon Redshift, and multi-tenant security provides the scalability, performance, and cost efficiency needed for modern analytics workloads.

Verisk’s journey has positioned them to scale confidently into the future, processing not just billions, but potentially trillions of records as their model resolution increases.


About the authors

Karthick Shanmugam

Karthick Shanmugam

Karthick is Head of Architecture at Verisk EES. Focused on scalability, security, and innovation, he drives the development of architectural blueprints that align technology direction with business objectives. He is dedicated to building a modern, adaptable foundation that accelerates Verisk’s digital transformation and enhances value delivery across global platforms.

Srinivasa Are

Srinivasa Are

Srinivasa is a Principal Data Architect at Verisk EES, with extensive experience driving cloud transformation and data modernization across global enterprises. Known for combining deep technical expertise with strategic vision, Srini helps organizations unlock the full potential of their data through scalable cost-optimized architectures on AWS—bridging innovation, efficiency, and meaningful business outcomes.

Raks Khare

Raks Khare

Raks is a Senior Analytics Specialist Solutions Architect at AWS based out of Pennsylvania. He helps customers across varying industries and regions architect data analytics solutions at scale on the AWS platform. Outside of work, he likes exploring new travel and food destinations and spending quality time with his family.

Duvan Segura-Camelo

Duvan Segura-Camelo

Duvan is a Senior Analytics & AI Solutions Architect at AWS based out of Michigan, he helps customers architect scalable data analytics and AI solutions. With over two decades of experience in Analytics, Big Data and AI, Duvan is passionate about helping organizations build advanced, highly scalable solutions on AWS. Outside of work, he enjoys spending time with his family, staying active, reading, and playing the guitar.

Ashish Agrawal

Ashish Agrawal

Ashish is a Principal Product Manager with Amazon Redshift, building cloud-based data warehouses and analytics cloud services. Ashish has over 25 years of experience in IT. Ashish has expertise in data warehouses, data lakes, and platform as a service. Ashish has been a speaker at worldwide technical conferences.

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

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

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

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

How OpenSearch Service stores indexes and performs queries

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

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

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

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

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

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

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

Storage planning

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

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

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

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

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

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

Data nodes:

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

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

Cluster manager nodes:

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

Sharding strategy

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

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

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

T-shirt sizing for an e-commerce workload

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

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

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

Scaling strategies for e-commerce workloads

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

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

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

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

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

Monitoring and operational best practices

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

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

Conclusion

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

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

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


About the authors

Raaga NG

Raaga NG

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

Harsh Bansal

Harsh Bansal

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

Aditya Challa

Aditya Challa

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

Abe Raghib

Abe Raghib

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

Amazon Athena adds 1-minute reservations and new capacity control features

Post Syndicated from Manan Nayar original https://aws.amazon.com/blogs/big-data/amazon-athena-adds-1-minute-reservations-and-new-capacity-control-features/

Many of you choose serverless services for your analytics workloads because of its simplicity and elasticity. But many of you running mission-critical queries face a common challenge: ensuring your high-priority workloads run when needed and without interference from other queries in your account.

Amazon Athena is a serverless interactive query service that makes it simple to analyze data using SQL. Capacity Reservations is a feature of Athena that addresses the need to run critical workloads by providing dedicated serverless capacity for the workloads you specify. With Capacity Reservations, you request capacity in the form of Data Processing Units (DPU) and you assign them to your workloads.

In this post, we highlight three new capabilities that make Capacity Reservations more flexible and easier to manage: reduced minimums for fine-grained capacity adjustments, an autoscaling solution for dynamic workloads, and capacity cost and performance controls.

Now available: 1-minute reservations and 4 DPU minimum

Yesterday, we announced a big change for Capacity Reservations: you can now reserve as few as 4 DPU (down from 24 DPU) for as little as 1 minute (down from 60 minutes). This update lets you make frequent, fine-grained capacity adjustments to closely match your workload patterns and hold less capacity, with savings up to 95% for workloads that complete in under an hour.

We’ve optimized Athena for interactive queries that need a quick response, but many of you use Athena for non-interactive queries as well. For example, you may have queries that run on a schedule to prepare data for downstream analysis or perform updates to Apache Iceberg tables. These queries often process larger volumes of data and run for longer than interactive queries. If you’re using Athena’s scan-based pricing option, all your queries count towards your account-level quota. This means that your latency sensitive interactive queries can sometimes end up queued behind non-interactive queries that are running in your account.

Capacity Reservations addresses prioritization problems like this by making it possible to assign dedicated capacity to Athena workgroups. For example, Twilio operates a query platform that serves 1,500+ users who run over 2.5 million queries per month. They use Capacity Reservations for important workloads that need to have dedicated capacity to run optimally and avoid competing with other workloads.

Capacity Reservations have worked well when your workloads have been large and predictable. For example, users accessing dashboards at the start of the workday, or data processing jobs that run continuously 24/7. However, you’ve told us that you wanted more flexibility to update your reservations more frequently, to better match changes in demand.

With Athena’s new 4 DPU and 1-minute minimums, you’re now able to adjust capacity more frequently and match demand more closely than before. We’re excited to see how these updates benefit your mission-critical query workloads.

Autoscaling for dynamic workloads

The reduced minimums enable frequent capacity adjustments, but making those adjustments requires effort. Consider a business intelligence workload that peaks in the morning as executives review dashboards but decreases throughout the day. You want this workload isolated so that high-priority queries aren’t queued behind less important queries.

With 1-minute minimums, you can now adjust capacity to closely track these patterns. However, manually adjusting capacity this frequently is tedious—you need to monitor utilization, decide when to scale, and periodically adjust DPU.

We recently launched an autoscaling solution that uses AWS Step Functions to orchestrate capacity adjustments. It monitors capacity utilization metrics that Athena emits to Amazon CloudWatch at 1-minute granularity, analyzes utilization signal over configurable intervals, then conditionally adds or removes DPU so you can maintain consistent performance even during traffic spikes.

We have made this available as a 1-click deployment from the Athena console: just click Set up autoscaling on the details page for your reservation. When you do, a AWS CloudFormation template sets up all the resources you need. Among the resources set up is the Step Functions state machine, which you can view by opening Athena’s left-side navigation menu and clicking Workflows.

You can also find the template and information on the configurable autoscaling parameters in our documentation. See Automatically adjust capacity in the Athena User Guide.

We chose Step Functions for this solution to enable extensibility and customization. Step Functions integrates tightly with AWS services and allows you to define sophisticated state machines in Amazon States Language, a JSON-based language for serverless workflows. This makes it straightforward to add conditional logic, integrate additional services, or modify the workflow to match your specific requirements.

Control DPU usage at the workgroup and query levels

Part of the ease of use and simplicity of Athena is that it allocates capacity to queries automatically based on their complexity. However, sometimes preventing a single query from using too much capacity or operating at a required level of concurrency is more important than individual query performance. We recently released new DPU cost and performance controls so you can set constraints on Athena’s capacity allocation behavior when you’re using Capacity Reservations.

You can control DPU allocation in two places: workgroup-level controls that apply to all queries in that workgroup, or per query using the StartQueryExecution API. Both approaches set a type of budget that Athena adheres to when planning queries and determining how much capacity to allocate.

You can set minimum and maximum DPU limits from 4 to 124 DPU in increments of 4. Setting a maximum prevents Athena from allocating more DPU than specified. When you set a minimum, you instruct Athena to allocate at least the specified DPU. This can be beneficial when you know that a specific query requires a specific number of DPU to run optimally for your use case. Set both to create a range. For example, a minimum of 4 and maximum of 16 lets Athena start with 4 DPU and scale to 16 if needed. Setting them to the same value forces queries to run on an exact number of DPU.

Controls that you set at the workgroup-level are visible in the workgroup details page and the reservation that the workgroup has been added to.

Last but not least: every query you run on reserved capacity now reports its DPU usage in the Athena console and GetQueryExecution / BatchGetQueryExecution APIs, giving you complete visibility into capacity utilization.

Moving to Capacity Reservations

Getting started with Capacity Reservations involves creating a reservation with your desired DPU count, then assigning workgroups to that reservation. For end users, nothing changes. You continue running queries as usual and no SQL changes are needed. For administrators, you create a Capacity Reservation with your desired DPU count, then assign workgroups to that reservation. Athena automatically routes queries from assigned workgroups to your reserved capacity, isolated from other queries in your account and no impact to your account-level concurrency quota.

Conclusion

These updates to Capacity Reservations give you greater flexibility and control over your Athena workloads. The reduced minimums let you adjust capacity in smaller increments and shorter time windows, allowing you to match your usage patterns more closely than before. Autoscaling eliminates the work of making those adjustments manually. And DPU controls give you fine-grained influence over how individual queries consume capacity. Together, these capabilities help you optimize costs, manage concurrency, and deliver predictable performance for your most critical workloads—all while preserving Athena’s serverless benefits.

To learn more, see Athena Capacity Reservations in the Athena User Guide, Athena pricing page, or create your first Capacity Reservation in the Athena console.


About the authors

Manan Nayar

Manan Nayar

Manan is a Software Engineer at AWS based in Vancouver with over 8 years of experience building high-scale distributed systems and data platforms. In his spare time, he’s an avid runner and hiker who enjoys exploring the outdoors and staying active.

Mario Alkhoury

Mario Alkhoury

Mario is a Software Engineer on the Athena team, where he works on distributed systems, Capacity Reservations, and drivers. Based in the San Francisco Bay Area, he enjoys reading and spending time with family and friends outside of work.

Saroj Yadav

Saroj Yadav

Saroj is a Software Development Manager with AWS, driving innovations in data analytics with Amazon Athena and previously AWS Glue. Over the last 25 years, she has scaled infrastructure and delivered software products for companies during periods of hypergrowth.

Pathik Shah

Pathik Shah

Pathik is a Sr. Analytics Architect at Amazon Web Services. He joined AWS in 2015 and has been focusing on the big data analytics space since then, helping customers build scalable and robust solutions using AWS analytics services.

Theo Tolv

Theo Tolv

Theo is a Principal Analytics Architect based in Stockholm, Sweden. He’s worked with small and big data for most of his career and has built applications running on AWS since 2008. In his spare time, he likes to tinker with electronics and read space opera.

Scott Rigney

Scott Rigney

Scott is a Principal Technical Product Manager with the Amazon Athena service and works out of Arlington, Virginia. Scott has worked in the data, analytics, and machine learning space for longer than he’d like to admit.

Amazon OpenSearch Ingestion 101: Set CloudWatch alarms for key metrics

Post Syndicated from Utkarsh Agarwal original https://aws.amazon.com/blogs/big-data/amazon-opensearch-ingestion-service-101-set-cloudwatch-alarms-for-key-metrics/

Amazon OpenSearch Ingestion is a fully managed, serverless data pipeline that simplifies the process of ingesting data into Amazon OpenSearch Service and OpenSearch Serverless collections. Some key concepts include:

  • Source – Input component that specifies how the pipeline ingests the data. Each pipeline has a single source which can be either push-based and pull-based.
  • Processors – Intermediate processing units that can filter, transform, and enrich records before delivery.
  • Sink – Output component that specifies the destination(s) to which the pipeline publishes data. It can publish records to one or more destinations.
  • Buffer – It is the layer between the source and the sink. It serves as temporary storage for events, decoupling the source from the downstream processors and sinks. Amazon OpenSearch Ingestion also offers a persistent buffer option for push-based sources
  • Dead-letter queues (DLQs) – Configures Amazon Simple Storage Service (Amazon S3) to capture records that fail to write to the sink, enabling error handling and troubleshooting.

This end-to-end data ingestion service can help you collect, process, and deliver data to your OpenSearch environments without the need to manage underlying infrastructure.

This post provides an in-depth look at setting up Amazon CloudWatch alarms for OpenSearch Ingestion pipelines. It goes beyond our recommended alarms to help identify bottlenecks in the pipeline, whether that’s in the sink, the OpenSearch clusters data is being sent to, the processors, or the pipeline not pulling or accepting enough from the source. This post will help you proactively monitor and troubleshoot your OpenSearch Ingestion pipelines.

Overview

Monitoring your OpenSearch Ingestion pipelines is crucial for catching and addressing issues early. By understanding the key metrics and setting up the right alarms, you can proactively manage the health and performance of your data ingestion workflows. In the following sections, we provide details about alarm metrics for different sources, monitors, and sinks. The specific values for the threshold, period, and datapoints to alarm used for alarms can vary based on the individual use case and requirements.

Prerequisites

To create an OpenSearch Ingestion pipeline, refer to Creating Amazon OpenSearch Ingestion pipelines. For creating CloudWatch alarms, refer to Create a CloudWatch alarm based on a static threshold.

You can enable logging for OpenSearch Ingestion Pipeline, which captures various log messages during pipeline operations and ingestion activity, including errors, warnings, and informational messages. For details on enabling and monitoring pipeline logs, refer to Monitoring pipeline logs

Sources

The entry point of your pipeline is often where monitoring should begin. By setting appropriate alarms for source components, you can quickly identify ingestion bottlenecks or connection issues. The following table summarizes key alarm metrics for different sources.

Source Alarm Description Recommended Action
HTTP/ OpenTelemetry requestsTooLarge.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The request payload size of the client (data producer) is greater than the maximum request payload size, resulting in the status code HTTP 413. The default maximum request payload size is 10 MB for HTTP sources and 4 MB for OpenTelemetry sources. The limit for the HTTP sources can be increased for the pipelines with persistent buffer enabled. The chunk size for the client can be reduced so that the request payload doesn’t exceed the maximum size. You can examine the distribution of payload sizes of incoming requests using the payloadSize.sum metric.
HTTP requestsRejected.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The request was sent to the HTTP endpoint of the OpenSearch Ingestion pipeline by the client (data producer), but the request wasn’t accepted by the pipeline, and it rejected the request with the status code 429 in the response. For persistent issues, consider increasing the minimum OCUs for the pipeline to allocate additional resources for request processing.
Amazon S3 s3ObjectsFailed.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The pipeline is unable to read some objects from the Amazon S3 source. Refer to REF-003 in Reference Guide below.
Amazon DynamoDB Difference for totalOpenShards.max - activeShardsInProcessing.value
Threshold: >0
Statistic: Maximum (totalOpenShards.max) and Sum (activeShardsInProcessing.value)
Datapoints to Alarm: 3 out of 3.Additional Note: refer REF-004 for more details on configuring this specific alarm.
It monitors alignment between total open shards that should be processed by the pipeline and active shards currently in processing. The activeShardsInProcessing.value will go down periodically as shards close but should never misalign from ‘totalOpenShards.max’ for longer than a couple of minutes. If the alarm is triggered, you can consider stopping and starting the pipeline, this option resets the pipeline’s state, and the pipeline will restart with a new full export. It is non-destructive, so it does not delete your index or any data in DynamoDB. If you don’t create a fresh index before you do this, you might see a high number of errors from version conflicts because the export tries to insert older documents than the current _version in the index. You can safely ignore these errors. For root cause analysis on the misalignment, you can reach out to AWS Support
Amazon DynamoDB dynamodb.changeEventsProcessingErrors.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The number of processing errors for change events for a pipeline with stream processing for DynamoDB. If the metrics report increasing values, refer to REF-002 in Reference Guide below
Amazon DocumentDB documentdb.exportJobFailure.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The attempt to trigger an export to Amazon S3 failed. Review ERROR-level logs in the pipeline logs for entries beginning with “Received an exception during export from DocumentDB, backing off and retrying.” These logs contain the complete exception details indicating the root cause of the failure.
Amazon DocumentDB documentdb.changeEventsProcessingErrors.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The number of processing errors for change events for a pipeline with stream processing for Amazon DocumentDB. Refer to REF-002 in Reference Guide below
Kafka kafka.numberOfDeserializationErrors.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The OpenSearch Ingestion pipeline encountered deserialization errors while consuming a record from Kafka. Review WARN-level logs in the pipeline logs and verify serde_format is configured correctly in the pipeline configuration and the pipeline role has access to the AWS Glue Schema Registry (if used).
OpenSearch opensearch.processingErrors.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
Processing errors were encountered while reading from the index. Ideally, the OpenSearch Ingestion pipeline would retry automatically, but for unknown exceptions, it might skip processing. Refer to REF-001 or REF-002 in Reference Guide below, to get the exception details that resulted in processing errors.
Amazon Kinesis Data Streams kinesis_data_streams.recordProcessingErrors.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The OpenSearch Ingestion pipeline encountered an error while processing the records. If the metrics report increasing values, refer to REF-002 in Reference Guide below, which can help in identifying the cause.
Amazon Kinesis Data Streams kinesis_data_streams.acknowledgementSetFailures.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The pipeline encountered a negative acknowledgment while processing the streams, causing it to reprocess the stream. Refer to REF-001 or REF-002 in Reference Guide below.
Confluence confluence.searchRequestsFailed.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
While trying to fetch the content, the pipeline encountered the exception. Review ERROR-level logs in the pipeline logs for entries beginning with “Error while fetching content.” These logs contain the complete exception details indicating the root cause of the failure.
Confluence confluence.authFailures.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The number of UNAUTHORIZED exceptions received while establishing the connection Although the service should automatically renew tokens, if the metrics show an increasing value, review ERROR-level logs in the pipeline logs to identify why the token refresh is failing.
Jira jira.ticketRequestsFailed.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
While trying to fetch the issue, the pipeline encountered an exception. Review ERROR-level logs in the pipeline logs for entries beginning with “Error while fetching issue.” These logs contain the complete exception details indicating the root cause of the failure.
Jira jira.authFailures.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The number of UNAUTHORIZED exceptions received while establishing the connection. Although the service should automatically renew tokens, if the metrics show an increasing value, review ERROR-level logs in the pipeline logs to identify why the token refresh is failing.

Processors

The following table provides details about alarm metrics for different processors.

Processor Alarm Description Recommended Action
AWS Lambda aws_lambda_processor.recordsFailedToSentLambda.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
Some of the records could not be sent to Lambda. In the case of high values for this metric, refer to REF-002 in Reference Guide below.
AWS Lambda aws_lambda_processor.numberOfRequestsFailed.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The pipeline was unable to invoke the Lambda function. Although this situation should not occur under normal conditions, if it does, review Lambda logs and refer to REF-002 in Reference Guide below.
AWS Lambda aws_lambda_processor.requestPayloadSize.max
Threshold: >= 6292536
Statistic: MAXIMUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The payload size is exceeding the 6 MB limit, so the Lambda function can’t be invoked. Consider revisiting the batching thresholds in the pipeline configuration for the aws_lambda processor.
Grok grok.grokProcessingMismatch.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The incoming data doesn’t match the Grok pattern defined in the pipeline configuration. In the case of high values for this metric, review the Grok processor configurations and make sure the defined pattern matches according to the incoming data.
Grok grok.grokProcessingErrors.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The pipeline encountered an exception when extracting the information from the incoming data according to the defined Grok pattern. In the case of high values for this metric, refer to REF-002 in Reference Guide below.
Grok grok.grokProcessingTime.max
Threshold: >= 1000
Statistic: MAXIMUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The maximum amount of time that each individual record takes to match against patterns from the match configuration option. If the time taken is equal to or more than 1 second, check the incoming data and the Grok pattern. The maximum amount of time during which matching occurs is 30,000 milliseconds, which is controlled by the timeout_millis parameter.

Sinks and DLQs

The following table contains details about alarm metrics for different sinks and DLQs.

Sink Alarm Description Recommended Action
OpenSearch opensearch.bulkRequestErrors.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The number of errors encountered while sending a bulk request. Refer to REF-002 in Reference Guide below which can help to identify the exception details.
OpenSearch opensearch.bulkRequestFailed.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The number of errors received after sending the bulk request to the OpenSearch domain. Refer to REF-001 in Reference Guide below which can help to identify the exception details.
Amazon S3 s3.s3SinkObjectsFailed.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
The OpenSearch Ingestion pipeline encountered a failure while writing the object to Amazon S3. Verify that the pipeline role has the necessary permissions to write objects to the specified S3 key. Review the pipeline logs to identify the specific keys where failures occurred.
Monitor the s3.s3SinkObjectsEventsFailed.count metric for granular details on the number of failed write operations.
Amazon S3 DLQ s3.dlqS3RecordsFailed.count
Threshold: >0
Statistic: SUM
Period: 5 minutes
Datapoints to alarm: 1 out 1
For a pipeline with DLQ enabled, the records are either sent to the sink or to the DLQ (if they are unable to send to the sink). This alarm indicates the pipeline was unable to send the records to the DLQ due to some error. Refer to REF-002 in Reference Guide below which can help to identify the exception details.

Buffer

The following table contains details about alarm metrics for buffers.

Buffer Alarm Description Recommended Action
BlockingBuffer BlockingBuffer.bufferUsage.value
Threshold: >80
Statistic: AVERAGE
Period: 5 minutes
Datapoints to alarm: 1 out 1
The percent usage, based on the number of records in the buffer. To investigate further, check if the Pipeline is bottlenecked due to processors or sink by comparing timeElapsed.max metrics and analyzing bulkRequestLatency.max
Persistent persistentBufferRead.recordsLagMax.value
Threshold: > 5000
Statistic: AVERAGE
Period: 5 minutes
Datapoints to alarm: 1 out 1
The maximum lag in terms of number of records stored in the persistent buffer. If the value for bufferUsage is low, increase the maximum OCUs. If bufferUsage is also high [>80], investigate if pipeline is bottlenecked by processors or sink.

Reference Guide

The following provide guidance for resolving common pipeline issues along with general reference.

REF-001: WARN-level Log Review

Review WARN-level logs in the pipeline logs to identify the exception details.

REF-002: ERROR-level Log Review

Review ERROR-level logs in the pipeline logs to identify the exception details.

REF-003: S3 Objects Failed

When troubleshooting increasing s3ObjectsFailed.count values, monitor these specific metrics to narrow down the root cause:

  • s3ObjectsAccessDenied.count – This metric increments when the pipeline encounters Access Denied or Forbidden errors while reading S3 objects. Common causes include:
  • Insufficient permissions in the pipeline role.
  • Restrictive S3 bucket policy not allowing the pipeline role access.
  • For cross-account S3 buckets, incorrectly configured bucket_owners mapping.
  • s3ObjectsNotFound.count – This metric increments when the pipeline receives Not Found errors while attempting to read S3 objects.

For further assistance with the recommended actions, contact AWS support.

REF-004: Configuring Alarm for difference in totalOpenShards.max and activeShardsInProcessing.value for Amazon DynamoDB source.

  1. Open the CloudWatch console at https://console.aws.amazon.com/cloudwatch/.
  2. In the navigation pane, choose Alarms, All alarms.
  3. Choose Create alarm.
  4. Choose Select Metric.
  5. Select Source.
  6. In source, following JSON can be used after updating the <sub-pipeline-name>, <pipeline-name> and <region>.
    {
        "metrics": [
            [ { "expression": "m1-e1", "label": "Expression2", "id": "e2", "period": 900 } ],
            [ { "expression": "FLOOR((m2/15)+0.5)", "label": "Expression1", "id": "activeShardsInProcessing", "visible": false, "period": 900 } ],
            [ "AWS/OSIS", "<sub-pipeline-name>.dynamodb.totalOpenShards.max", "PipelineName", "<pipeline-name>", { "stat": "Maximum", "id": "m1", "visible": false } ],
            [ ".", "<sub-pipeline name>.dynamodb.activeShardsInProcessing.value", ".", ".", { "stat": "Average", "id": "m2", "visible": false } ]
        ],
        "view": "timeSeries",
        "stacked": false,
        "period": 900,
        "region": "<region>"
    }

Let’s review couple of scenarios based on the above metrics.

Scenario 1 – Understand and Lower Pipeline Latency

Latency within a pipeline is built up of three main components:

  • The time it takes to send documents via bulk requests to OpenSearch,
  • the time it takes for data to go through the pipeline processors, and
  • the time that data sits in the pipeline buffer

Bulk requests and processors (last two items in the previous list) are the root causes for why the buffer builds up and leads to latency.

To monitor how much data is being stored in the buffer, monitor the bufferUsage.value metric. The only way to lower latency within the buffer is to optimize the pipeline processors and sink bulk request latency, depending on which of those is the bottleneck.

The bulkRequestLatency metric measures the time taken to execute bulk requests, including retries, and can be used to monitor write performance to the OpenSearch sink. If this metric reports an unusually high value, it indicates that the OpenSearch sink may be overloaded, causing increased processing time. To troubleshoot further, review the bulkRequestNumberOfRetries.count metric to confirm whether the high latency is due to rejections from OpenSearch that are leading to retries, such as throttling (429 errors) or other reasons. If document errors are present, examine the configured DLQ to identify the failed document details. Additionally, the max_retries parameter can be configured in the pipeline configuration to limit the number of retries. However, if the documentErrors metric reports zero, the bulkRequestNumberOfRetries.count is also zero, and the bulkRequestLatency remains high, it is likely an indicator that the OpenSearch sink is overloaded. In this case, review the destination metrics for additional details.

If the bulkRequestLatency metric is low (for example, less than 1.5 seconds) and the bulkRequestNumberOfRetries metric is reported as 0, then the bottleneck is likely within the pipeline processors. To monitor the performance of the processors, review the <processorName>.timeElapsed.avg metric. This metric reports the time taken for the processor to complete processing of a batch of records. For example, if a grok processor is reporting a much higher value than other processors for timeElapsed, it may be due to a slow grok pattern that can be optimized or even replaced with a more performant processor, depending on the use case.

Scenario 2 – Understanding and Resolving Document Errors to OpenSearch

The documentErrors.count metric tracks the number of documents that failed to be sent by bulk requests. The failure can happen due to various reasons such as mapping conflicts, invalid data formats, or schema mismatches. When this metric reports a non-zero value, it indicates that some documents are being rejected by OpenSearch. To identify the root cause, examine the configured Dead Letter Queue (DLQ), which captures the failed documents along with error details. The DLQ provides information about why specific documents failed, enabling you to identify patterns such as incorrect field types, missing required fields, or data that exceeds size limits. For example, find the sample DLQ objects for common issues below:

Mapper parsing exception:

{"dlqObjects": [{
        "pluginId": "opensearch",
        "pluginName": "opensearch",
        "pipelineName": "<PipelineName>",
        "failedData": {
            "index": "<IndexName>",
            "indexId": null,
            "status": 400,
            "message": "failed to parse field [<fieldname>] of type [integer] in document with id '<DocumentId>'. Preview of field's value: 'N/A' caused by For input string: \"N/A\"",
            "document": {<OriginalDocument>}
        },
        "timestamp": "…"
    }]}

Here, OpenSearch cannot store the text string “N/A” in a field that is only for numbers, so it rejects the document and stores it in the DLQ.

Limit of total fields exceeded:

{"dlqObjects": [{
        "pluginId": "opensearch",
        "pluginName": "opensearch",
        "pipelineName": "<PipelineName>",
        "failedData": {
            "index": "<IndexName>",
            "indexId": null,
            "status": 400,
            "message": "Limit of total fields [<field limit>] has been exceeded",
            "document": {<OriginalDocument>}
        },
        "timestamp": "…"
    }]}

The index.mapping.total_fields.limit setting is the parameter that controls the maximum number of fields allowed in an index mapping, and exceeding this limit will cause indexing operations to fail. You can check if all those fields are required or leverage various processors provided by OpenSearch Ingestion to transform the data.

Once these issues are identified, you can either correct the source data, adjust the pipeline configuration to transform the data appropriately, or modify the OpenSearch index mapping to accommodate the incoming data format.

Clean up

When setting up alarms for monitoring your OpenSearch Ingestion pipelines, it’s important to be mindful of the potential costs involved. Each alarm you configure will incur charges based on the CloudWatch pricing model.

To avoid unnecessary expenses, we recommend carefully evaluating your alarm requirements and configuring them accordingly. Only set up the alarms that are essential for your use case, and regularly review your alarm configurations to identify and remove unused or redundant alarms.

Conclusion

In this post, we explored the comprehensive monitoring capabilities for OpenSearch Ingestion pipelines through CloudWatch alarms, covering key metrics across various sources, processors, and sinks. Although this post highlights the most critical metrics, there’s more to discover. For a deeper dive, refer to the following resources:

Effective monitoring through CloudWatch alarms is crucial for maintaining healthy ingestion pipelines and maintaining optimal data flow.


About the authors

Utkarsh Agarwal

Utkarsh Agarwal

Utkarsh is a Cloud Support Engineer in the Support Engineering team at AWS. He provides guidance and technical assistance to customers, helping them build scalable, highly available, and secure solutions in the AWS Cloud. In his free time, he enjoys watching movies, TV series, and of course cricket! Lately, he is also attempting to master foosball.

Ramesh Chirumamilla

Ramesh Chirumamilla

Ramesh is a Technical Manager with Amazon Web Services. In his role, Ramesh works proactively to help craft and execute strategies to drive customers’ adoption and use of AWS services. He uses his experience working with Amazon OpenSearch Service to help customers cost-optimize their OpenSearch domains by helping them right-size and implement best practices.

Taylor Gray

Taylor Gray

Taylor is a Software Engineer in the Amazon OpenSearch Ingestion team at Amazon Web Services. He has contributed many features within both Data Prepper and OpenSearch Ingestion to enable scalable solutions for customers. In his free time, he enjoys pickle ball, reading, and playing Rocket League.

Sovereign failover – Design for digital sovereignty using the AWS European Sovereign Cloud

Post Syndicated from Ivo Kammerath original https://aws.amazon.com/blogs/architecture/sovereign-failover-design-for-digital-sovereignty-using-the-aws-european-sovereign-cloud/

Organizations operating across multiple jurisdictions need to consider the impact of regulatory changes or geopolitical events on their access to cloud infrastructure. This post explains how to design failover architectures that span AWS partitions—including the AWS European Sovereign Cloud, AWS GovCloud (US) and other AWS Regions in the global infrastructure — so workloads can continue operating when sovereignty requirements shift.

Although the AWS European Sovereign Cloud is designed to help customers with operational autonomy and data residency requirements, it can also be used to address broader geopolitical and sovereignty risks. This post explores the architectural patterns, challenges, and best practices for building cross-partition failover, covering network connectivity, authentication, and governance. By understanding these constraints, you can design resilient cloud-native applications that balance regulatory compliance with operational continuity.

Understanding sovereignty risks

Digital sovereignty entails managing digital dependencies — deciding how data, technologies, and infrastructure are used, and reducing the risk of loss of access, control, or connectivity. As with any disaster recovery strategy, there are several means to provide continuity for the systems to be designed. Most of them involve some form of failover architecture, i.e. providing a second set of infrastructure to be used when the disaster incapacitates the original infrastructure. What differs for sovereign disaster recovery are the control mechanics and structures of the target to fail over to. Incorporating the AWS European Sovereign Cloud into your workload design adds failover capabilities that help you to reestablish or maintain enhanced sovereignty if the primary environment becomes unavailable.

As regulatory requirements evolve, modern failover architectures must account for sovereign environments such as the AWS European Sovereign Cloud, AWS GovCloud (US), and multi-vendor deployments. This post focuses on three core areas for incorporating sovereignty requirements into failover design: failover strategy, network connectivity across isolated partitions, and authentication and authorization in cross-partition architectures. These patterns apply to both short regional outages and long-term partition failures.

Understanding AWS partitions

As a global cloud provider, AWS operates multiple infrastructure partitions tailored to meet specific operational and regulatory requirements. In addition to its AWS global infrastructure, AWS offers specialized partitions such as AWS GovCloud for US government agencies, the AWS China Regions, and the AWS European Sovereign Cloud for customers that require stringent data residency and control within the EU.

Each partition is a logically isolated group of AWS Regions with its own set of resources, including AWS Identity and Access Management (IAM). Because of this separation, partitions act as hard boundaries. Credentials don’t carry over, and services such as Amazon S3 and features like S3 Cross-Region Replication or AWS Transit Gateway inter-region peering cannot function across partitions. These limitations are intentional, providing operational isolation. AWS GovCloud (US), launched in 2011, supports US public sector customers with compliance needs such as FedRAMP and ITAR. The AWS China regions are operated through local partnerships to meet Chinese data sovereignty laws. Similarly, the AWS European Sovereign Cloud is a partition built entirely within the EU, launched in 2026.

These partitions provide enhanced data control and physical infrastructure isolation, making them essential if you operate in regulated sensitive sectors and need to satisfy strict compliance requirements.

Key benefits of AWS partitions

AWS introduced partitions for several reasons. They are key to helping customers meet country-specific compliance and regulatory requirements, whether in AWS GovCloud (US), AWS China, or the AWS European Sovereign Cloud. This is underpinned by multiple safeguards and controls, including physical, logical, and operational separation of the cloud infrastructure between partitions. This directly corresponds to the security aspects of partitions. Partitions allow AWS to provide a complete isolation of resources, which helps manage security, especially for architectures running sensitive workloads.

Another important point to keep in mind when talking about partitions is service availability. Not all AWS services are available in every partition. To learn more about the AWS services available by Region, refer to AWS Capabilities by Region.

Cross-partition architectures

A cross-partition architecture enables partition failover by deploying resources and infrastructure across multiple isolated AWS partitions. Because partitions are fully separated by identity, networking, and service boundaries, failover can’t simply switch between them as within a single partition or region. Instead, environments must be pre-provisioned and kept in sync through internal or external tooling. Without such an architecture, failover between partitions is impractical. Cross-partition architectures make failover possible but require duplicate infrastructure, separate identity systems, and custom data synchronization.

Figure 1: Different reasons for failover and their possible locations

When designing cross-Region or cross-partition failover strategies, the choice of Regions depends on the type of disaster you want to mitigate:

  • Natural disasters – select Regions in different geographic zones or with distinct geographic features.
  • Technical disasters – separate workloads across independent parts of the global technical infrastructure, such as power grids, networks, and other shared resources.
  • Human-driven disasters – consider political, socioeconomic, and legal factors that might affect operations.

Figure 2: Active-active failover scenario including a sovereign failover option

Partition failover

Cross-partition workloads arise from industry needs to maintain continuity across sovereign domains while meeting regional regulations. Examples include military and defense connecting specialized clouds (such as AWS GovCloud (US)) with commercial environments, and emergency response systems requiring secure partition isolation combined with unified management (a single pane of glass approach). Control planes managing workloads across partitions are critical for handling multi-tenant structures, enabling centralized metrics, log aggregation, onboarding, security management and more.

However, cross-partition connections increase operational complexity, security and compliance overhead, costs, and governance challenges. These factors make it important to implement such architectures only when they are truly required. Standard cloud resilience models range from simple backups to multi-site setups, and can be implemented across multiple Availability Zones as well as multiple Regions. The same concept equally applies across multiple partitions. We can move backups into a second partition to be able to recover into that partition. Equally we can run an application pilot light in another partition. This greatly reduces the cost of the infrastructure required in the second partition because it will only be built up when needed. Finally, warm standby or multi-site active-active setups mainly differ in the need for more complex network synchronization across partitions.

Figure 3: Different types of disaster recovery scenarios

You might also consider vendor independence as an additional sovereignty requirement when planning failover. One way to achieve vendor independence is to use another cloud provider. However, failing over to another AWS partition is simpler than switching cloud providers because you can reuse your infrastructure as code templates across partitions.

Reasons to connect partitions

Although partitions are designed for isolation, some workloads within a partition might need to communicate with workloads in less regulated partitions or with external systems accessible over the public internet. For such instances several architectural strategies and the corresponding architectural decisions should be considered. There might be use cases where you need AWS Services to communicate across partitions and orchestrate actions spanning multiple partitions, such as:

  • Cross-domain applications
  • Feature parity and service availability
  • Cost-optimization while meeting security demands
  • Infrastructure consolidations
  • Control plane patterns

Implementing these use cases requires a deeper look into the technical aspects of connecting partitions from both a network standpoint and a security standpoint.

Regional connections vs. connected partitions

Regional connections let you link AWS Regions within the same partition using features like S3 Cross-Region Replication and Transit Gateway peering, facilitating relatively seamless workload distribution and failover within the partition’s global infrastructure. Understanding the distinction between regional connections and connected partitions is crucial for designing resilient, compliant architectures that meet both operational and regulatory demands.

Connecting partition networks

You can connect AWS partitions in three ways: internet connectivity secured by TLS, IPsec Site-to-Site VPN over the internet, or through an AWS Direct Connect gateway to on-premises routers or using Direct Connect point of presence (PoP) partner connections to another Direct Connect PoP. Each approach offers different trade-offs in terms of security complexity and recovery. For more information about connectivity patterns between AWS GovCloud (US) and the global AWS infrastructure, see Connectivity patterns between AWS GovCloud (US) and AWS commercial partition. In addition to the customer gateway solution shown previously, partners located in Direct Connect PoPs can provide cross-partition connectivity services. These services can move traffic from one Direct Connect PoP to another. Such a setup enables dedicated lines between the AWS European Sovereign Cloud Direct Connect PoPs and the Direct Connect locations in other partitions.

Because IAM credentials don’t work across partitions, you need to create separate roles or use external identity providers. Common approaches include using IAM roles with trust relationships and external IDs, AWS Security Token Service (AWS STS) regional endpoints, resource-based policies, or cross-account roles managed through AWS Organizations. A modern best practice is to federate identities from a single, centralized identity provider to multiple partitions, avoiding the need for IAM users wherever possible. If IAM users are still used, credentials can be stored in AWS Secrets Manager, rotated using Lambda, and a backup user can improve availability. These patterns are often combined with standard access controls, such as Amazon API Gateway with authorizers, to secure cross-partition interactions. For a deeper dive into cross-partition authentication and authorization with AWS IAM, see IAM Identity Center for AWS environments spanning AWS GovCloud (US) and standard Regions.

When securing communication between AWS partitions, certificate-based approaches present both opportunities and challenges. Because AWS Certificate Manager (ACM) certificates and AWS Private Certificate Authority (AWS Private CA) are bound to individual partitions, you must typically deploy and manage separate public key infrastructure (PKI) infrastructures in each environment, including dedicated root CAs and manual handling of private key transfers. To establish secure cross-partition communication, a more advanced solution involves using double-signed certificates, where root CAs in each partition cross-sign each other’s certificates, creating a bidirectional chain of trust. Implementing this requires setting up root CAs with AWS Certificate Manager Private CA, establishing cross-signing agreements, managing trust stores across partitions, and handling complex certificate validation and revocation checks. You must also comply with differing regulatory requirements and maintain detailed audit trails. Although this approach adds operational complexity, it is essential for enabling authenticated, encrypted communication across isolated partitions, particularly in regulated environments where security and compliance are paramount.

Managing AWS Organizations across partitions

Setting up AWS European Sovereign Cloud accounts within your AWS Organization must be done in a completely separate organization. In the AWS GovCloud (US) partition, accounts can be paired into a commercial organization, as described in Inviting Accounts into an Organization for AWS GovCloud. With sovereignty as the main goal, failing over to an AWS European Sovereign Cloud-only state is simpler if the AWS Organizations setup is separate from the start. This doesn’t require starting from scratch. Instead, you can manage the same organizational units (OUs) and policies for the AWS European Sovereign Cloud by reusing your existing deployment automation.

Ideally, AWS Organizations account structures should be separated to make it straightforward to use the AWS landscape within the AWS European Sovereign Cloud without relying on the other partitions.

Figure 4: connectivity and service distribution across AWS partitions like the AWS European Sovereign Cloud

Security controls should be tailored per partition using distinct Service Control Policies (SCPs), with AWS Control Tower managing the commercial side. Networking requires isolated Transit Gateways, separate Amazon Route 53 DNS zones, and secure cross-partition communication using AWS PrivateLink. For monitoring, AWS Config aggregators and AWS Security Hub instances must be configured separately in each partition, while consolidated billing can be managed through Organizations. It’s important to consider limitations (for example, AWS Control Tower can’t directly manage AWS GovCloud (US) or AWS European Sovereign Cloud accounts), and the limited availability of some AWS Organizations features in these partitions. Overall, this approach supports governance, security, and operational clarity across partitions.

Conclusion

Navigating sovereignty-driven cloud architectures requires a strategy that addresses partition isolation, network connectivity, and secure cross-partition authentication. Prioritizing sovereignty in failover design adds complexity, but it might be worth the trade-off if your workloads need protection against geopolitical risks or regulatory changes. Start by identifying the disaster scenarios that matter most to your business, then select the simplest architecture that addresses those risks. By designing proactively for evolving regulations, you can maintain both compliance and resilience in the cloud.


About the authors

More room to build: serverless services now support payloads up to 1 MB

Post Syndicated from Anton Aleksandrov original https://aws.amazon.com/blogs/compute/more-room-to-build-serverless-services-now-support-payloads-up-to-1-mb/

To support cloud applications that increasingly depend on rich contextual data, AWS has raised the maximum payload size from 256 KB to 1 MB for asynchronous AWS Lambda function invocations, Amazon Simple Queue Service (Amazon SQS), and Amazon EventBridge. Developers can use this enhancement to build and maintain context-rich event-driven systems and reduce the need for complex workarounds such as data chunking or external large object storage.

Overview

Modern cloud applications rely on context-rich, structured data to drive intelligent behavior. Large language model (LLM) prompts, telemetry signals, personalization data, machine learning (ML) outputs, and user interaction logs are no longer simple strings. Instead, they’re typically complex, nested JSON or YAML objects carrying meaningful context. Previously, developers working with serverless services such as Amazon SQS, Lambda (asynchronous invocations and Amazon SQS event-source mapping), or EventBridge had to carefully manage their data to fit within the 256 KB payload size limit. This commonly meant chunking larger payloads, externalizing payloads to object stores such as Amazon S3, or using data compression. These workarounds added complexity and latency, creating edge cases that were difficult to monitor and debug.

With the recent launches, you can now transmit payloads up to 1 MB, significantly reducing the need for complex data chunking and architectural workarounds. This increased capacity streamlines design patterns, reduces operational overhead, and makes event-driven systems more intuitive to build and maintain. Developers can now include richer data in single payloads—from detailed LLM prompts and full system states to comprehensive context and complete transaction histories.

The new 1 MB payload size limit applies to asynchronous Lambda function invocations, whether you trigger them using either SQS event-source mapping, AWS Command Line Interface (AWS CLI), AWS SDKs, Lambda Invoke API, or AWS services such as EventBridge. The increased limit also extends to all messages and events flowing through Amazon SQS queues and EventBridge Event Buses.

Getting started

There’s nothing you need to do to get started. This enhancement is automatically applied to all new and existing Lambda functions, SQS queues, and EventBridge Event Buses.

If you were previously chunking data at 256KB (or lower) threshold, then you might need to make changes to your service configurations or business logic code to start using the new limit. For example, if you’ve explicitly set Amazon SQS MaximumMessageSize attribute, then you might need to adjust it to a new desired value. Larger payloads might also result in higher costs, as described in the following section.

Real-world example: rich event context in agentic event-driven architectures

Event-driven architectures allow services to operate independently without centralized coordination. In these systems, comprehensive event context is essential. With the increased 1 MB payload limit, events can now carry more comprehensive data—from user profiles and order details to historical interactions. This enables services such as inventory, shipping, and notifications to act autonomously.

Consider the following example. In hospitality and quick-service industries, customer satisfaction depends on timely, thoughtful service recovery. When a guest submits negative feedback through a survey, review, or complaint form, service teams must gather context, interpret the issue, and craft a response. Traditionally, this meant manually piecing together visit logs, loyalty data, and prior complaints. Now, this can be fully automated using an AI agent powered by AWS serverless services and Amazon Bedrock, as shown in the following figure.

Figure 1: Customer feedback processing pipeline

The workflow:

  1. Receive: A new review is submitted through the Review application and emitted as an event to EventBridge Event Bus.
  2. Detect: Event Bus delivers the event to downstream Feedback analysis agent. The agent running in a Lambda function recognizes the review as low-rating or complaint.
  3. Enrich: The agent collects the guest’s visit metadata, booking details, loyalty activity, and complaint history using attached MCP tools into a single structured JSON payload (up to 1 MB).
  4. Queue: The payload is sent to an SQS queue for further asynchronous processing by downstream components.
  5. Generate: A separate Lambda function polls messages from Amazon SQS and invokes an Amazon Bedrock model to analyze the full complaint context, draft a personalized response, suggest a gesture (such as a refund or credit), and classify issue severity.
  6. Deliver: The message is logged and sent to the customer, and to the service team for further analysis.

This use case demonstrates the importance of having a rich context: current and previous visits details, loyalty tier, prior interactions, and feedback history. Previously, teams had to offload pieces of context to Amazon S3 and reference them externally, adding latency and architectural complexity. The new 1 MB payload size means that all this information can be transported together, improving the serverless agentic workflow efficiency and streamlining maintenance.

Best practices when using large payloads

The following sections outline best practices that you should apply when using larger payloads.

Performance considerations

Monitor Lambda function memory usage carefully when working with larger payloads, because parsing and processing complex JSON objects can increase memory usage and execution duration. Test your systems thoroughly under load, especially for high-throughput applications, by benchmarking with realistic payload sizes and traffic patterns. Although the payload limit has increased to 1 MB, the Lambda 15-minute timeout and memory limits remain unchanged. When applicable, you can use compression to process even larger datasets efficiently, but remember to account for the added CPU overhead of compression and decompression in your performance calculations. Read the Monitoring best practices for event delivery with Amazon EventBridge post for more best practices to tune your event-driven architectures performances.

Operational guidelines

Configure dead-letter-queues (DLQ) to make sure that failed messages are retained for inspection and troubleshooting. This becomes especially important with larger payloads, because debugging complex data structures necessitates access to the complete message context. Implement robust error handling and retries to manage transient failures, particularly when processing rich payload content that may contain nested structures or complex relationships.

To further optimize throughput, you can batch similar smaller events together into a single payload. However, avoid mixing unrelated events and maintain clear boundaries between different business domains and processes.

Always make sure that your downstream dependencies are capable of handling larger payloads.

When to use external storage

Even with the increased 1 MB payload limit, there are scenarios where patterns such as claim check remain a sound architectural choice. These patterns involve storing a full payload in an external system, such as Amazon S3, and passing a lightweight reference through your event stream. This approach continues to provide value when payloads exceed the new limit, when data needs to be reused by multiple consumers, or when strict governance, traceability, and security requirements are involved. For example, audit logs, image metadata, or large ML inference inputs may still surpass the 1 MB boundary, even when compressed. Instead of risking truncation or fragmentation, a claim check enables consistent, scalable access to the complete data set.

You can use open source libraries such as the Kafka sink connector for EventBridge and Amazon SQS Extended Client Library (available for Python and Java) that abstract complexities of storing large objects in external storage.

Cost management

Although larger payloads enable richer context in your applications, logging full payloads can increase storage and processing costs. Services such as CloudWatch Logs charge based on data volume, thus implementing selective logging, payload truncation, or sampling becomes crucial for high-volume events. Consider logging only essential fields or implementing smart sampling strategies based on business importance.

For full payload archival and retention, evaluate cost-effective storage solutions such as Amazon S3 with appropriate lifecycle policies. This can include moving older logs to cheaper storage tiers or implementing automated cleanup procedures for non-critical data. Balance your retention needs with cost optimization by defining clear policies for what data needs to be kept and for how long.

Review the pricing pages for AWS Lambda, Amazon EventBridge, and Amazon SQS to learn about the costs of delivering and processing events and messages.

Conclusion

The increase in maximum payload size from 256 KB to 1 MB enables developers to build more efficient distributed architectures. You can use this enhancement to transport richer context in event and message payloads, reducing the need for complex workarounds that previously added architectural complexity and operational overhead. This added room to transmit rich context means that you can streamline your workflows, improve observability, and reduce architectural complexity whether using choreography or orchestration patterns.

Go to the developer guides for AWS Lambda, Amazon EventBridge, and Amazon SQS, to learn more about how to take advantage of this update.

To learn more about serverless architectures, visit Serverless Land.

File integrity monitoring with AWS Systems Manager and Amazon Security Lake 

Post Syndicated from Adam Nemeth original https://aws.amazon.com/blogs/security/file-integrity-monitoring-with-aws-systems-manager-and-amazon-security-lake/

Customers need solutions to track inventory data such as files and software across Amazon Elastic Compute Cloud (Amazon EC2) instances, detect unauthorized changes, and integrate alerts into their existing security workflows.

In this blog post, I walk you through a highly scalable serverless file integrity monitoring solution. It uses AWS Systems Manager Inventory to collect file metadata from Amazon EC2 instances. The metadata is sent through the Systems Manager Resource Data Sync feature to a versioned Amazon Simple Storage Service (Amazon S3) bucket, storing one inventory object for each EC2 instance. Each time a new object is created in Amazon S3, an Amazon S3 Event Notification triggers a custom AWS Lambda function. This Lambda function compares the latest inventory version with the previous one to detect file changes. If a file that isn’t expected to change has been created, modified, or deleted, the function creates an actionable finding in AWS Security Hub. Findings are then ingested by Amazon Security Lake in a standard OCSF format, which centralizes and normalizes the data. Finally, the data can be analyzed using Amazon Athena for one-time queries, or by building visual dashboards with Amazon QuickSight and Amazon OpenSearch Service. Figure 1 summarizes this flow:

Figure 1: File integrity monitoring workflow

Figure 1: File integrity monitoring workflow

This integration offers an alternative to the default AWS Config and Security Hub integration, which relies on limited data (for example, no file modification timestamps). The solution presented in this post provides control and flexibility to implement custom logic tailored to your operational needs and support security-related efforts.

This flexible solution can also be used with other Systems Manager Inventory metadata, such as installed applications, network configurations, or Windows registry entries, enabling custom detection logic across a wide range of operational and security use cases.

Now let’s build the file integrity monitoring solution.

Prerequisites

Before you get started, you need an AWS account with permissions to create and manage AWS resources such as Amazon EC2, AWS Systems Manager, Amazon S3, and Lambda.

Step 1: Start an EC2 instance

Start by launching an EC2 instance and creating a file that you will later modify to simulate an unauthorized change.

Create an AWS Identity and Access Management (IAM) role to allow the EC2 instance to communicate with Systems Manager:

  1. Open the AWS Management Console and go to IAM, choose Roles from the navigation pane, and then choose Create role.
  2. Under Trusted entity type, select AWS service, select EC2 as the use case, and choose Next.
  3. On the Add permissions page, search for and select the AmazonSSMManagedInstanceCore IAM policy, then choose Next.
  4. Enter SSMAccessRole as the role name and choose Create role.
  5. The new SSMAccessRole should now appear in your list of IAM roles:
Figure 2: Create an IAM role for communication with Systems Manager

Figure 2: Create an IAM role for communication with Systems Manager

Start an EC2 instance:

  1. Open the Amazon EC2 console and choose Launch Instance.
  2. Enter a Name, keep the default Linux Amazon Machine Image (AMI), and select an Instance type (for example, t3.micro).
  3. Under Advanced details:
    1. IAM instance profile, select the previously created SSMAccessRole
    2. Create a fictitious payment application configuration file in the /etc/paymentapp/ folder on the EC2 instance. Later, you will modify it to demonstrate a file-change event for integrity monitoring. To create this file during EC2 startup, copy and paste the following script into User data.
#!/bin/bash
mkdir -p /etc/paymentapp
echo "db_password=initial123" > /etc/paymentapp/config.yaml

Figure 3: Adding the application configuration file

Figure 3: Adding the application configuration file

  1. Leave the remaining settings as default, choose Proceed without key pair, and then select Launch Instance. A key pair isn’t required for this demo because you use Session Manager for access.

Step 2: Enable Security Hub and Security Lake

If Security Hub and Security Lake are already enabled, you can skip to Step 3.
To start, enable Security Hub, which collects and aggregates security findings. AWS Security Hub CSPM adds continuous monitoring and automated checks against best practices.

  1. Open the Security Hub console.
  2. Choose Security Hub CSPM from the navigation pane and then select Enable AWS Security Hub CSPM and choose Enable Security Hub CSPM at the bottom of the page.

Note: For this demo, you don’t need the Security standards options and can clear them.

Figure 4: Enable Security Hub CSP

Figure 4: Enable Security Hub CSP

Next, activate Security Lake to start collecting actionable findings from Security Hub:

  1. Open the Amazon Security Lake console and choose Get Started.
  2. Under Data sources, select Ingest specific AWS sources.
  3. Under Log and event sources, select Security Hub (you will use this only for this demo):
Figure 5: Select log and event sources

Figure 5: Select log and event sources

  1. Under Select Regions, choose Specific Regions and make sure you select the AWS Region that you’re using.
  2. Use the default option to Create and use a new service role.
  3. Choose Next and Next again, then choose Create.

Step 3: Configure Systems Manager Inventory and sync to Amazon S3

With Security Hub and Security Lake enabled, the next step is to enable Systems Manager Inventory to collect file metadata and configure a Resource Data Sync to export this data to S3 for analysis.

  1. Create an S3 bucket by carefully following the instructions in the section To create and configure an Amazon S3 bucket for resource data sync.
  2. After you created the bucket, enable versioning in the Amazon S3 console by opening the bucket’s Properties tab, choosing Edit under Bucket Versioning, selecting Enable, and saving your changes. Versioning causes each new inventory snapshot to be saved as a separate version, so that you can track file changes over time.

Note: In production, enable S3 server access logging on the inventory bucket to keep an audit trail of access requests, enforce HTTPS-only access, and enable CloudTrail data events for S3 to record who accessed or modified inventory files.

The next step is to enable Systems Manager Inventory and set up the resource data sync:

  1. In the Systems Manager console, go to Fleet Manager, choose Account management, and select Set up inventory.
  2. Keep the default values but deselect every inventory type except File. Set a Path to limit collection to the files relevant for this demo and your security requirements. Under File, set the Path to: /etc/paymentapp/.
Figure 6: Set the parameters and path

Figure 6: Set the parameters and path

  1. Choose Setup Inventory.
  2. In Fleet Manager, choose Account management and select Resource Data Syncs.
  3. Choose Create resource data sync, enter a Sync name, and enter the name of the versioned S3 bucket you created earlier.
  4. Select This Region and then choose Create.

Step 4: Implement the Lambda function

Next, complete the setup to detect changes and create findings. Each time Systems Manager Inventory writes a new object to Amazon S3, an S3 Event Notification triggers a Lambda function that compares the latest and previous object versions. If it finds created, modified, or deleted files, it creates a security finding. To accomplish this, you will create the Lambda function, set its environment variables, add the helper layer, and attach the required permissions.

The following is an example finding generated in AWS Security Finding Format (ASFF) and sent to Security Hub. In this example, you see a notification about a file change on the EC2 instance listed under the Resources section.

{
	...
"Id": "fim-i-0b8f40f4de065deba-2025-07-12T13:48:31.741Z",
	"AwsAccountId": "XXXXXXXXXXXX",
	"Types": [
		"Software and Configuration Checks/File Integrity Monitoring"
	],
	"Severity": {
		"Label": "MEDIUM"
	},
	"Title": "File changes detected via SSM Inventory",
	"Description": "0 created, 1 modified, 0 deleted file(s) on instance i-0b8f40f4de065deba",
	"Resources": [
		{
			"Type": "AwsEc2Instance",
			"Id": "i-0b8f40f4de065deba"
		}
	],
	...
}

Create the Lambda function

This function detects file changes, reports findings, and removes unused Amazon S3 object versions to reduce costs.

  1. Open the Lambda console and choose Create function in the navigation pane.
  2. For Function Name enter fim-change-detector.
  3. Select Author from scratch, enter a function name, select the latest Python runtime, and choose Create function.
  4. On the Code tab, paste the following main function and choose Deploy.
import boto3, os, json, re
from datetime import datetime, UTC
from urllib.parse import unquote_plus
from helpers import is_critical, load_file_metadata, is_modified, extract_instance_id

s3 = boto3.client('s3')
securityhub = boto3.client('securityhub')

CRITICAL_FILE_PATTERNS = os.environ["CRITICAL_FILE_PATTERNS"].split(",")
SEVERITY_LABEL = os.environ["SEVERITY_LABEL"]
	
def lambda_handler(event, context):
	# Safe event handling
	if "Records" not in event or not event["Records"]:
		return

	# Extract S3 event
	record = event['Records'][0]
	bucket = record['s3']['bucket']['name']
	key = unquote_plus(record['s3']['object']['key'])
	current_version = record['s3']['object'].get('versionId')
	if not current_version:
		return

	# Fetching the region name
	account_id = context.invoked_function_arn.split(":")[4]
	region = boto3.session.Session().region_name

	# Get object versions (latest first)
	versions = s3.list_object_versions(Bucket=bucket, Prefix=key).get('Versions', [])
	versions = sorted(versions, key=lambda v: v['LastModified'], reverse=True)

	# Find previous version
	idx = next((i for i,v in enumerate(versions) if v["VersionId"] == current_version), None)
	if idx is None or idx + 1 >= len(versions):
		return
	prev_version = versions[idx+1]["VersionId"]

	# Load both versions
	current = load_file_metadata(bucket, key, current_version)
	previous = load_file_metadata(bucket, key, prev_version)

	# Compare
	created = {p for p in set(current) - set(previous) if is_critical(p)}
	deleted = {p for p in set(previous) - set(current) if is_critical(p)}
	modified = {p for p in set(current) & set(previous) if is_critical(p) and is_modified(p, current, previous)}

	# Report if changes were found
	if created or deleted or modified:
		instance_id = extract_instance_id(bucket, key, current_version)
		now = datetime.now(UTC).isoformat(timespec='milliseconds').replace('+00:00', 'Z')
		finding = {
			"SchemaVersion": "2018-10-08",
			"Id": f"fim-{instance_id}-{now}",
			"ProductArn": f"arn:aws:securityhub:{region}:{account_id}:product/{account_id}/default",
			"AwsAccountId": account_id,
			"GeneratorId": "ssm-inventory-fim",
			"CreatedAt": now,
			"UpdatedAt": now,
			"Types": ["Software and Configuration Checks/File Integrity Monitoring"],
			"Severity": {"Label": SEVERITY_LABEL},
			"Title": "File changes detected via SSM Inventory",
			"Description": (
				f"{len(created)} created, {len(modified)} modified, "
				f"{len(deleted)} deleted file(s) on instance {instance_id}"
			),
			"Resources": [{"Type": "AwsEc2Instance", "Id": instance_id}]
		}
		securityhub.batch_import_findings(Findings=[finding])

	# No change – delete older S3 version
	else:
		if prev_version != current_version:
			try:
				s3.delete_object(Bucket=bucket, Key=key, VersionId=prev_version)
			except Exception as e:
				print(f"Delete previous S3 object version failed: {e}")

Note: In production, set Lambda reserved concurrency to prevent unbounded scaling, configure a dead letter queue (DLQ) to capture failed invocations, and optionally attach the function to an Amazon VPC for network isolation.

Configure environment variables

Configure the two required environment variables in the Lambda console. These two variables (one for critical paths to monitor and one for security finding severity) must be set or the function will fail.

  1. Open the Lambda console and choose Configuration and then select Environment variables.
  2. Choose Edit and then choose Add environment variable.
  3. Under Key, choose CRITICAL_FILE_PATTERNS
    1. Enter ^/etc/paymentapp/config.*$ as the value.
    2. Set the SEVERITY_LABEL to MEDIUM.
Figure 7: CRITICAL_FILE_PATTERNS and SEVERITY_LABEL configuration

Figure 7: CRITICAL_FILE_PATTERNS and SEVERITY_LABEL configuration

Set up permissions

The next step is to attach permissions to the Lambda function

  1. In your Lambda function, choose Configuration and then select Permissions.
  2. Under Execution role, select the role name that will lead to the role in IAM.
  3. Choose Add permissions and select Create inline policy. Select JSON view.
  4. Paste the following policy, and make sure to replace <bucket-name> with the name of your S3 bucket, and you also update <region> and <account-id> with your AWS Region and Account ID:
{
"Version": "2012-10-17",
"Statement": [
	{
		"Effect": "Allow",
		"Action": "securityhub:BatchImportFindings",
		"Resource": "arn:aws:securityhub:<region>:<account-id>:product/<account-id>/default"
	},
	{
		"Effect": "Allow",
		"Action": [
			"s3:GetObject",
			"s3:GetObjectVersion",
			"s3:ListBucketVersions",
			"s3:DeleteObjectVersion"
		],
		"Resource": [
			"arn:aws:s3:::<bucket-name>",
			"arn:aws:s3:::<bucket-name>/*"
			]
		}
	]
}

  1. To finalize, enter a Policy name and choose Create policy.

Add functions to the Lambda layer

For better modularity, add some helper functions to a Lambda layer. These functions are already referenced in the import section of the preceding Lambda function’s Python code. The helper functions check critical paths, load file metadata, compare modification times, and extract the EC2 instance ID.

Open AWS CloudShell from the top-right corner of the AWS console header, then copy and paste the following script and press Enter. It creates the helper layer and attaches it to your Lambda function.

#!/bin/bash
set -e
FUNCTION_NAME="fim-change-detector"
LAYER_NAME="fim-change-detector-layer"

mkdir -p python
cat > python/helpers.py << 'EOF'
import json, re, os
from dateutil.parser import parse as parse_dt
import boto3
s3 = boto3.client('s3')
CRITICAL_FILE_PATTERNS = os.environ.get("CRITICAL_FILE_PATTERNS", "").split(",")

def is_critical(path):
	return any(re.match(p.strip(), path) for p in CRITICAL_FILE_PATTERNS if p.strip())

def load_file_metadata(bucket, key, version_id):
	obj = s3.get_object(Bucket=bucket, Key=key, VersionId=version_id)
	data = {}
	for line in obj['Body'].read().decode().splitlines():
		if line.strip():
			i = json.loads(line)
			n, d, m = i.get("Name","").strip(), i.get("InstalledDir","").strip(), i.get("ModificationTime","").strip()
			if n and d and m: data[f"{d.rstrip('/')}/{n}"] = m
	return data

def is_modified(path, current, previous):
	try: return parse_dt(current[path]) != parse_dt(previous[path])
	except: return current[path] != previous[path]

def extract_instance_id(bucket, key, version_id):
	obj = s3.get_object(Bucket=bucket, Key=key, VersionId=version_id)
	for line in obj['Body'].read().decode().splitlines():
		if line.strip():
			r = json.loads(line)
			if "resourceId" in r: return r["resourceId"]
	return None
EOF

zip -r helpers_layer.zip python >/dev/null
LAYER_VERSION_ARN=$(aws lambda publish-layer-version \
	--layer-name "$LAYER_NAME" \
	--description "Helper functions for File Integrity Monitoring" \
	--zip-file fileb://helpers_layer.zip \
	--compatible-runtimes python3.13 \
	--query 'LayerVersionArn' \
	--output text)

aws lambda update-function-configuration \
	--function-name "$FUNCTION_NAME" \
	--layers "$LAYER_VERSION_ARN" >/dev/null
echo "Layer created and attached to the Lambda function."

Step 5: Set up S3 Event Notifications

Finally, set up S3 Event Notifications to trigger the Lambda function when new inventory data arrives.

  1. Open the S3 console and select the Systems Manager Inventory bucket that you created.
  2. Choose Properties and select Event notifications.
  3. Choose Create event notification.
    1. Enter an Event name.
    2. In the Prefix field, enter AWS%3AFile/ to limit Lambda triggers to file inventory objects only.
      Note: The prefix contains a : character, which must be URL-encoded as %3A.
    3. Under Event types, select Put.
    4. At the bottom, select your newly created Lambda function, and choose Save changes.

In this example, inventory collection runs every 30 minutes (48 times each day) but can be adjusted based on security requirements to optimize costs. The Lambda function is triggered once for each instance whenever a new inventory object is created. You can further reduce event volume by filtering EC2 instances through S3 Event Notification prefixes, enabling focused monitoring of high-value instances.

Step 6: Test the file change detection flow

Now that the EC2 instance is running and the sample configuration file /etc/paymentapp/config.yaml has been initialized, you’re ready to simulate an unauthorized change to test the file integrity monitoring setup.

  1. Open the Systems Manager console.
  2. Go to Session Manager and choose Start session.
  3. Select your EC2 instance and choose Start Session.
  4. Run the following command to modify the file:

echo “db_password=hacked456" | sudo tee /etc/paymentapp/config.yaml

This simulates a configuration tampering event. During the next Systems Manager Inventory run, the updated metadata will be saved to Amazon S3.

To manually trigger this:

  1. Open the Systems Manager console and choose State Manager.
  2. Select your association and choose Apply association now to start the inventory update.
  3. After the association status changes to Success, check your SSM Inventory S3 bucket in the AWS:File folder and review the inventory object and its versions.
  4. Open the Security Hub console and choose Findings. After a short delay, you should see a new finding like the one shown in Figure 8:
Figure 8: View file change findings

Figure 8: View file change findings

Step 7: Query and visualize findings

While Security Hub provides a centralized view of findings, you can deepen your analysis using Amazon Athena to run SQL queries directly on the normalized Security Lake data in Amazon S3. This data follows the Open Cybersecurity Schema Framework (OCSF), which is a vendor-neutral standard that simplifies integration and analysis of security data across different tools and services.

The following is an example Athena query:

SELECT
	finding_info.desc AS description,
	class_uid AS class_id,
	severity AS severity_label,
	type_name AS finding_type,
	time_dt AS event_time,
	region,
	accountid
FROM amazon_security_lake_table_us_east_1_sh_findings_2_0

Note: Be sure to adjust the FROM clause for other Regions. Security Lake processes findings before they appear in Athena, so expect a short delay between ingestion and data availability.
You will see a similar result for the preceding query, shown in Figure 9:

Figure 9: Athena query result in the Amazon Athena query editor

Figure 9: Athena query result in the Amazon Athena query editor

Security Lake classifies this finding as an OCSF 2004 Class, Detection Finding. You can explore the full schema definitions at OCSF Categories. For more query examples, see the Security Lake query examples.
For visual exploration and real-time insights, you can integrate Security Lake with OpenSearch Service and QuickSight, both of which now offer extensive generative AI support. For a guided walkthrough using QuickSight, see How to visualize Amazon Security Lake findings with Amazon QuickSight.

Clean up

After testing the step-by-step guide, make sure to clean up the resources you created for this post to avoid ongoing costs.

  1. Terminate the EC2 instance
  2. Delete the Resource Data Sync and Inventory Association
  3. Remove the Lambda function.
  4. Disable Security Lake and Security Hub CSPM
  5. Delete IAM roles created for this post
  6. Delete the associated SSM Resource Data Sync and Security Lake S3 buckets.

Conclusion

In this post, you learned how to use Systems Manager Inventory to track file integrity, report findings to Security Hub, and analyze them using Security Lake.
You can access the full sample code to set up this solution in the AWS Samples repository.
While this post uses a single-account, single-Region setup for simplicity, Security Lake supports collecting data across multiple accounts and Regions using AWS Organizations. You can also use a Systems Manager resource data sync to send inventory data to a central S3 bucket.

Getting Started with Amazon Security Lake and Systems Manager Inventory provides guidance for enabling scalable, cloud-centric monitoring with full operational context.

Adam Nemeth
Adam Nemeth

Adam is a Senior Solutions Architect and generative AI enthusiast at AWS, helping financial services customers by embracing the Day 1 culture and customer obsession of Amazon. With over 24 years of IT experience, Adam previously worked at UBS as an architect and has also served as a delivery lead, consultant, and entrepreneur. He lives in Switzerland with his wife and their three children.