← Financial Cloud Cloud 云端社群 · Amarathon 2025 回顾

极点宏观|Financial Cloud Cloud · AWS Amarathon 2025

极点宏观|Amarathon 2025 回顾 33:构建流式 Iceberg 表以进行实时物流分析

演讲者: Fahad Shah

场次: 33

场次
Summit Dev Lounge2026 Re:cap
01 把架构写成 Steering,引导 AI Agent
Summit Dev Lounge2026 Re:cap
02 Agent Harness 才是真正的工程护城河
Summit Dev Lounge2026 Re:cap
03 用白话询问可观测性数据
Summit Dev Lounge2026 Re:cap
04 以 Bedrock AgentCore 打造 Serverless AR 游戏
Summit Dev Lounge2026 Re:cap
05 AgentCore 上的多 Agent 量化回测
Summit Dev Lounge2026 Re:cap
06 三分钟用 Kiro 把博客变成幻灯片
Summit Dev Lounge2026 Re:cap
AWS Community Day Hong Kong 2025 Re:cap
02 使用 Terraform 实现 AWS 合规
AWS Community Day Hong Kong 2025 Re:cap
03 从初学者到构建者:一段精彩的 AWS 云之旅
AWS Community Day Hong Kong 2025 Re:cap
04 以团队为先:使用 Laravel 与 Bref 开展无服务器工程
AWS Community Day Hong Kong 2025 Re:cap
05 活动开幕式
AWS Community Day Hong Kong 2025 Re:cap
06 智能体到智能体:在 AWS 上构建可互操作的 AI
AWS Community Day Hong Kong 2025 Re:cap
07 利用另一类遥测数据,借助 AI 智能体更快改进
AWS Community Day Hong Kong 2025 Re:cap
08 告别氛围编程:使用 Kiro 进行规格驱动开发
AWS Community Day Hong Kong 2025 Re:cap
09 使用 MCP 与 AI 智能体进行自动化测试
AWS Community Day Hong Kong 2025 Re:cap
10 使用机器学习方法实现电信安全现代化
AWS Community Day Hong Kong 2025 Re:cap
11 重新思考生成式 AI 智能体:RAG 与 MCP
AWS Community Day Hong Kong 2025 Re:cap
12 使用 TAK 和 AWS 开展灾难与应急响应
AWS Community Day Hong Kong 2025 Re:cap
13 从测试视角重新思考无服务器应用程序工作流
AWS Community Day Hong Kong 2025 Re:cap
14 Practical AWS FinOps for Cloud Success
AWS Community Day Hong Kong 2025 Re:cap
15 基于 AWS 的 AI 驱动全球纯 Alpha 宏观交易:重塑风险调整后资产收益
AWS Community Day Hong Kong 2025 Re:cap
FSI Recap
01 现代交易生命周期:从交易到结算
FSI Recap
02 Goldman Sachs:通过 Fast Track 加速应用程序上云 - AWS Re:cap Q1/2023
FSI Recap
03 Zurich Insurance Group:在 AWS 上构建高效的日志管理解决方案
FSI Recap
04 FSI Meetup 2025 年第四季度 - Brex 数据库灾难恢复
FSI Recap
05 FSI Meetup 2025 Q4 - Graviton 迁移成功案例
FSI Recap
06 FSI Meetup 2025 Q4 - Stifel 现代数据平台
FSI Recap
07 FSI Meetup 2025 Q4 - PayPal 金融交易数据对账系统
FSI Recap
08 FSI Meetup 2025 Q4 - 规模化提升韧性
FSI Recap
09 最大限度提高 AI 推理成本效益:战略性采用 AWS GPU 实例
FSI Recap
10 高级智能体 AI 设计模式
FSI Recap
11 在 AWS 上构建全新的现代化应用
FSI Recap
AWS re:Invent 2025
01 Coinbase re:Invent 回顾 (IND3312)
AWS re:Invent 2025
02 利用 AI 和 AWS 构建未来交易平台
AWS re:Invent 2025
03 交易创新:Jefferies 基于 Amazon Bedrock 构建的 AI 助手 (IND3315)
AWS re:Invent 2025
04 FSI 如何通过 Agentic AI (GBL302) 彻底改变 HFT 分析
AWS re:Invent 2025
05 使用 Amazon Time Sync 改进分布式系统(采用 Nasdaq)
AWS re:Invent 2025
06 Amazon Aurora HA 和 DR 全球弹性设计模式 (DAT442)
AWS re:Invent 2025
07 构建智能体式 AI:Amazon Nova Act 与 Strands Agents 实践 (DEV327)
AWS re:Invent 2025
08 深入探讨 Amazon Aurora 及其创新 (DAT441)
AWS re:Invent 2025
09 深入探讨 Amazon S3(STG407)
AWS re:Invent 2025
10 Nasdaq:为全球金融服务构建弹性基础设施 (HMC327)
AWS re:Invent 2025
11 AWS Lambda 新功能 (CNS376)
AWS re:Invent 2025
12 使用 Kiro 进行规范驱动开发 (DEV314)
AWS re:Invent 2025
13 Amazon 的 FinOps:全球电商巨头的云成本管理经验 (AMZ308)
AWS re:Invent 2025
14 AWS 上交易平台的 Tick-to-Trade 延迟
AWS re:Invent 2025
政务数据
01 The AI Era: The Boundary Between Development and Design Is Disappearing
政务数据
02 端侧多模态 AI 与智慧城市实践
政务数据
03 大模型能力评测与 AI 项目落地方法论
政务数据
04 基于云代理的政府开发全链路受控自动化
政务数据
05 从多智能体看 Agent 时代软件新生态
政务数据
06 AI 驱动的宏观量化研究与智慧治理
政务数据
07 公共数据授权运营与智慧政务实践
政务数据
08 数据资产化落地实践:确权合规、工程治理与数字政府案例
政务数据
09 AI技术赋能心理健康公益:可信平台的治理、架构与实践
政务数据
Amarathon 2025 回顾
01 开发者的智能体架构设计路线图
Donnie Prakoso
02 Amazon Bedrock 数据自动化
Hafiz Syed Ashir Hassan
03 AgentCore 上的多智能体
Tan Xin
04 实践中构建智能体式 AI:Nova Act 与 Strands Agents
Haowen Huang
04 使用规格驱动开发,通过 Kiro 加速迁移项目
Sanchit Dilip Jain
06 从「匹配」到「理解」:由 AgentCore Memory 驱动的个性化 AI 搜索实践
Liu Cao
07 从观察到优化:从 LLM 可观测性迈向 AIOps,将实时洞察转化为智慧自动化
Jimmy Soh
08 部署 TEAM 并打造最佳工程团队
Yuji Oshima
09 五年来所谓无服务器数据库带来的五个惨痛教训
Renato Losio
14 如果 AI 替我工作会怎样:Q Developer CLI 与 Kiro 如何改变我的日常工作
Miguel Angel Muñoz
16 兼顾速度与警觉:Amazon Bedrock Agent 开发的安全要点
Brian Tarbox
26 在单张 H100 上运行 OSS LLM:更智能、更便宜、更快速
Adit Modi Adit Modi
28 现代统一元数据架构:打破数据孤岛的新方法
Shaofeng Shi
29 无服务器 MediaOps:使用 Amazon Web Services 上的 AI 自动化视频工作流
Luis Valdivia
30 通过大规模性能测试构建兼具效率与可靠性的架构
Luis Guirigay
31 通过开源连接世界:技术、社区与全球开发者关系的实践历程
Richard Lin
33 构建流式 Iceberg 表以进行实时物流分析
Fahad Shah
34 加速大规模机器人策略训练:基于 Kiro、Trainium 和 EKS 的自动化闭环架构
Junjie Tang
35 通过规格驱动开发,从 Vibe 走向可行方案
Ricardo Sueiras
36 让云成本分析更智能:使用 Strands 和 AgentCore 构建 FinOps 智能体
Xiaofei Li
37 使用 CNCF Kagent、K8sGPT 和 Nova Sonic 转型 K8s 对话式智能体 AIOps
Shaoyi Li

