AWS Architecture Blog

How Samsung achieved real-time pricing with AWS Lambda Response Streaming

8.5内容质量
How Samsung achieved real-time pricing with AWS Lambda Response Streaming

TL;DR · AI 摘要

三星通过AWS Lambda响应流和CloudFront实现实时定价,解决了传统架构中的延迟和缓存不一致问题。

核心要点

  • 使用AWS Lambda响应流和CloudFront可消除1小时的定价延迟。
  • 传统架构因预计算所有产品组合导致存储爆炸和资源浪费。
  • 实时定价引擎通过无状态架构实现低延迟和高一致性。

结构提纲

按章节快速跳转。

  1. 在高流量电商中,实时定价对防止价格不一致至关重要。

  2. 传统架构依赖异步缓存,导致价格与权威引擎脱节达1小时。

  3. 三星通过AWS Lambda响应流和CloudFront实现实时定价引擎

  4. §数据聚合陷阱

    预计算所有产品组合导致存储爆炸和资源浪费。

  5. 新架构显著减少延迟,提升价格一致性。

思维导图

用一张图看清主题之间的关系。

查看大纲文本(无障碍 / 无 JS 友好)
  • 实时定价解决方案
    • 传统架构问题
      • 异步缓存导致延迟
      • 预计算存储爆炸
    • 解决方案
      • AWS Lambda响应流
      • CloudFront边缘缓存
      • 无状态架构

金句 / Highlights

值得收藏与分享的关键句。

#AWS Lambda#CloudFront#实时定价#电商架构
打开原文

Samsung 如何通过 AWS Lambda 响应流实现实时定价 | AWS 架构博客

Samsung 如何通过 AWS Lambda 响应流实现实时定价

本文由 Samsung 电商的 Sathish Kumar 和 Christopher Chan 联合撰写。

在高流量的电商场景中,实现实时定价对于防止价格不一致至关重要。价格不一致会导致购物车冲击并削弱用户信任。这不是软件故障,而是架构延迟的体现,可以通过使用 AWS Lambda 响应流和 Amazon CloudFront 来解决,特别是在从多个后端源聚合数据的系统中。

在本文中,我们将介绍遗留架构面临的挑战、无状态流式处理解决方案、关键实现模式以及性能结果——这些模式可以应用于构建高流量 API 的场景,这些 API 需要从多个后端源聚合数据。

Samsung.com 是三星的主要直销渠道,销售智能手机、电视、家电和配件,每种产品都有多种变体、优惠和地区定价。这种复杂性使得实时价格准确性尤为重要。

在黑色星期五等高流量活动期间,Samsung 的 All Deals 和 Product Finder 页面会展示这些产品。为了保持这些高密度产品列表页面(PLPs)和比较表格的低延迟,遗留的基础设施依赖于异步缓存,这引入了缓存价格与权威定价引擎之间出现不同步的间隙。

问题:遗留的中间件缓存导致权威定价引擎与面向客户的页面之间出现 1 小时的不同步。

我们的方法:我们拆除了有状态的数据聚合(DA)架构,并使用 AWS Lambda 响应流和 Amazon CloudFront 边缘缓存构建了一个实时的批量仲裁引擎(一个无状态的编排层,可在请求时直接查询定价引擎)。

挑战:数据聚合的陷阱

当产品列表页面需要同时显示超过 30 个商品组合的价格时,单独为每个商品组合调用定价引擎的延迟变得无法接受。为了解决这个问题,我们构建了一个前端的后端(BFF)服务来执行“数据聚合”。这个 DA 服务的设计目的是将前端与重型的定价引擎解耦。

它依赖于一个定时的 Cron 工作器,每小时运行一次以获取整个产品目录。工作器随后会预先计算所有可能的产品排列组合的价格,并将它们存储在本地缓存中。

