实时数仓架构实战:Kafka + Flink + ClickHouse 流批一体

从零构建基于 Kafka + Flink + ClickHouse 的实时数仓,涵盖数据采集、流式计算、OLAP 存储与实时大盘的全链路实战。

在数据驱动决策日益重要的今天,传统的 T+1 离线数仓已无法满足业务对实时性的苛刻要求。电商大促的秒级监控、金融风控的毫秒级拦截、物联网设备的实时告警,都要求数据在产生后的秒级甚至毫秒级内完成采集、处理与呈现。本文将带你从零构建一套基于 Kafka + Flink + ClickHouse 的实时数仓架构,实现真正的流批一体(Streaming-Batch Unification)。


一、实时数仓架构概览

一套成熟的实时数仓通常采用 Lambda 或 Kappa 架构。为了兼顾历史数据回溯与实时增量处理,我们采用改良后的 Kappa+ 架构:以流处理为主链路,批处理作为异常补偿与历史回溯的兜底方案。

┌─────────────┐     ┌─────────────┐     ┌─────────────┐     ┌─────────────┐
│  业务系统    │────▶│   Kafka     │────▶│   Flink     │────▶│ ClickHouse  │
│ (MySQL/APP) │     │  消息队列    │     │  实时计算    │     │  OLAP 引擎   │
└─────────────┘     └─────────────┘     └─────────────┘     └─────────────┘
                                              │
                                              ▼
                                       ┌─────────────┐
                                       │  实时仪表盘  │
                                       │ (Grafana/BI)│
                                       └─────────────┘

各层职责如下:

  • 数据采集层:通过 Canal/Debezium 采集 MySQL Binlog,或业务服务直接埋点上报,统一打入 Kafka Topic。
  • 消息中间件层:Kafka 作为高吞吐量、低延迟的分布式日志流平台,承担数据暂存与削峰填谷的职责。
  • 实时计算层:Apache Flink 负责 Exactly-Once 语义的流式计算,包括数据清洗(ETL)、维度关联、窗口聚合与异常检测。
  • OLAP 存储层:ClickHouse 利用 MergeTree 引擎家族提供毫秒级查询响应,支撑高并发实时分析场景。
  • 应用展示层:Grafana、Superset 或自研 BI 系统直接查询 ClickHouse,构建实时数据大屏与自助分析。
# docker-compose.yml 核心服务定义
version: "3.8"
services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.6.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181

  kafka:
    image: confluentinc/cp-kafka:7.6.0
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1

  jobmanager:
    image: flink:1.19-scala_2.12
    command: jobmanager
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager

  taskmanager:
    image: flink:1.19-scala_2.12
    command: taskmanager
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager

  clickhouse:
    image: clickhouse/clickhouse-server:24.3
    ports:
      - "8123:8123"
      - "9000:9000"

二、Kafka 数据源层:统一数据入口

Kafka 在实时数仓中扮演着"数据总线"的角色。所有业务系统的变更数据(CDC)与事件数据(Event)都应统一接入 Kafka,形成标准化的数据流。

2.1 使用 Debezium 采集 MySQL CDC

Debezium 是 CDC 领域的事实标准,能够捕获 MySQL、PostgreSQL 等数据库的增删改操作,并以 Avro/JSON 格式输出到 Kafka。

// Debezium MySQL Connector 配置
{
  "name": "mysql-cdc-connector",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "database.hostname": "mysql",
    "database.port": "3306",
    "database.user": "debezium",
    "database.password": "dbz",
    "database.server.id": "184054",
    "database.server.name": "dbserver1",
    "database.include.list": "ecommerce",
    "table.include.list": "ecommerce.orders,ecommerce.order_items",
    "database.history.kafka.bootstrap.servers": "kafka:9092",
    "database.history.kafka.topic": "schema-changes.ecommerce",
    "include.schema.changes": "true",
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.drop.tombstones": "false",
    "transforms.unwrap.delete.handling.mode": "rewrite"
  }
}
# 注册 Debezium Connector
curl -i -X POST http://localhost:8083/connectors \
  -H "Content-Type: application/json" \
  -d @debezium-mysql-connector.json

# 查看 Topic 是否自动生成
kafka-topics --list --bootstrap-server localhost:9092
# 预期输出包含: dbserver1.ecommerce.orders, dbserver1.ecommerce.order_items

2.2 业务埋点直接上报

对于前端点击、APP 行为等事件,通常由业务服务直接序列化为 JSON 后发送到 Kafka。