现代物流面临的挑战

● 管理卡车、司机、路线、燃料、维护、货运和仓库等多个数据流。

● 需要实时运营视图和长期分析。

数据存储要求

● 提供最新的连接视图,以支持即时运营。

● 使用 Apache Iceberg 进行长期分析。

技术栈

● RisingWave:提供流处理能力的数据平台。

● Lakekeeper:用于数据管理的开放式 REST catalog。

● Kafka:流数据的事件骨干。

● 对象存储(例如 MinIO):数据存储解决方案。

目标

● 演示如何使用指定的开放式技术栈构建流式 Iceberg 表。

● 为现代物流数据管理提供简单有效的解决方案。

物流分析问题

● 当今的物流平台会生成:

● 卡车:车队清单和位置

● 司机:名册和任务分配

● 货运:始发地、目的地和重量

● 仓库:容量和站点

● 路线:预计到达时间(ETA)和距离

● 燃料与维护:成本和可靠性信号

● 挑战:

● 运营团队需要涵盖所有这些数据流的最新连接视图。

● 数据团队需要将同一份数据存入 Iceberg,以进行 BI、AI 和历史分析。

我们将构建的内容(流式 Iceberg 模式)

● Kafka 将七个物流主题送入 RisingWave。

● 使用 SQL 表达多路流连接,并在 RisingWave 内持续进行物化。

