
客户: 西门子能源 · 数据部门
行业: 能源 · 制造业
服务: 云架构、无服务器开发、AI/ML 集成、基础设施即代码(IaC)
所用 AWS 服务: S3、S3 Access Grants、Lambda、Step Functions、EventBridge、API Gateway、Bedrock Data Automation、Amplify、SQS、CloudFormation(使用 Python 编写的 CDK)
挑战
西门子能源的数据部门管理着海量的产品测量数据,例如在全球制造运营、涡轮机或其他设施设备中产生的文档、图像、音频录音和视频文件。
该团队需要一套能够满足以下要求的系统:
- 摄入任意大小的文件——从几 KB 到数 GB 的测量数据集以及高分辨率图像,同时不触及负载上限,也不降低性能。
- 自动提取丰富的元数据——借助 AI,从每一个上传的文件中提取,无论其模态如何(文本、图像、音频、视频)。不过,当某些元数据事先已知时,他们希望能够将这些元数据混入自动生成的元数据中。
- 实施分区级别的访问控制——使不同的团队和用户只能查看和修改各自指定范围内的数据,并与西门子能源现有的 Microsoft Entra ID(Azure AD)SSO 身份提供商相集成。
- 扩展至数十万个文件——且不受 OneDrive 或 SharePoint 等工具在生命周期管理上的限制(约 30 万个文件的上限)。
- 提供一个基于 Web 的界面——供非技术用户在通过全公司范围的 SSO 完成认证后,上传、浏览、下载和预览文件。
该系统必须达到生产级别、经过压力测试,并能够借助完整的基础设施即代码(IaC)在多个环境(沙盒、开发、UAT、生产)中部署。
解决方案
我们设计并交付了产品测量数据管道(PMDP)——一套相互连接的微服务和接口,全部以无服务器方式运行在 AWS 上。
架构概览
该平台由以下部分组成:
- 文件管理器 API:一个 RESTful 服务,负责处理文件的 CRUD、分块上传、预签名 URL 的生成、元数据操作,以及以分区为范围的访问控制。
- 文件元数据增强微服务:一条事件驱动的多模态 AI 管道,借助 Amazon Bedrock Data Automation,从上传的文件中自动提取结构化的元数据。
- Web 应用:一个由 AWS Amplify 托管的 React 应用,提供一个集成了西门子能源 SSO 的 Storage Browser 界面。
- 压力测试套件:一个 CLI 工具,通过对最大达数 GB 的文件进行并发上传和下载,来验证系统在负载下的表现。
所有基础设施均以 Python CDK 定义,并通过 CI/CD 流水线(自托管的 GitLab)进行验证和自动部署。
深入剖析:它是如何运作的
1. Hive 风格的对象键分区
每一个通过 API 上传的文件,都会采用 Apache Hive 风格的分区方案存储在 S3 中。这不仅仅是为了组织归类。它实现了高性能的查询,并直接映射到访问控制模型上。
# construct_object_key.py — 确定性的、可查询的 S3 键生成
class Capability(str, Enum):
EDAA = "edaa"
CUSTOMER_FACING = "customer-facing"
MANUFACTURING = "manufacturing"
QUALITY = "quality"
# ...
def construct_object_key(
capability: Capability | str | None = None,
file_name: str = "",
sub_capability: str | None = None,
year: str | None = None,
month: str | None = None,
day: str | None = None,
) -> str:
now = datetime.now()
generated_hash = hashlib.md5(
f"{capability}{sub_capability}{file_name}".encode()
).hexdigest()[:2]
return (
f"capability={capability}/"
f"sub-capability={sub_capability or 'unknown'}/"
f"year={year or now.year}/"
f"month={month or f'{now.month:02d}'}/"
f"day={day or f'{now.day:02d}'}/"
f"hash={generated_hash}/"
f"{file_name}"
)一个以 report.pdf 为名、在 manufacturing 能力域下于 2025 年 12 月 23 日上传的文件,会变成:
capability=manufacturing/sub-capability=unknown/year=2025/month=12/day=23/hash=a3/report.pdf
这种结构使得 Athena、Glue 或任何与 Hive 兼容的工具,都能按能力域、日期范围或分区键的任意组合,高效地查询数据湖。两个字符的哈希值可以防止 S3 在高强度写入时出现热点。
2. 结合 Entra ID 联合身份验证的 S3 Access Grants
我们没有为每个用户逐一管理 IAM 策略,而是实现了 S3 Access Grants——这是 AWS 一项相对较新的功能,它可以把身份提供商的声明(claims)直接映射到 S3 前缀级别的权限上。
每一个 Entra ID 用户或组都会被映射到一个特定的 S3 分区。当用户完成认证时,其 JWT 令牌中的 oid 声明会被用来查找其对应的 IAM 角色,随后再通过 sts:AssumeRoleWithWebIdentity 来担任该角色。该角色的权限通过一个 Access Grant 被限定在其所属分区的范围内。
# access_grants.py — 通过 S3 Access Grants 实现按用户的分区隔离
CfnAccessGrant(
self,
f"AccessGrantUser{user_idx}",
access_grants_location_id=user_location.ref,
permission="READWRITE",
grantee=CfnAccessGrant.GranteeProperty(
grantee_identifier=user_role.role_arn,
grantee_type="IAM",
),
)随后,Lambda 集成会在每次请求时,用 Entra ID 令牌换取受限范围的 AWS 凭证,并通过缓存来避免多余的 STS 调用:
# get_cached_or_exchange_credentials.py — 带缓存的令牌交换
def get_cached_or_exchange_credentials(id_token: str) -> CredentialsTypeDef:
key = _cache_key_from_token(id_token)
now = time.time()
with _CACHE_LOCK:
entry = _CACHE.get(key)
if entry and entry["expires_at"] > now + _SKEW_SECONDS:
return entry["credentials"]
new_credentials = _exchange_and_assume_with_expiry(id_token)
with _CACHE_LOCK:
_CACHE[key] = {
"credentials": new_credentials,
"expires_at": new_credentials["Expiration"].timestamp(),
}
return new_credentials这意味着,处于 manufacturing 分区中的用户 A,无法读取或写入用户 B 所在的 quality 分区中的文件。这一限制是在 S3 层面强制执行的,而不仅仅是在应用层面。
3. 多模态 AI 元数据增强
当一个文件落入 S3 时,一条事件驱动的管道便会自动启动。系统会检测文件类型,生成一份量身定制的 Bedrock Data Automation 蓝图,运行提取作业,并将结构化的结果发布回 EventBridge。
整个工作流由 Step Functions 使用 JSONata 表达式进行编排:
# workflow.py — 使用 Bedrock Data Automation 的 Step Functions 编排
definition_body = sfn.DefinitionBody.from_chainable(
extract_event_data
.next(generate_blueprint) # 基于元数据要求的动态 blueprint
.next(start_data_automation_job) # Bedrock Data Automation
.next(wait_for_job_completion) # 使用任务令牌的异步等待
.next(normalize_event_data) # 经过 EventBridge 模式校验的输出
.next(publish_output_event) # EventBridge 完成事件
.next(sfn.Succeed(self, "FileMetadataEnrichmentCompletion"))
)蓝图生成器会根据文件模态以及用户自定义指定的元数据需求,动态创建提取用的模式(schema)。图像会得到边界框检测和分类。文档会得到摘要和要点提取。音频则会得到转录和说话人识别:
# bedrock_data_automation_blueprint_generator.py
def generate_blueprint_schema(enrichments, blueprint_type):
properties = {}
# 基础属性——始终提取
properties["data_classification"] = {
"type": "string",
"instruction": "The data classification level (public, internal, confidential, restricted)",
}
properties["summary"] = {
"type": "string",
"instruction": "A brief summary of the file content",
}
properties["keywords"] = {
"type": "array",
"items": {"type": "string"},
"instruction": "Key terms and keywords extracted from the document",
}
# 特定于模态的增强
if enrichments.get("audio", {}).get("transcribe"):
properties["transcript"] = {
"type": "string",
"instruction": "Full transcript of the audio content",
}
if enrichments.get("video", {}).get("scenes"):
properties["scenes"] = {
"type": "array",
"items": {
"type": "object",
"properties": {
"scene_number": {"type": "number"},
"start_time": {"type": "number"},
"end_time": {"type": "number"},
"description": {"type": "string"},
},
},
"instruction": "Scene changes detected in the video with timestamps",
}
return {"class": f"{blueprint_type.capitalize()}Metadata", "properties": properties}增强后的元数据会作为一个 .metadata.json 附属文件(sidecar),与原始文件一并存储在同一个 Hive 分区位置,从而使其可以立即被查询。
4. 对数 GB 级文件上传的支持
API Gateway 有 10MB 的负载上限。而产品测量文件可能有数 GB 之大。我们通过一套基于预签名 URL 的分段上传流程解决了这一问题——在处理繁重的传输任务时,该流程完全绕开了 API Gateway。
API 会计算出最优的分块大小,发起一次分段上传,并为每个分块返回预签名 URL。客户端则直接向 S3 上传:
# multi_chunk.py — 预签名分段上传编排
def initiate_multi_chunk_upload_presigned(body):
chunk_size, num_chunks = _calculate_chunk_size(request.fileSize)
response = s3.create_multipart_upload(
Bucket=bucket, Key=key, ServerSideEncryption="aws:kms"
)
presigned_urls = []
for chunk_number in range(1, num_chunks + 1):
presigned_url = s3.generate_presigned_url(
"upload_part",
Params={
"Bucket": bucket, "Key": key,
"UploadId": response["UploadId"],
"PartNumber": chunk_number,
},
ExpiresIn=expires_in,
)
presigned_urls.append({
"chunkNumber": chunk_number,
"url": presigned_url,
"startByte": (chunk_number - 1) * chunk_size,
"endByte": min(chunk_number * chunk_size - 1, request.fileSize - 1),
})
return {"uploadId": response["UploadId"], "chunks": presigned_urls}在前端,Web 应用会透明地处理这一切;小文件走 API,大文件则会自动切换为分段上传:
// api.service.ts — 自动上传策略选择
async upload(file: File, request: InitiateUploadRequest, onProgress?) {
if (file.size > FILE_SIZE_THRESHOLD) {
return this.multiChunkUploadService.uploadFile(file, request, { onProgress });
}
const base64Content = await readFileAsBase64(file);
return this.httpService.request('/files', 'POST', {
body: { content: base64Content, fileName: request.fileName },
});
}5. Web 应用
前端是一个基于 AWS Amplify 的 Storage Browser 组件构建的 React 应用,我们对其进行了定制,使各项操作经由我们的 API 路由,而非直接访问 S3。这让我们能够完全掌控访问控制、元数据操作和上传策略,同时提供一种精致而熟悉的文件管理体验。
// storage-browser.provider.tsx — 自定义存储浏览器,带有基于 API 的操作
const { StorageBrowser } = createStorageBrowser({
config: {
registerAuthListener: async (onAuthStateChange) => {
const authService = getAuthService();
authService.registerAuthListener(onAuthStateChange);
},
listLocations: async ({ options }) => {
return await apiService.getLocations({
options: { pageSize: 30, nextToken: options?.nextToken },
});
},
},
actions: actionsBuilder.buildActions(),
});用户使用其西门子能源的 Entra ID 凭证登录后,立即就只能看到自己有权访问的分区。他们可以上传任意大小的文件、浏览 Hive 分区的文件夹结构、内嵌预览文档和图像、通过预签名 URL 下载,以及编辑元数据——所有这些都无需离开浏览器。
6. 跨账户的事件驱动架构
文件管理器与元数据增强微服务运行在各自独立的 AWS 账户中。当一个文件被上传时,文件管理器的 S3 存储桶会触发一个 Lambda,由它向一个跨账户的 EventBridge 总线发布一个 FileMetadataEnrichmentRequest 事件。
# file_metadata_enrichment_processor.py — 跨账户事件发布
def put_event(bucket, key, size=None, etag=None, enrichments=None):
detail = create_event_detail(bucket, key, size, etag, enrichments)
return events_client.put_events(Entries=[{
"Source": "com.siemens-energy.pmdp.file-metadata-enrichment",
"DetailType": "FileMetadataEnrichmentRequest",
"Detail": json.dumps(detail),
"EventBusName": EVENT_BUS_ARN, # 跨账户 ARN
}])在增强这一侧,EventBridge 规则会将事件经由 SQS(并配备一个 DLQ 以增强弹性)路由进入一个 EventBridge Pipe;该 Pipe 会先依据一个模式注册表(schema registry)校验事件,然后才调用 Step Functions 工作流。完成事件和异常事件会被转发回最初发起的账户。
这种解耦的架构意味着,这套元数据增强微服务可以被西门子能源的任何团队复用。他们只需向总线发布事件即可。
压力测试:在规模化场景下加以验证
该系统必须处理巨量的数据进出。他们的某些文件仓库仅仅是删除就可能耗时数天。我们构建了一个专门的压力测试 CLI,它可以生成大小可配置的文件(从 100KB 到 5GB 以上),对它们进行并发上传,通过预签名 URL 下载,并用 MD5 校验和来验证其完整性。
测试运行验证了以下方面:
- 并发上传——同时上传 20 个以上的文件,其中包括通过分段方式上传的数 GB 级负载
- 100% 的成功率——在单次测试运行中横跨数百个文件均是如此
- 下载验证——在数据往返经过整条管道之后,确认其逐字节的完整性
- 自动清理——从 S3 和本地存储中同时清理测试产物
基础设施即代码:一切尽在 CDK 中
整个平台——两个 AWS 账户、所有服务、所有 IAM 角色以及所有事件路由——都以 Python CDK 定义。特定环境的配置通过 Hydra/OmegaConf 进行管理,从而使得启用一个新环境或接入一个新团队都变得轻而易举。
# config.py — 类型安全的、特定于环境的配置
config_environment = make_config(
env=zf(cdk.Environment),
environment=zf(Literal["sandbox", "dev", "uat", "prd"]),
s3explorer=zf(config_s3explorer),
access_grants=zf(config_access_grants),
project_name=zf(str, default="mfg-product-measurement-data-pipeline"),
file_metadata_enrichment_event_bus_arn=zf(str),
)Web 应用的 CI/CD 使用 GitLab OIDC 联合身份验证,无需管理任何长期有效的凭证或需要轮换的密钥。该 CDK 堆栈在单个构件中即可配置好 OIDC 提供商、部署角色以及 Amplify 源存储桶。
成果
- 文件大小支持:不设上限(已测试至 5GB 以上)
- 上传并发数:20 个以上的同时上传
- 元数据提取:对文档、图像、音频和视频均自动进行
- 访问控制:通过 S3 Access Grants 实现按用户、按组的分区隔离
- 环境:由单个 CDK 代码库支撑的 4 个环境(沙盒、开发、UAT、生产)
- 身份集成:结合 OIDC 联合身份验证的 Microsoft Entra ID SSO
- 基础设施:100% 基础设施即代码(Python CDK)
技术栈
- 计算: AWS Lambda(Python 3.14,ARM64)
- 编排: AWS Step Functions(JSONata)
- AI/ML: Amazon Bedrock Data Automation
- 存储: Amazon S3(智能分层、KMS 加密、传输加速)
- API: Amazon API Gateway(REST),搭配 Lambda Powertools + Swagger
- 事件: Amazon EventBridge、EventBridge Pipes、SQS
- 身份: Microsoft Entra ID、OIDC 联合身份验证、S3 Access Grants
- 前端: React、Vite、AWS Amplify、Amplify UI Storage Browser
- IaC: AWS CDK(Python)、Hydra/OmegaConf
- CI/CD: GitLab CI
结论
这个项目需要贯穿整个 AWS 技术栈的深厚专业能力——从底层的 IAM 策略设计和跨账户事件路由,到前沿的 Bedrock Data Automation 集成以及 Amplify Storage Browser 的定制。
我们竭尽所能地超越了最初的规格要求,并推荐了 AWS 云创新中最新、最出色的技术。最终的成果是这样一套系统设计:上传一个文件便会触发一条 AI 管道,为其增添结构化的元数据,将其存储在一个可查询的分区中,并使其立即可以通过 Web 界面浏览——而这一切,用户除了拖放操作之外无需做任何事情。
许多组织都需要在 AWS 上构建稳健、现代的云应用——那种能够应对真实世界规模、与企业身份提供商相集成、并在关键之处充分利用 AI 的系统。你也是其中之一吗? 我们聊聊吧。