// 业务埋点 Kafka Producer 示例 (Spring Boot)
@Component
public class EventTracker {
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    public void trackOrderSubmit(OrderEvent event) {
        String payload = JsonUtils.toJson(Map.of(
            "event_type", "order_submit",
            "user_id", event.getUserId(),
            "order_id", event.getOrderId(),
            "amount", event.getAmount(),
            "timestamp", System.currentTimeMillis(),
            "properties", event.getExtra()
        ));

        kafkaTemplate.send("user_behavior_events", event.getUserId(), payload)
            .whenComplete((result, ex) -> {
                if (ex != null) {
                    log.error("Failed to send event", ex);
                }
            });
    }
}

2.3 Kafka Topic 命名规范

良好的 Topic 命名规范是后续数据治理的基础。建议采用三层结构:数据源.业务域.数据类型

db.ecommerce.orders          -- 数据库 CDC: 电商域订单表
db.ecommerce.order_items     -- 数据库 CDC: 电商域订单商品表
app.user.behavior_events     -- 应用埋点: 用户行为事件
app.fraud.risk_events        -- 应用埋点: 风控事件
log.nginx.access             -- 日志数据: Nginx 访问日志

Apache Flink 是当前业界最流行的实时计算引擎,其基于 Checkpoint 的 Exactly-Once 语义与强大的 Window 算子,使其成为实时数仓计算层的不二之选。

3.1 基础 ETL 作业:订单数据清洗

以下是一个完整的 Flink SQL 作业,从 Kafka 读取订单数据,进行清洗、过滤与字段标准化后写入 ClickHouse。

-- Flink SQL: 创建 Kafka Source 表
CREATE TABLE kafka_orders (
    order_id        BIGINT,
    user_id         BIGINT,
    total_amount    DECIMAL(18, 2),
    order_status    STRING,
    create_time     TIMESTAMP(3),
    proctime AS PROCTIME(),
    event_time AS create_time,
    WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
    'connector' = 'kafka',
    'topic' = 'db.ecommerce.orders',
    'properties.bootstrap.servers' = 'kafka:9092',
    'properties.group.id' = 'flink-etl-orders',
    'format' = 'json',
    'json.fail-on-missing-field' = 'false',
    'json.ignore-parse-errors' = 'true',
    'scan.startup.mode' = 'latest-offset'
);

-- 创建 ClickHouse Sink 表
CREATE TABLE ch_orders (
    order_id        UInt64,
    user_id         UInt64,
    total_amount    Decimal(18, 2),
    order_status    LowCardinality(String),
    create_time     DateTime64(3),
    etl_time        DateTime64(3)
) WITH (
    'connector' = 'jdbc',
    'url' = 'jdbc:clickhouse://clickhouse:8123/default',
    'table-name' = 'orders_realtime',
    'username' = 'default',
    'password' = '',
    'driver' = 'com.clickhouse.jdbc.ClickHouseDriver',
    'sink.buffer-flush.max-rows' = '10000',
    'sink.buffer-flush.interval' = '5s'
);

-- 清洗与写入
INSERT INTO ch_orders
SELECT
    order_id,
    user_id,
    total_amount,
    COALESCE(NULLIF(TRIM(order_status), ''), 'UNKNOWN') AS order_status,
    create_time,
    NOW() AS etl_time
FROM kafka_orders
WHERE total_amount > 0
  AND order_id IS NOT NULL;

3.2 窗口聚合:分钟级 GMV 统计

实时数仓的核心价值之一在于低延迟的指标聚合。以下示例展示了如何使用 Flink 的 Tumble 窗口计算每分钟的 GMV(成交总额)。

CREATE TABLE ch_minute_gmv (
    window_start    DateTime64(3),
    window_end      DateTime64(3),
    order_count     UInt64,
    total_gmv       Decimal(18, 2),
    paid_gmv        Decimal(18, 2)
) WITH (
    'connector' = 'jdbc',
    'url' = 'jdbc:clickhouse://clickhouse:8123/default',
    'table-name' = 'minute_gmv_realtime'
);

INSERT INTO ch_minute_gmv
SELECT
    TUMBLE_START(event_time, INTERVAL '1' MINUTE) AS window_start,
    TUMBLE_END(event_time, INTERVAL '1' MINUTE)   AS window_end,
    COUNT(*)                                      AS order_count,
    SUM(total_amount)                             AS total_gmv,
    SUM(CASE WHEN order_status = 'PAID' THEN total_amount ELSE 0 END) AS paid_gmv