● 结果从 RisingWave 持久化为原生 Apache Iceberg 表,并存储在 MinIO 等兼容 S3 的对象存储中。

● Spark、Trino 和 DuckDB 等引擎通过开放式 REST catalog 查询相同的 Iceberg 表。

为何使用 RisingWave 构建流式 Iceberg 表?

[ 1 ] 批处理优先的工作流:

● 周期性作业、过时的连接结果和繁重的数据流水线。

● 需要使用单独的 ETL 工具写入 Iceberg。

[ 2 ] RisingWave + 流式 Iceberg 表:

● 在 RisingWave MV 中持续更新连接和聚合结果。

● 始终“近乎实时”的 Iceberg 快照。

● 一条 RisingWave 数据流水线同时支持实时仪表盘和离线分析。

● 目标:由 RisingWave 负责流数据流水线和 Iceberg 写入,让 Iceberg 的使用体验如同数据库。


高层架构

● 我们的端到端技术栈:

● Kafka — 7 个物流主题的事件骨干。

● RisingWave(流数据库)— 使用 SQL 摄取、连接和聚合数据;管理物化视图。

● RisingWave Iceberg Table Engine + Lakekeeper — Iceberg 表的开放式 REST catalog。

● MinIO — 兼容 S3 的对象存储。

● 模式:Kafka → RisingWave → MinIO 中的 Iceberg → 任何引擎均可通过 REST catalog 查询。

RisingWave 中的物流数据流与多路流连接

● RisingWave 中的七个物流数据流

● 本示例使用七个 Kafka 主题,这些主题会成为 RisingWave 中的 source:

● trucks — 车队清单、容量、当前位置。

● driver — 司机详细信息和 assigned_truck_id。

● shipments — 始发地、目的地、重量、卡车绑定关系。

● warehouses — 仓库位置和容量。

● route — route_id、truck_id、driver_id、ETD/ETA、distance_km。

● fuel — 加油事件(时间、升数、加油站)。

● maint — 维护历史和成本。

● RisingWave 将每一个数据流都视为流式表,可直接使用简单的 PostgreSQL 风格 SQL 进行连接。

模式 1:RisingWave 中的多路流连接

● 在 RisingWave 中,我们将核心物流逻辑表达为一个多路流连接。

● LEFT JOIN drivers → trucks,以保留未匹配的司机数据。

● JOIN shipments,以附加工作负载和目的地。

● JOIN warehouses,以加入容量和位置信息。

● JOIN route,以获取 ETD/ETA 和距离。

● JOIN fuel 和 maint,以获取成本和可靠性信号。

● 这会生成 logistics_joined_mv — RisingWave 内针对每辆卡车/每位司机/每条路线持续更新的反规范化物流记录。


车队 KPI、原生 Iceberg 表与跨引擎读取 模式 2

RisingWave 中的车队 KPI 视图

