Lambda 架构由批处理层、加速层和服务层构成:批处理层用 MapReduce/Spark 处理全量历史数据(分钟/小时级延迟),加速层用 Flink/Storm 处理增量数据(秒/毫秒级延迟),服务层合并两者结果。其核心问题是维护两套代码,开发成本高且数据口径易不一致。
Kappa 架构删除了批处理层,仅保留流处理系统,通过消息队列重播实现历史数据处理。代码统一、便于维护,但消息中间件存在性能瓶颈,历史数据处理能力较弱。
湖仓一体则将数据湖的灵活性与数据仓库的可靠性融合——通过 Iceberg、Paimon 等开放表格式,一套存储同时支持流式写入、批量分析和 ACID 事务。这是当前架构演进的主流方向。
无论采用哪种架构风格,数据分层是数仓设计的基石。以电商实时数仓为例:
分层 | 职责 | 技术选型 |
|---|---|---|
ODS | 多源异构数据实时接入 | Kafka + Flink CDC |
DWD | 实时 ETL 清洗、维度关联 | Flink SQL |
DWS | 预聚合、宽表构建 | Flink + Paimon |
ADS | 多维分析 / 点查服务 | Doris / StarRocks |
ODS 层通过 Kafka 实现多源数据接入,支持 JSON/Avro/Protobuf 格式解析。DWD 层基于 Flink SQL 实现实时 ETL。DWS 层通过预聚合降低下游压力。
-- 创建 Kafka 源表(ODS 层)
CREATE TABLE ods_user_behavior (
user_id BIGINT,
item_id BIGINT,
behavior STRING,
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'raw_events',
'properties.bootstrap.servers' = 'kafka:9092',
'properties.group.id' = 'flink_ods_consumer',
'scan.startup.mode' = 'latest-offset',
'format' = 'json'
);WATERMARK 用于处理乱序数据,设置 5 秒延迟阈值。scan.startup.mode 控制消费起始位置。
-- 创建 DWD 层明细表(写入 Paimon)
CREATE TABLE dwd_user_behavior (
user_id BIGINT,
item_id BIGINT,
behavior STRING,
event_time TIMESTAMP(3),
user_region STRING,
item_category STRING,
PRIMARY KEY (user_id, event_time) NOT ENFORCED
) WITH (
'connector' = 'paimon',
'path' = 'hdfs://warehouse/paimon/dwd_user_behavior',
'bucket' = '8',
'merge-engine' = 'aggregation'
);
-- 实时 ETL:清洗 + 维表关联
INSERT INTO dwd_user_behavior
SELECT
o.user_id,
o.item_id,
o.behavior,
o.event_time,
u.region AS user_region,
i.category AS item_category
FROM ods_user_behavior o
LEFT JOIN dim_user u ON o.user_id = u.user_id
LEFT JOIN dim_item i ON o.item_id = i.item_id
WHERE o.behavior IS NOT NULL;Paimon 表支持 aggregation 合并引擎,可实现 PartialUpdate 等增量聚合。
-- 创建 DWS 层聚合表
CREATE TABLE dws_hourly_agg (
hour_start TIMESTAMP(3),
item_id BIGINT,
pv BIGINT,
uv BIGINT,
PRIMARY KEY (hour_start, item_id) NOT ENFORCED
) WITH (
'connector' = 'paimon',
'path' = 'hdfs://warehouse/paimon/dws_hourly_agg',
'bucket' = '16',
'merge-engine' = 'aggregation'
);
-- 小时级聚合(Flink 流式写入)
INSERT INTO dws_hourly_agg
SELECT
DATE_FORMAT(event_time, 'yyyy-MM-dd HH:00:00') AS hour_start,
item_id,
COUNT(*) AS pv,
COUNT(DISTINCT user_id) AS uv
FROM dwd_user_behavior
GROUP BY DATE_FORMAT(event_time, 'yyyy-MM-dd HH:00:00'), item_id;# Flink Checkpoint 配置(保障 Exactly-Once)
flink.state.backend: rocksdb
flink.checkpoints.interval: 60s
flink.checkpoints.dir: hdfs://checkpoint_path
flink.state.backend.rocksdb.localdir: /data/rocksdbCheckpoint 机制保障 Exactly-Once 语义。RocksDB 作为状态后端可支撑大状态场景。
实时任务与离线任务混部会导致资源抢占引发延迟抖动。推荐通过 K8s 实现弹性资源调度,根据实时负载自动调整 TaskManager 数量。
大数据架构师的核⼼价值在于在能力上限与系统开销之间找到最优平衡点。从 Lambda 的双系统维护,到 Kappa 的流式统一,再到湖仓一体的存储计算分离,每一次演进都在降低复杂度、提升时效性。
本文提供的电商实时数仓代码(Flink SQL + Paimon + Kafka)可直接在本地环境运行验证。实践中,建议从单 Flink 作业 + Paimon 存储起步建立性能基线,再根据业务需求逐步引入 Doris 加速层或多源 CDC 采集。架构设计的本质不是堆砌组件,而是用最简的链路解决最核心的业务问题。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。