FROM kafka_orders
GROUP BY TUMBLE(event_time, INTERVAL '1' MINUTE);

3.3 DataStream API:自定义 Watermark 与状态计算

对于复杂逻辑,Flink 的 DataStream API 提供了比 SQL 更精细的控制能力。以下示例实现了一个带状态的用户会话识别算法。

public class UserSessionJob {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(60000, CheckpointingMode.EXACTLY_ONCE);
        env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000);

        KafkaSource<Event> source = KafkaSource.<Event>builder()
            .setBootstrapServers("kafka:9092")
            .setTopics("app.user.behavior_events")
            .setGroupId("flink-session-processor")
            .setStartingOffsets(OffsetsInitializer.latest())
            .setValueOnlyDeserializer(new EventDeserializationSchema())
            .build();

        DataStream<SessionResult> sessions = env.fromSource(
                source, WatermarkStrategy
                    .<Event>forBoundedOutOfOrderness(Duration.ofSeconds(30))
                    .withTimestampAssigner((event, ts) -> event.getTimestamp()),
                "Kafka Source"
            )
            .keyBy(Event::getUserId)
            .process(new SessionTimeoutFunction(Time.minutes(30)));

        sessions.addSink(new ClickHouseSink<>());
        env.execute("User Session Identification");
    }
}

public class SessionTimeoutFunction extends KeyedProcessFunction<Long, Event, SessionResult> {
    private ValueState<SessionState> sessionState;
    private final long sessionGapMs;

    @Override
    public void open(Configuration parameters) {
        sessionState = getRuntimeContext().getState(
            new ValueStateDescriptor<>("session", SessionState.class)
        );
    }

    @Override
    public void processElement(Event event, Context ctx, Collector<SessionResult> out) throws Exception {
        SessionState state = sessionState.value();
        long currentTime = ctx.timestamp();

        if (state == null || (currentTime - state.lastEventTime) > sessionGapMs) {
            if (state != null) {
                out.collect(new SessionResult(state.userId, state.startTime, state.lastEventTime, state.eventCount));
            }
            state = new SessionState(event.getUserId(), currentTime, currentTime, 1);
        } else {
            state.lastEventTime = currentTime;
            state.eventCount++;
        }
        sessionState.update(state);
        ctx.timerService().registerEventTimeTimer(currentTime + sessionGapMs);
    }

    @Override
    public void onTimer(long timestamp, OnTimerContext ctx, Collector<SessionResult> out) throws Exception {
        SessionState state = sessionState.value();
        if (state != null && timestamp >= state.lastEventTime + sessionGapMs) {
            out.collect(new SessionResult(state.userId, state.startTime, state.lastEventTime, state.eventCount));
            sessionState.clear();
        }
    }
}

四、ClickHouse 存储层:高性能 OLAP 引擎

ClickHouse 是一个面向列的 OLAP 数据库,其向量化执行引擎与数据压缩能力,使其在实时分析场景下表现卓越。在实时数仓中,我们需要根据查询模式精心设计表结构。

4.1 订单明细表设计

明细表通常采用 MergeTree 引擎,并以时间戳作为分区键,以查询维度作为排序键。

-- 订单实时明细表
CREATE TABLE IF NOT EXISTS orders_realtime (
    order_id        UInt64,
    user_id         UInt64,
    total_amount    Decimal(18, 2),
    order_status    LowCardinality(String),
    create_time     DateTime64(3),
    etl_time        DateTime64(3),
    sign            Int8 DEFAULT 1
) ENGINE = MergeTree()
PARTITION BY toYYYYMMDD(create_time)
ORDER BY (order_status, intHash32(user_id), create_time)
TTL create_time + INTERVAL 90 DAY
SETTINGS index_granularity = 8192;

关键设计要点:

  • PARTITION BY:按天分区,便于 TTL 自动过期与分区裁剪。
  • ORDER BY:将查询过滤频繁的 order_status 放在最前;使用 intHash32(user_id) 避免单个用户的写入热点。
  • LowCardinality:对枚举值类型(如订单状态)使用该修饰符,可显著降低存储体积并提升查询速度。
  • TTL:设置 90 天过期,冷数据可自动下沉到对象存储(如 S3)。

4.2 聚合表设计:ReplacingMergeTree

对于需要频繁更新状态的表(如最新订单状态),使用 ReplacingMergeTree 并在查询时加 FINAL 关键字。