● 在连接后的 MV 之上,我们另外定义一个用于车队 KPI 的 RisingWave MV:

● 每辆卡车的容量利用率(%)。

● 每辆卡车的燃料总成本和维护成本。

● 合计运营总成本。

● 当前路线信息(ID、ETD、ETA、distance_km)。

● 相关司机详细信息。

RisingWave 中的 overview 会成为实时车队绩效表 供 Grafana 和运营仪表盘使用。模式 3:从 RisingWave 流式写入原生 Iceberg

● 不使用自定义 writer 服务,而是:

● [ 1 ] 将 logistics_joined_iceberg 定义为由 RisingWave 管理的原生 Iceberg 表。

● [ 2 ] 其 schema 与 logistics_joined_mv 一致。

● [ 3 ] RisingWave 中的一小段配置可控制将流式变更提交为 Iceberg 快照的频率。

模式 4:通过 REST catalog 进行跨引擎读取

● 使用由 RisingWave 创建并注册到 Lakekeeper REST catalog 的 Iceberg 表:

● [ 1 ] Spark 将 lakekeeper 挂载为 catalog

● [ 2 ] Trino / DuckDB / Dremio 可以使用各自的 Iceberg connector 读取同一张表。

[ 3 ] 所有引擎都能看到 RisingWave 持续更新的同一份 Iceberg 数据。

● 无需副本,无需专有表格式 — 只有由 RisingWave 写入的标准 Iceberg。


从本地笔记本电脑到生产集群:部署选项

● 部署选项:从笔记本电脑到集群

[ 1 ] 本地(用于学习和原型开发):

● 使用 Docker 运行 RisingWave、Kafka、MinIO 和 Lakekeeper。

● 非常适合在笔记本电脑上试验流连接和 Iceberg 表。

[ 2 ] 生产环境(用于实际工作负载):

● 通过 Kubernetes + Helm 部署 RisingWave 和技术栈的其余部分。

● 使用适合你的环境的存储类、资源限制和持久化配置。

● 在 RisingWave 中使用相同的 SQL 和模式 只不过更加持久、可扩展且自动化程度更高。

简化传统 Iceberg 技术栈

● 传统 Iceberg 部署通常需要:

● 单独的流处理引擎。

● 独立运行的 Iceberg writer 作业。

● 外部压缩和维护工作流。

● 额外的衔接机制,以确保 catalog、writer 和存储保持一致。

● 使用 RisingWave:

● [ 1 ] 流数据库负责数据摄取、连接、物化视图和 Iceberg 写入。

● [ 2 ] REST catalog + MinIO 让一切保持完全开放且可互操作。

● 更少的活动组件、更低的运维负担。

使用 RisingWave 的参考架构

● 可以将系统视为以 RisingWave 为中心的三个层级:

[ 1 ] 数据流 → RisingWave 表。

● Kafka 主题会成为 RisingWave 中的流式表。

[ 2 ] 表 → RisingWave 物化视图。

● 流连接和聚合会成为实时 MV(logistics_joined_mv、truck_fleet_overview)。

[ 3 ] 视图 → 流式 Iceberg 表。

● RisingWave 只需少量配置和一条 INSERT....SELECT,即可将 MV 转换为流式 Iceberg 表。

● 一旦将 RisingWave 视为“流式 SQL + Iceberg 引擎”,就能在许多领域复用这一模型。

物流以外的可复用模式

● RisingWave + Iceberg 模式适用于:

● 电子商务:订单、库存、定价、客户事件。

● FinTech:交易、余额、风险信号。

● 工业 IoT:机器、传感器、警报、维护。

● 电信:会话、使用量、QoS 指标。

● 只要有多个实时数据流,并且需要开放式长期存储,就可以用同样的方式使用 RisingWave MV 和 Iceberg 表。

要点总结(RisingWave + Iceberg)

● 结合 Kafka、RisingWave、REST catalog、MinIO 和 Iceberg 的参考架构。

● 实用模式:多路流连接、KPI 视图,以及从 RisingWave 原生写入 Iceberg。

● 无需自定义 writer、临时拼凑的压缩作业或严重的供应商锁定,即可获得实时物流分析。