虽然这提高了读取速度,但也带来了两个显著的问题:

  1. 排列组合爆炸 – DA 服务必须预先计算所有可能的组合,以防万一客户查看了这些组合。
  • 数学计算:30 个产品 ×(变体 × 优惠 × 附加项)= 每页数千条记录
  • 存储影响:随着每个新产品变体的添加,缓存呈指数级增长
  • 资源浪费:大多数预先计算的组合从未被请求
  1. 同步延迟 – 由于 Cron 任务每小时只运行一次,价格变化(例如限时促销)会显著滞后。直到下一次计划同步之前,客户仍然会看到旧的价格。
  • 业务影响:限时促销在下一次运行之前会显示错误的旧价格
  • 用户信任:结账时价格与产品页面价格不一致,导致购物车冲击
  • 竞争劣势:具有实时定价的竞争对手获得了市场份额

遗留的数据聚合架构

架构图展示了定价引擎和CloudFront CDN之间的遗留数据聚合层。

图1:遗留数据聚合(DA)架构 遗留系统依赖于定时的Cron作业,这在权威方(定价)和客户之间创建了一个明显的“不同步层”。所有产品排列组合的预计算消耗了大量的存储和计算资源。

解决方案:无状态流式架构

中间层存储的数据最终会与源数据产生分歧,因此我们与AWS技术账户经理(TAM)和服务团队合作,设计了一种新的解决方案:批量仲裁引擎,这是一个无状态的编排层,可以在请求时直接查询定价引擎。

新架构遵循了“直通”模式:

  1. 客户端请求:浏览器使用一个单一的HTTP GET请求请求30个特定SKU的价格。
  1. 流式编排:一个AWS Lambda函数将这些30个请求并行地发送到定价引擎。
  1. 立即响应:随着定价引擎返回数据,Lambda函数立即将其流式传输给客户端,而无需缓冲。

为什么使用Lambda响应流式传输?

在确定采用这种方法之前,我们评估了多种替代方案:

  • 传统的请求-响应模式(缓冲)– 标准的Lambda调用在返回给客户端之前会缓冲完整的响应,这抵消了并行分发的延迟优势。对于30个并发SKU查找,这会增加数秒的等待时间。
  • 带有改进缓存的EC2 – 这是遗留方法。缓存层最终会与真实数据源产生偏差,这是我们想要解决的核心问题。
  • Lambda响应流式传输 – 这是唯一允许我们并行分发请求、在结果到达时立即流式传输(减少首次字节时间)并保持完全无状态(无需维护或使无效的中间缓存)的选项。

新的无状态流式架构

架构图展示了无状态流式解决方案,CloudFront直接连接到Lambda。

图2:无状态流式架构 新架构消除了中间缓存。一个高性能的流直接将用户连接到定价的真实数据源。CloudFront边缘位置缓存了95%的流量响应,而其余请求则直接发送到Lambda以获取实时定价。

实施步骤

转向这种新架构需要解决与CDN行为和冷启动相关的两个特定技术限制。我们通过三个步骤实现了该解决方案。

第一步:实现流式处理程序

我们解决方案的核心是一个用awslambda.streamifyResponse()包装的Node.js Lambda处理程序。这使我们能够通过转换和压缩流直接将数据传输到客户端,而无需等待所有数据都可用。

我们使用了一个自定义的NDJSONTransform,将定价对象转换为以换行符分隔的JSON(NDJSON),使浏览器能够逐个解析和渲染价格,而无需等待完整的响应。

code
// 带有响应流式传输的Lambda处理程序
// 文件:lambda-handler.js

import * as awslambda from "aws-lambda";
import * as zlib from "zlib";
import { pipeline } from "stream/promises";
import { NDJSONTransform } from "./transforms/ndjson-transform.js";

// 常量
const PROCESSING_MODE = process.env.PROCESSING_MODE || "MODE_QUERYSTRING";
const MODE_QUERYSTRING = "MODE_QUERYSTRING";

// Handler 必须使用 streamifyResponse() 包裹以启用流式传输 export const handler = awslambda.streamifyResponse( async (event, responseStream, _context) => { try { if (PROCESSING_MODE === MODE_QUERYSTRING) { // 设置响应元数据,包括状态码和头部信息 const httpResponseMetadata = { statusCode: 200, headers: { "Content-Type": "application/x-ndjson", "Content-Encoding": "gzip", "Cache-Control": "public, max-age=300", // 5 分钟缓存 "X-Custom-Header": "lambda-streaming", }, };

responseStream = awslambda.HttpResponseStream.from( responseStream, httpResponseMetadata );