CREATE TABLE IF NOT EXISTS orders_latest (
    order_id        UInt64,
    user_id         UInt64,
    total_amount    Decimal(18, 2),
    order_status    LowCardinality(String),
    update_time     DateTime64(3),
    version         UInt64
) ENGINE = ReplacingMergeTree(version)
PARTITION BY toYYYYMMDD(update_time)
ORDER BY order_id;

INSERT INTO orders_latest VALUES (1001, 42, 199.99, 'PAID', now64(), 1);
INSERT INTO orders_latest VALUES (1001, 42, 199.99, 'SHIPPED', now64(), 2);

-- 查询时必须使用 FINAL 获取最新版本
SELECT * FROM orders_latest FINAL WHERE order_id = 1001;

4.3 分布式表与集群写入

生产环境中 ClickHouse 通常以分片+副本的高可用集群部署。通过分布式表(Distributed Engine)实现数据的自动分片与查询路由。

-- 在每个分片上创建本地表
CREATE TABLE IF NOT EXISTS orders_local ON CLUSTER 'prod_cluster' (
    order_id        UInt64,
    user_id         UInt64,
    total_amount    Decimal(18, 2),
    order_status    LowCardinality(String),
    create_time     DateTime64(3)
) ENGINE = MergeTree()
PARTITION BY toYYYYMMDD(create_time)
ORDER BY (order_status, intHash32(user_id), create_time);

-- 创建分布式代理表
CREATE TABLE IF NOT EXISTS orders_distributed ON CLUSTER 'prod_cluster' AS orders_local
ENGINE = Distributed('prod_cluster', 'default', 'orders_local', intHash32(user_id));

Flink 写入时直接面向分布式表,由 ClickHouse 根据分片键(此处为 intHash32(user_id))自动路由到对应节点。


五、ClickHouse 物化视图:实时预聚合加速

物化视图(Materialized View)是 ClickHouse 的杀手级特性。它能够在数据写入时自动触发预聚合计算,将聚合结果实时同步到目标表,从而将复杂的 OLAP 查询转化为简单的点查。

5.1 构建分钟级 GMV 预聚合表

-- 1. 创建目标聚合表
CREATE TABLE IF NOT EXISTS mv_minute_gmv (
    minute          DateTime,
    order_status    LowCardinality(String),
    order_count     UInt64,
    total_amount    AggregateFunction(sum, Decimal(18, 2)),
    paid_amount     AggregateFunction(sumIf, Decimal(18, 2), UInt8)
) ENGINE = AggregatingMergeTree()
PARTITION BY toYYYYMMDD(minute)
ORDER BY (minute, order_status);

-- 2. 创建物化视图,自动从明细表聚合
CREATE MATERIALIZED VIEW IF NOT EXISTS mv_minute_gmv_view
TO mv_minute_gmv
AS SELECT
    toStartOfMinute(create_time) AS minute,
    order_status,
    countState()                 AS order_count,
    sumState(total_amount)       AS total_amount,
    sumIfState(total_amount, order_status = 'PAID') AS paid_amount
FROM orders_realtime
GROUP BY minute, order_status;

-- 3. 查询时使用 -Merge 合并状态
SELECT
    minute,
    order_status,
    countMerge(order_count)     AS order_count,
    sumMerge(total_amount)      AS total_gmv,
    sumMerge(paid_amount)       AS paid_gmv
FROM mv_minute_gmv
WHERE minute >= now() - INTERVAL 1 HOUR
GROUP BY minute, order_status
ORDER BY minute;

5.2 用户行为漏斗的物化视图

电商场景中常用的漏斗分析也可以通过物化视图实时构建。

CREATE TABLE IF NOT EXISTS mv_funnel_step (
    window_start    DateTime,
    user_id         UInt64,
    step1_pv        UInt8,
    step2_cart      UInt8,
    step3_order     UInt8,
    step4_pay       UInt8
) ENGINE = SummingMergeTree()
PARTITION BY toYYYYMMDD(window_start)
ORDER BY (window_start, user_id);

CREATE MATERIALIZED VIEW IF NOT EXISTS mv_funnel_step_view
TO mv_funnel_step
AS SELECT
    toStartOfHour(event_time) AS window_start,
    user_id,
    maxIf(1, event_type = 'page_view')   AS step1_pv,
    maxIf(1, event_type = 'add_cart')    AS step2_cart,
    maxIf(1, event_type = 'submit_order') AS step3_order,
    maxIf(1, event_type = 'pay_success') AS step4_pay
