在数据实时化需求日益增长的背景下,传统 T+1 离线数仓已难以满足秒级分析诉求。实时数仓的构建不仅是流式计算引擎的选型,更涉及端到端的低延迟、高吞吐、可容错系统工程。本文聚焦于开源技术栈(Apache Flink、ClickHouse、Kafka),系统阐述实时数仓的分层模型(ODS → DWD → DWS → ADS)、各层核心实现代码、关键优化策略以及容错机制。文中所有代码均经过简化验证,可直接作为工程参考,但不包含特定业务场景的案例分析,旨在为读者提供一套可复用的技术框架和调优思路。
组件 | 角色 | 核心优势 |
|---|---|---|
Apache Flink | 实时计算引擎 | 精确一次语义、状态后端、事件时间处理、丰富的窗口 API、支持 SQL 和 DataStream 混合编程 |
ClickHouse | 实时 OLAP 数据库 | 列式存储、向量化执行、主键稀疏索引、极速聚合查询,适合高并发点查与大规模聚合 |
Apache Kafka | 消息中间件 | 高吞吐、持久化、回溯消费,配合 Flink Checkpoint 实现端到端一致性 |
实时数仓的分层模型(ODS → DWD → DWS → ADS)在流式场景下依然成立,但每一层必须具备可回放和可修正能力。Flink 的 Checkpoint + Kafka 的 offset 管理天然满足此要求。
业务日志 → Nginx → Filebeat → Kafka (ODS)
↓
Flink SQL (ETL + 维表 JOIN)
↓
Kafka (DWD 明细流)
↓
Flink SQL (窗口聚合 / CEP / 多流关联)
↓
Kafka (DWS 汇总流)
↓
Flink Sink → ClickHouse (ADS)
↓
业务查询 / BI 看板所有 Flink 作业均开启 Checkpoint(间隔 60s,exactly-once),状态后端使用 RocksDB 以支持大状态(TB 级)。Checkpoint 目录建议配置为高可用 HDFS 或 S3。
假设 ODS 层 Kafka 主题 ods_order 包含 JSON 格式订单数据,需解析、过滤脏数据,并关联商品维表(MySQL 缓存于 Redis,此处采用 Flink 异步 I/O + 本地缓存)。
// Flink 1.16 Java API
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(4);
env.enableCheckpointing(60000, CheckpointingMode.EXACTLY_ONCE);
// 1. 从 Kafka 消费 ODS
KafkaSource<String> source = KafkaSource.<String>builder()
.setBootstrapServers("kafka:9092")
.setTopics("ods_order")
.setGroupId("flink_dwd_group")
.setStartingOffsets(OffsetsInitializer.latest())
.setValueOnlyDeserializer(new SimpleStringSchema())
.build();
DataStream<String> odsStream = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source");
// 2. 解析 JSON 并过滤(使用 Jackson)
SingleOutputStreamOperator<OrderEvent> orderStream = odsStream
.map(new MapFunction<String, OrderEvent>() {
@Override
public OrderEvent map(String value) throws Exception {
ObjectMapper mapper = new ObjectMapper();
JsonNode node = mapper.readTree(value);
if (node.get("order_id") == null || node.get("user_id") == null) {
return null; // 脏数据
}
return new OrderEvent(
node.get("order_id").asLong(),
node.get("user_id").asLong(),
node.get("product_id").asLong(),
node.get("amount").asDouble(),
node.get("create_time").asLong()
);
}
})
.filter(Objects::nonNull)
.assignTimestampsAndWatermarks(
WatermarkStrategy.<OrderEvent>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((event, timestamp) -> event.getCreateTime())
);
// 3. 异步关联维表(商品信息缓存)
DataStream<EnrichedOrder> enrichedStream = AsyncDataStream.unorderedWait(
orderStream,
new AsyncDimJoinFunction(), // 实现 RichAsyncFunction,从 Redis/MySQL 查询
1000, TimeUnit.MILLISECONDS,
100 // 最大并发请求数
);
// 4. 写出到 Kafka DWD
KafkaSink<String> dwdSink = KafkaSink.<String>builder()
.setBootstrapServers("kafka:9092")
.setRecordSerializer(KafkaRecordSerializationSchema.builder()
.setTopic("dwd_order_enriched")
.setValueSerializationSchema(new SimpleStringSchema())
.build()
)
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.build();
enrichedStream.map(EnrichedOrder::toJson).sinkTo(dwdSink);关键设计要点:
assignTimestampsAndWatermarks 容忍 5 秒乱序,基于事件时间。从 dwd_order_enriched 读取流,按 1 分钟滚动窗口计算每个商品 ID 的订单金额总和、订单数,并输出到 Kafka dws_product_minute。
-- 使用 Flink SQL(生产环境可混合 Table API 做更复杂的逻辑)
CREATE TABLE dwd_order (
order_id BIGINT,
user_id BIGINT,
product_id BIGINT,
amount DOUBLE,
create_time TIMESTAMP(3),
product_name STRING,
category STRING,
WATERMARK FOR create_time AS create_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'dwd_order_enriched',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'json'
);
CREATE TABLE dws_product_minute (
product_id BIGINT,
window_start TIMESTAMP(3),
window_end TIMESTAMP(3),
total_amount DOUBLE,
order_count BIGINT,
PRIMARY KEY (product_id, window_start) NOT ENFORCED
) WITH (
'connector' = 'kafka',
'topic' = 'dws_product_minute',
'format' = 'json'
);
-- 开启 MiniBatch 优化(Flink 1.15+ 默认开启,但可显式设置)
SET table.exec.mini-batch.enabled = true;
SET table.exec.mini-batch.size = 5000;
SET table.exec.mini-batch.allow-latency = '5s';
INSERT INTO dws_product_minute
SELECT
product_id,
TUMBLE_START(create_time, INTERVAL '1' MINUTE) AS window_start,
TUMBLE_END(create_time, INTERVAL '1' MINUTE) AS window_end,
SUM(amount) AS total_amount,
COUNT(*) AS order_count
FROM dwd_order
GROUP BY product_id, TUMBLE(create_time, INTERVAL '1' MINUTE);优化措施:
state.ttl(如 10 分钟),避免无限膨胀。optimizer.aggregate-phase-strategy 控制。ClickHouse 建表(使用 ReplacingMergeTree 以支持去重,因为窗口结果可能因延迟数据而重复发送):
CREATE TABLE dws_product_minute_local ON CLUSTER clickhouse_cluster
(
product_id UInt64,
window_start DateTime,
window_end DateTime,
total_amount Float64,
order_count UInt64,
update_time DateTime DEFAULT now()
)
ENGINE = ReplacingMergeTree(update_time)
PARTITION BY toYYYYMMDD(window_start)
ORDER BY (product_id, window_start)
SETTINGS index_granularity = 8192;
-- 分布式表(查询用)
CREATE TABLE dws_product_minute_dist ON CLUSTER clickhouse_cluster
AS dws_product_minute_local
ENGINE = Distributed(clickhouse_cluster, default, dws_product_minute_local, rand());Flink Sink 实现(批量 Buffer + JDBC):
public class ClickHouseBatchSink extends RichSinkFunction<DwsResult> {
private transient ClickHouseConnection conn;
private transient PreparedStatement stmt;
private final List<DwsResult> buffer = new ArrayList<>();
private static final int BATCH_SIZE = 10000;
private static final int FLUSH_INTERVAL_MS = 5000;
@Override
public void open(Configuration parameters) throws Exception {
Class.forName("com.clickhouse.jdbc.ClickHouseDriver");
String url = "jdbc:clickhouse://clickhouse-host:8123/default?socket_timeout=120000&connect_timeout=10000&async_insert=1&wait_for_async_insert=0";
conn = DriverManager.getConnection(url);
String sql = "INSERT INTO dws_product_minute_dist (product_id, window_start, window_end, total_amount, order_count) VALUES (?, ?, ?, ?, ?)";
stmt = conn.prepareStatement(sql);
}
@Override
public void invoke(DwsResult value, Context context) throws Exception {
buffer.add(value);
if (buffer.size() >= BATCH_SIZE) {
flush();
}
}
private void flush() throws SQLException {
if (buffer.isEmpty()) return;
for (DwsResult r : buffer) {
stmt.setLong(1, r.productId);
stmt.setTimestamp(2, new Timestamp(r.windowStart));
stmt.setTimestamp(3, new Timestamp(r.windowEnd));
stmt.setDouble(4, r.totalAmount);
stmt.setLong(5, r.orderCount);
stmt.addBatch();
}
stmt.executeBatch();
buffer.clear();
}
@Override
public void close() throws Exception {
flush();
if (stmt != null) stmt.close();
if (conn != null) conn.close();
}
}写入优化关键点:
wait_for_async_insert=0 可异步返回)。TwoPhaseCommitSinkFunction 或使用 Flink 的 ExactlyOnce 语义,配合死信队列(DLQ)记录写入失败的数据,定时补偿。业务常见查询如“近 5 分钟热门商品 TOP 10”,直接在分布式表上做 GROUP BY 可能较慢。推荐使用 物化视图 + 投影预先聚合:
-- 创建物化视图,按分钟预聚合
CREATE MATERIALIZED VIEW mv_product_minute_agg
ENGINE = SummingMergeTree()
PARTITION BY toYYYYMMDD(window_start)
ORDER BY (product_id, window_start)
AS SELECT
product_id,
window_start,
SUM(total_amount) AS sum_amount,
SUM(order_count) AS sum_count
FROM dws_product_minute_dist
GROUP BY product_id, window_start;
-- 添加跳数索引加速过滤
ALTER TABLE mv_product_minute_agg ADD INDEX idx_amount sum_amount TYPE minmax GRANULARITY 2;查询示例(要求 500ms 内返回):
SELECT
product_id,
sum(sum_amount) AS total,
argMax(sum_count, window_start) AS last_count
FROM mv_product_minute_agg
WHERE window_start > now() - INTERVAL 5 MINUTE
GROUP BY product_id
ORDER BY total DESC
LIMIT 10
SETTINGS max_threads = 8, optimize_aggregation_in_order = 1;查询优化技巧:
optimize_aggregation_in_order 利用主键顺序加速。ORDER BY,将过滤字段放在最前。PROJECTION(ClickHouse 原生支持)自动选择最优聚合方式。场景 | 解决方案 |
|---|---|
Flink 作业重启 | 从最近完成的 Checkpoint 恢复,Kafka 消费位点自动回滚,保证 exactly-once(需 Sink 支持两阶段提交) |
ClickHouse 写入失败 | 使用死信队列记录失败数据,定时重试;Flink Sink 实现 CheckpointedFunction 保存未确认的批次 |
延迟数据(迟到超过 5 分钟) | Flink 侧设置 allowedLateness(如 1 分钟)并输出到侧输出流,独立修正;DWS 层使用 ReplacingMergeTree 最终去重 |
维表变更(商品分类调整) | 使用 CDC(Debezium)监控 MySQL binlog,广播变更到所有 TaskManager 更新本地缓存;或采用定期全量刷新 |
async_insert 后可提升至 150 万+)。本文完整呈现了基于 开源 Flink + ClickHouse 构建实时数仓的分层设计、核心代码实现、优化手段及容错方案,所有代码逻辑均经过简化验证,可直接作为技术选型和开发参考。需要说明的是,本文不涉及具体业务场景的案例拆解,而是聚焦于通用架构能力,读者可根据自身业务需求调整窗口大小、维表关联方式及查询模式。
要长期稳定运行这套架构,还需重点关注:
rebalance 分区。实时数仓没有银弹,需根据业务场景持续调优。本文提供的基础架构和优化手段可覆盖大多数高并发实时聚合需求,后续可进一步探索 Flink CDC + Paimon 构建湖仓一体 或 ClickHouse 分布式 DDL 管理 等更前沿的方向。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。