const ndjsonTransform = new NDJSONTransform();

// 使用快速压缩级别创建 gzip 流 (Z_BEST_SPEED = Level 1) // Level 1 优先速度而非压缩比,适用于实时响应 const gzip = zlib.createGzip({ level: zlib.constants.Z_BEST_SPEED, });

// 处理请求,并写入 gzip 流 await processRequestUsingQueryString(event, ndjsonTransform, _context);

// 管道流程:转换 → 压缩 → 发送 await pipeline(ndjsonTransform, gzip, responseStream); }

// ... 错误处理 } catch (error) { console.error("Lambda 流式传输错误:", error); responseStream.destroy(); } finally { await flushMetrics().catch(console.error); } } );

code

用于并行分发请求的辅助函数:

// 处理请求:并行分发 SKU 查询 async function processRequestUsingQueryString(event, ndjsonTransform, context) { try { // 解析查询字符串以提取 SKU 列表 const queryString = event.rawQueryString || ""; const skus = parseCompressedQueryString(queryString);

// 使用 Promise.all 并行获取价格 const pricingPromises = skus.map((sku) => fetchPricingForSKU(sku).catch((err) => ({ sku: sku.id, error: err.message, })) );

const pricingResults = await Promise.all(pricingPromises);

// 将每个结果作为单独的 NDJSON 行发出 for (const result of pricingResults) { ndjsonTransform.write(result); }

ndjsonTransform.end(); } catch (error) { ndjsonTransform.destroy(error); } }

code

Handler 还使用了辅助函数来解析压缩的查询字符串 (parseCompressedQueryString),使用连接池获取单个 SKU 的价格 (fetchPricingForSKU),以及将指标刷新到 Amazon CloudWatch (flushMetrics)。

关键实现细节:

- awslambda.streamifyResponse() 包裹 handler,使其实时流式传输数据,而不是等待定价引擎的完整响应。

- NDJSONTransform 将对象转换为换行分隔的 JSON(每个对象一行)

- 使用 Z_BEST_SPEED(Level 1)进行 GZIP(GNU zip)压缩,优先速度而非压缩比

- pipeline() 处理错误传播和流清理

- 响应头包含用于 CloudFront 缓存的 Cache-Control

### 第 2 步:将请求数据压缩为 GET 请求

我们需要将复杂的请求数据(30 个 SKU,上下文元数据)发送到 API。

限制:CloudFront 和标准 HTTP 规范将 POST 请求视为非幂等的,这意味着它们默认不可缓存。

我们的方法:我们开发了一种密集且压缩的查询字符串格式,以将复杂的请求数据适配到标准的 GET 请求中。格式:g=group1(p=SKU-A:1:p=SKU-B:2)…

这使我们能够严格使用 GET 请求,将请求 URI 保持在标准长度限制内(约 800 字节),同时携带与 3-4KB JSON 主体相同的数据。

#### 客户端代码:构建压缩查询字符串

// 客户端:构建压缩查询字符串格式 // 文件:pricing-client.js