FROM user_behavior_events
GROUP BY window_start, user_id;

通过 -MergeSummingMergeTree 的自动折叠,查询漏斗总体转化率时只需对预聚合结果做简单汇总,耗时从秒级降至毫秒级。


六、实时仪表盘与可视化

数据只有被看到才有价值。实时数仓的最后一公里是将 ClickHouse 中的数据通过可视化工具呈现给业务人员。

6.1 Grafana + ClickHouse 数据源

Grafana 通过安装 ClickHouse 插件即可直接对接。以下是一个展示实时 GMV 趋势的 Grafana Query:

-- Grafana Panel Query: 实时 GMV 趋势
SELECT
    toStartOfMinute(create_time) AS time,
    sumIf(total_amount, order_status = 'PAID') AS "已支付 GMV",
    sum(total_amount) AS "总 GMV"
FROM orders_realtime
WHERE create_time >= $__from AND create_time <= $__to
GROUP BY time
ORDER BY time;
-- Grafana Panel Query: Top 10 热销商品 (关联订单商品表)
SELECT
    p.product_name,
    sum(oi.quantity) AS total_qty,
    sum(oi.quantity * oi.unit_price) AS total_revenue
FROM order_items_realtime oi
JOIN products p ON oi.product_id = p.product_id
WHERE oi.create_time >= today()
GROUP BY p.product_name
ORDER BY total_revenue DESC
LIMIT 10;

6.2 基于 WebSocket 的实时推送大屏

对于需要秒级刷新的指挥中心大屏,可以通过后端服务轮询 ClickHouse 并向 WebSocket 客户端推送增量数据。

# Python FastAPI + ClickHouse 实时推流示例
import asyncio
from datetime import datetime, timedelta
from fastapi import FastAPI, WebSocket
from clickhouse_driver import Client

app = FastAPI()
ch_client = Client(host='localhost', port=9000)

@app.websocket("/ws/realtime-dashboard")
async def realtime_dashboard(websocket: WebSocket):
    await websocket.accept()
    try:
        while True:
            result = ch_client.execute("""
                SELECT
                    toStartOfMinute(now()) AS minute,
                    count() AS order_count,
                    sum(total_amount) AS gmv
                FROM orders_realtime
                WHERE create_time >= now() - INTERVAL 1 MINUTE
            """)
            payload = {
                "timestamp": datetime.utcnow().isoformat(),
                "data": result
            }
            await websocket.send_json(payload)
            await asyncio.sleep(5)
    except Exception:
        await websocket.close()

# 前端 ECharts 实时折线图即可订阅此 WebSocket 并动态 appendData

七、性能优化与生产调优

实时数仓上线后,持续的性能调优是保障 SLA 的关键。以下为常见优化项对比表:

优化维度优化前优化后实施手段
Kafka 吞吐10 MB/s80 MB/s调整 batch.size=65536, linger.ms=10, 启用 Snappy 压缩
Flink 延迟5-10s<2s启用 Mini-Batch + Local-Global 聚合,优化 Watermark 间隔
ClickHouse 写入5k 行/s50k+ 行/s使用 Batch Insert(单次>=10k 行),关闭 fsync,异步写入
ClickHouse 查询3-5s50-200ms引入物化视图预聚合,使用 PREWHERE 替代 WHERE,分区裁剪
存储成本100%35%启用列压缩(LZ4/ZSTD),TTL 下沉冷数据到 S3,删除冗余字段
端到端延迟2-5 min5-15s缩短 Flink Checkpoint 间隔,ClickHouse 使用内存 Buffer 表
// 生产级 Checkpoint 配置
env.enableCheckpointing(30000); // 30s 间隔
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(15000);
env.getCheckpointConfig().setCheckpointTimeout(600000);
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
env.getCheckpointConfig().enableExternalizedCheckpoints(
    ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION
);
// 状态后端建议 RocksDB + Incremental Checkpoint
env.setStateBackend(new EmbeddedRocksDBStateBackend(true));
env.getCheckpointConfig().setCheckpointStorage("hdfs://namenode:8020/flink-checkpoints");

ClickHouse 写入调优

-- 调整写入相关参数 (users.xml 或会话级)
SET max_insert_threads = 8;
SET max_execution_time = 300;
SET max_memory_usage = 16_000_000_000;

-- 使用 Buffer 表做异步聚合写入
CREATE TABLE orders_buffer AS orders_realtime
ENGINE = Buffer('default', 'orders_realtime',
    16,            -- num_layers
    10,            -- min_time (秒)
    100,           -- max_time (秒)
    100000,        -- min_rows
    1000000,       -- max_rows
    10000000,      -- min_bytes
    100000000      -- max_bytes
);
-- 业务写入 orders_buffer,ClickHouse 自动按阈值批量刷入 orders_realtime

八、常见问题解答(FAQ)

Q1: Flink 消费 Kafka 时出现数据延迟抖动,如何排查?

首先通过 Flink Web UI 的 Backpressure 标签页查看是否存在反压。若 Source 端出现 backpressure,通常是下游算子处理能力不足,可通过增加 TaskManager 的 slot 数量或并行度解决。若 Kafka Consumer Lag 持续增长但无反压,则需要检查 Kafka 分区数是否过少导致单并行度瓶颈。另一个常见原因是 Watermark 久久不推进(如某分区数据稀疏),此时可考虑使用 IdleTimeout 标记空闲分区。

WatermarkStrategy
    .<Event>forBoundedOutOfOrderness(Duration.ofSeconds(30))
    .withTimestampAssigner(...)
    .withIdleness(Duration.ofMinutes(2)); // 关键配置

Q2: ClickHouse 的 ReplacingMergeTree 在查询时为什么要加 FINAL?有没有替代方案?

ReplacingMergeTree 的合并操作是异步后台执行的,写入相同主键的数据在物理上可能仍存在于不同 part 中。FINAL 关键字会在查询时实时合并这些 part,确保返回最新版本,但这会牺牲查询性能。

替代方案有两种:一是使用 optimize table ... final 强制后台合并(不建议在生产库上频繁执行);二是通过业务层保证幂等写入,例如只插入最终状态而不保留中间状态;三是改用 VersionedCollapsingMergeTree 并在查询时自行根据版本号过滤。

Q3: 实时数仓与离线数仓的指标口径不一致,该如何治理?

口径不一致是 Lambda 架构的经典痛点。推荐采取以下措施:

  1. 统一维度层:将用户、商品、地区等维度表统一维护在 Hive/Iceberg 中,由 Flink 通过 Async Lookup Join 实时关联,确保维度属性一致。
  2. 统一指标定义:使用指标平台(如 Apache Doris/StarRocks 的视图层,或自研指标元数据中心)统一定义 SUM、COUNT、UV 等计算逻辑,实时与离线复用同一套 SQL 模板。
  3. DWD 层对齐:将 Kafka 中的实时流与离线 ODS 层数据通过一致性校验工具(如 Great Expectations)定期比对,发现差异及时溯源。

Q4: 数据量持续增长后,ClickHouse 单节点成为瓶颈,如何水平扩展?

ClickHouse 原生支持分布式集群。水平扩展的核心步骤如下:

  1. 搭建包含 ZooKeeper 的 ClickHouse 集群,定义 <remote_servers> 中的分片与副本配置。
  2. 所有本地表使用 ON CLUSTER 语法创建,确保元数据在所有节点间同步。
  3. 写入时面向 Distributed 表,由 ClickHouse 根据分片键(如 rand()user_id 的哈希)自动路由。
  4. 查询时也通过 Distributed 表发起,ClickHouse 会在各分片并行执行后汇总结果返回。
  5. 扩容时只需增加新节点并修改集群配置,历史数据可通过 ALTER TABLE ... MOVE PARTITION 或重新分布式写入实现再平衡。

九、总结

本文完整展示了从数据采集到实时可视化的全链路实时数仓构建过程:

  • Kafka 作为统一数据总线,承载高吞吐的事件流与 CDC 数据;
  • Flink 提供 Exactly-Once 的流式计算能力,完成清洗、关联与窗口聚合;
  • ClickHouse 以列存与向量化引擎支撑毫秒级 OLAP 查询,物化视图进一步将复杂分析转化为实时点查;
  • Grafana / 自研大屏 完成数据的最终消费,实现业务决策的实时化。

该架构已在多个电商与金融场景中落地验证,端到端延迟可稳定控制在 10 秒以内,ClickHouse 查询 P99 延迟低于 200ms。随着 Apache Flink 的 Table Store(Paimon)与 ClickHouse 的 S3 表引擎持续演进,未来实时与离线架构将进一步走向统一,流批一体的愿景正在加速成为现实。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「数据工程」更多文章

  1. BI 可视化工具深度对比:Tableau、Power BI、Superset 与 Metabase