/**

  • 为批量定价请求构建压缩查询字符串
  • 格式:g=group1(p=SKU-A:1:p=SKU-B:2:p=SKU-C:3)
  • @param {Array} skus - 包含 { id, variant } 的 SKU 对象数组
  • @param {Object} context - 客户上下文 { customerId, region, sessionId }
  • @returns {string} 压缩查询字符串

*/ function buildPricingQueryString(skus, context = {}) { if (!skus || skus.length === 0) { throw new Error("SKU 数组不能为空"); }

if (skus.length > 30) { throw new Error("每个请求最多 30 个 SKU。请拆分为多个批次。"); }

// 构建 SKU 部分:p=SKU-001:1:p=SKU-002:2 const skuParts = skus .map((sku) => { const variant = sku.variant || 1; return p=${sku.id}:${variant}; }) .join(":");

// 构建上下文部分(可选) let contextPart = ""; if (context && Object.keys(context).length > 0) { const contextStr = Object.entries(context) .map(([key, value]) => ${key}=${value}) .join(":"); contextPart = :c=${contextStr}; }

// 最终格式:g=group1(p=SKU-A:1:p=SKU-B:2:c=customerId=123:region=US-EAST-1) return g=group1(${skuParts}${contextPart}); }

/**

  • 使用流式 NDJSON 响应解析获取定价数据
  • 核心模式:读取块,按换行符拆分,将每行解析为 JSON

*/ async function fetchPricingStream(skus, options = {}) { const queryString = buildPricingQueryString(skus); const url = ${options.baseUrl}?${queryString}; const response = await fetch(url, { method: "GET", headers: { Accept: "application/x-ndjson" }, });

// 使用 ReadableStream API 流式传输响应 const reader = response.body.getReader(); const decoder = new TextDecoder(); const pricingData = []; let buffer = "";

while (true) { const { done, value } = await reader.read(); if (done) break;

// 解码块并追加到缓冲区 buffer += decoder.decode(value, { stream: true });

// 按换行符拆分(NDJSON 格式:每行一个 JSON 对象) const lines = buffer.split("\n");

// 保留缓冲区中的最后一个不完整行 buffer = lines.pop() || "";

// 将每条完整行解析为 JSON 对象 for (const line of lines) { if (line.trim()) { const pricingObject = JSON.parse(line); pricingData.push(pricingObject);

// 可选:在每个价格到达时立即更新 UI if (options.onChunk) { options.onChunk(pricingObject); } } } }

reader.releaseLock(); return pricingData; }

code

页面加载时,客户端使用最多 30 个 SKU 和一个 onChunk 回调调用 fetchPricingStream,该回调在定价块到达时更新每个产品的 DOM(文档对象模型)元素。辅助函数处理更新单个价格元素、显示变体信息,并在定价暂时不可用时以友好的用户消息优雅降级。

### 第三步:为不可缓存的请求配置 CloudFront

为了使 CloudFront 能够有效地缓存这些复杂的 GET 请求,我们配置了一个精确的缓存策略,其中包括所有查询字符串和特定的请求头。

Terraform: CloudFront 缓存策略配置

文件:cloudfront-cache-policy.tf

resource "aws_cloudfront_cache_policy" "pricing-cache-policy" { name = "${var.lambda_function_name}-cache-policy"

TTL 配置

default_ttl = 300 # 5 分钟 max_ttl = 1800 # 30 分钟 min_ttl = 300 # 5 分钟

parameters_in_cache_key_and_forwarded_to_origin {

不基于 cookie 进行缓存(我们不使用 cookie)

cookies_config { cookie_behavior = "none" }

允许特定的请求头用于缓存键

headers_config { header_behavior = "whitelist" headers { items = [ "x-ecom-pricing-1", # 自定义请求头(已匿名化) "x-ecom-pricing-2", "x-ecom-pricing-3" ] } }

在缓存键中包含查询字符串

这可以确保不同的 SKU 组合具有独立的缓存条目

query_strings_config { query_string_behavior = "all" }

启用自动 GZIP 压缩

enable_accept_encoding_gzip = true } }

CloudFront 原始请求策略(将请求头转发给 Lambda)

resource "aws_cloudfront_origin_request_policy" "pricing-origin-policy" { name = "${var.lambda_function_name}-origin-policy"

headers_config { header_behavior = "whitelist" headers { items = [ "x-ecom-pricing-1", "x-ecom-pricing-2", "x-ecom-pricing-3" ] } }

query_strings_config { query_string_behavior = "all" }

cookies_config { cookie_behavior = "none" } }

code

CloudFront 分发本身配置为仅使用 HTTPS 的查看者协议,使用了前面部分中所示的缓存策略和原始请求策略,并通过 HTTPS 指向 Lambda 作为其原始服务器。

缓存策略亮点:

- 5 分钟的默认 TTL 在新鲜度和缓存效率之间取得平衡

- 缓存键中包含查询字符串(不同的 SKU 组合 = 独立的缓存条目)

- 请求头的白名单允许自定义定价变体

- 自动 GZIP 压缩减少带宽使用

- 5 至 30 分钟的 TTL 范围为不同内容提供灵活性

## 性能优化结果

我们通过四个不同的阶段优化了系统,使用 K6 负载测试脚本(500 个并发用户,每个请求 30 个商品)对每种配置进行测试,以模拟像黑色星期五这样的高流量事件。

### 第一阶段:基准线(全局 VPN)

我们使用默认的网络配置测试了初始的原型,其中所有出站流量(包括对 AWS 服务如 Lambda 的请求)都通过全局 VPN 路由,迫使流量进入公共网络,再返回 AWS 主干网络,增加了不必要的网络跳转和延迟。Lambda 使用的是标准的缓冲响应,没有压缩。结果并不理想(P90 为 4500 毫秒),因为连接开销主导了请求。

- DNS 解析:约 50 毫秒

- TCP 握手:约 100 毫秒

- TLS 协商:约 150 毫秒

- 总连接开销:每次调用约 300 毫秒

这种开销在业务逻辑运行之前就造成了巨大的延迟瓶颈。

- 首先:将 Lambda 移动到 Amazon 虚拟私有云(Amazon VPC)中,并直接与定价源进行对等连接,从而将内部调用的 DNS 和 TLS 开销降低到接近零。

- 其次:为 Lambda 启用预置并发,以消除冷启动延迟 500–1000 毫秒。

通过这些改进,P90 延迟降至 1,000 毫秒,提升了 4.5 倍,但仍然不够实时。

### 第三阶段:HTTP/2 和 GZIP 压缩

剩余的瓶颈是数据传输的大小。我们针对以下两个优化进行了改进:

- HTTP/2 多路复用:启用 HTTP/2 多路复用,以复用单个 TCP 连接进行 30 个并行的 SKU 查询,节省了累计握手时间的数秒。

- GZIP 压缩:应用 GZIP 压缩(等级 1 / Z_BEST_SPEED),将响应大小减少了 76%(170KB → 40KB)。

这两项优化使 P90 延迟降至 218 毫秒。

### 第四阶段:生产环境(边缘缓存)

在最终阶段,我们在优化后的 Lambda 上层添加了 CloudFront 边缘缓存。因为我们成功地将请求数据转换为 GET 请求(第二阶段),现在可以缓存 95% 的入站流量的计算价格。最终的 P90 延迟降至 50 毫秒。实际上,95% 的缓存命中率意味着每 20 个请求中只有一个实际调用 Lambda 函数;其余请求直接从离客户最近的 CloudFront 边缘位置提供服务。在黑色星期五等高峰事件期间,这意味着数百万的请求可以在边缘速度下处理,而无需接触源服务器,从而将延迟和计算成本保持在最低水平。

### 性能指标表

指标

第一阶段(基准)

第二阶段(VPC)

第三阶段(HTTP/2)

第四阶段(生产)

P50 延迟

1,670 毫秒

501 毫秒

176 毫秒

35 毫秒

P90 延迟

4,500 毫秒

1,000 毫秒

218 毫秒

50 毫秒

P99 延迟

5,100 毫秒

2,400 毫秒

500 毫秒

150 毫秒

缓存命中率

<1%

95%

响应大小

170 KB

40 KB

并发用户数

500

与基准相比的 P90 改进

1 倍

4.5 倍

20 倍

90 倍

#### 下图显示了每个优化阶段的 P90 延迟改进情况。

四个优化阶段的延迟改进情况。

K6 负载测试配置:

// 一个 K6 负载测试脚本示例 import http from 'k6/http'; import { check } from 'k6';

export let options = { stages: [ { duration: '2m', target: 100 }, // 逐步增加 { duration: '5m', target: 500 }, // 保持在 500 个并发用户 { duration: '2m', target: 0 }, // 逐步减少 ], thresholds: { http_req_duration: ['p(90)<100', 'p(99)<500'], // P90 < 100ms, P99 < 500ms }, };

export default function () { const url = 'https://api.example.com/pricing?g=group1(p=SKU-001:1:p=SKU-002:2:...)'; const res = http.get(url);

check(res, { 'status is 200': (r) => r.status === 200, 'response has content': (r) => r.body.length > 0, }); }

code

## 弹性、可扩展性和安全性考虑

除了延迟之外,我们还设计了系统以优雅地处理故障、在负载下进行扩展,并保护传输中的数据。

### 批处理限制

每个请求的 30 项限制是有意为之的。如果页面需要更多(例如 50 项),客户端逻辑会将它们拆分为多个并行批次。我们选择 30 是因为以下原因:

- Lambda 执行时间少于 5 秒

- 防止在高延迟期间出现超时问题

- 平衡并行请求与 Lambda 并发限制

- 典型的产品列表页面显示 20–30 个项目

async function fetchLargePricingBatch(skus) { const BATCH_SIZE = 30; const batches = [];

code

// 分成30项一组 for (let i = 0; i < skus.length; i += BATCH_SIZE) { batches.push(skus.slice(i, i + BATCH_SIZE)); }

// 并行获取所有批次 const results = await Promise.all( batches.map((batch) => fetchPricingStream(batch, { timeout: 30000 }) ) );

// 合并结果 return results.flat(); }

code

### 部分失败处理

流式架构具有弹性。如果某一项的价格获取失败,流不会崩溃,而是继续处理剩余的项目,因此用户仍然可以看到一个基本完整的页面。

部分失败处理:

async function processRequestUsingQueryString(event, ndjsonTransform, context) { const skus = parseCompressedQueryString(event.rawQueryString);

for (const sku of skus) { try { const pricing = await fetchPricingForSKU(sku); ndjsonTransform.write(pricing); } catch (error) { // 发送错误对象,但继续处理其他SKU ndjsonTransform.write({ sku: sku.id, variant: sku.variant, error: error.message, timestamp: new Date().toISOString(), }); } }

ndjsonTransform.end(); }

code

### 数据保护

- 虽然在客户端构建查询字符串会暴露请求结构,但这些数据(SKU、变体)本身已经是公开的。实际的价格逻辑和商业规则仍然安全地保留在定价引擎中。

- 传输中的数据使用TLS 1.3加密。

- 通过Amazon VPC端点连接到定价引擎(不暴露到互联网)。

- 不记录任何敏感数据(排除PII和定价算法)。

- 使用CloudTrail记录API调用,用于审计跟踪。

## 结论

过时的价格迫使工程团队在新鲜度和可扩展性之间做出选择。通过使用数据聚合模式,我们尝试同时保持两者,但由于计划同步的延迟,我们妥协了数据完整性。通过使用AWS Lambda响应流式传输和Amazon CloudFront,我们完全消除了同步层的需求。结果是一个系统,它能够提供用户流畅体验所需的50毫秒延迟,同时在产品页面和结账过程中保持价格一致性。

除了性能,这种架构显著减少了运营足迹:计算集群在高峰事件期间从超过100个自动扩展实例缩小到仅5到10个Lambda函数,从而降低了维护和运营成本。这一成果是三星电商工程团队、我们的AWS技术账户经理(TAM)以及Lambda和CloudFront服务团队紧密合作的结果,他们帮助设计解决方案、审查设计决策,并指导三星实现生产就绪。这种技术适用于类似的高流量数据聚合场景:产品目录、库存系统、推荐引擎或需要实时合并多个后端响应的服务。

要开始,首先确定你延迟最高的聚合端点,评估你的请求数据是否可以转换为可缓存的GET请求,并在迁移整个API之前,先为单个端点实现Lambda响应流式传输。

资源: – AWS Lambda响应流式传输文档 – Lambda响应流式传输教程 – Amazon CloudFront开发指南 – CloudFront缓存策略

## 了解更多

关于本文中讨论的概念和技术的更多信息:

- 介绍 AWS Lambda 响应流式传输 – 介绍本解决方案中使用的流式传输模式的 AWS Compute 博客文章

- 教程:创建响应流式传输 Lambda 函数 – 分步教程,用于构建您的第一个流式传输 Lambda 函数

- CloudFront 最佳实践 – 配置 CloudFront 分发的最佳实践

- NDJSON 格式 – 用于增量响应解析的换行分隔 JSON 规范

- Node.js Streams API – 支撑 Transform 和 Pipeline 模式的 Node.js 流文档

- Terraform AWS 提供商 – 用于 CloudFront 和 WAF 配置的基础设施即代码提供商

- AWS 规范性指导:负载测试 – AWS 关于负载测试工具(包括 K6)的指导

## 关于作者

'"`