首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >五层架构实战:基于腾讯云EMR与Flink CDC构建企业级实时离线一体化大数据平台

五层架构实战:基于腾讯云EMR与Flink CDC构建企业级实时离线一体化大数据平台

原创
作者头像
用户12678265
发布2026-08-16 14:03:37
发布2026-08-16 14:03:37
420
举报

五层架构实战:基于腾讯云EMR与Flink CDC构建企业级实时离线一体化大数据平台

本文源于已完结的19章企业级大数据训练营终极项目,核心解决数据孤岛任务血缘混乱SLA无法保障三大痛点。文章全程代码驱动,涵盖MySQL CDC实时采集、Hive 数仓分层(ODS/DWD/DWS/ADS)、Flink流式ETL、DolphinScheduler跨层依赖调度以及元数据质量校验,拒绝泛泛而谈。

1. 多层次架构全景设计(第一层:总体架构)

企业级大数据平台绝非工具的堆砌,而是分层复用职责分离的体现。我们设计的五层架构如下:

层级

功能定位

腾讯云技术选型

核心挑战

采集层

业务库日志/CDC实时接入

Canal + Flink CDC + Kafka

全量/增量同步无丢失、断点续传

存储层

原始数据冷备与索引

COS(对象存储)+ Hive Metastore

分区治理与小文件合并

计算层

离线T+1批处理 & 实时流处理

EMR(Spark SQL)+ Flink on EMR

数据倾斜根治、状态后端优化

调度层

复杂任务依赖编排与重试

DolphinScheduler(工作流)

跨层级递归依赖、死锁检测

治理层

数据质量监控与血缘分析

DataHub(元数据)+ 自研DQC规则

行数波动监控、空值率告警


2. 采集与存储层:Flink CDC 实时入湖(COS)

企业级平台的第一要务是毫秒级捕捉业务变更。我们使用 Flink CDC 2.4 直接读取 MySQL Binlog,同步至腾讯云 COS 并以 Hudi(或Parquet格式) 存储,保留变更历史(CDC日志)。

2.1 Flink CDC 启动脚本(含断点续传)

代码语言:javascript
复制
// 使用Flink SQL DDL 定义 MySQL CDC 源表
CREATE TABLE source_orders (
    order_id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),
    order_status STRING,
    cdc_timestamp TIMESTAMP(3) METADATA FROM 'op_ts' VIRTUAL,
    PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = 'your-mysql.com',
    'port' = '3306',
    'username' = 'cdc_user',
    'password' = '***',
    'database-name' = 'ecom_db',
    'table-name' = 'orders',
    'scan.startup.mode' = 'initial',  -- 全量+增量,记录binlog位点
    'server-id' = '5401-5404',        -- 多并行度需分配不同server-id
    'debezium.snapshot.mode' = 'initial'
);

-- 写入 COS(使用 Hive 格式,分区按天)
CREATE TABLE dwd_orders (
    order_id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),
    status STRING,
    event_time STRING
) PARTITIONED BY (dt STRING)
WITH (
    'connector' = 'filesystem',
    'path' = 'cosn://ecommerce-warehouse/dwd/orders', -- 腾讯云COS路径
    'format' = 'parquet',
    'sink.partition-commit.trigger' = 'process-time',
    'sink.partition-commit.delay' = '1 min'
);

-- 插入实时流
INSERT INTO dwd_orders 
SELECT order_id, user_id, amount, order_status, 
       DATE_FORMAT(cdc_timestamp, 'yyyy-MM-dd') AS dt
FROM source_orders;

2.2 解决小文件堆积(关键代码)

实时入湖会产生大量小时级甚至分钟级小文件,拖垮 NameNode。我们在 EMR 侧通过 Spark 定时合并(Coalesce)处理:

代码语言:javascript
复制
// 每天凌晨2点执行,合并前一天分区的小文件
val df = spark.read.parquet(s"cosn://warehouse/dwd/orders/dt=$yesterday")
// 重分区为合理并行度(每个分区文件大小约为256MB)
val optimizedDF = df.coalesce(10) 
optimizedDF.write
  .mode("overwrite")
  .option("compression", "snappy")
  .parquet(s"cosn://warehouse/dwd/orders_merged/dt=$yesterday")
// 原子性替换(使用Hive MSCK REPAIR TABLE 或 直接覆盖)

3. 离线计算层(第三层):Spark SQL 四层模型(ODS->DWD->DWS->ADS)

多层次的核心在于数据逐层建模。我们在 EMR 上使用 Spark SQL 执行批处理 ETL,引入拉链表处理缓慢变化维度(SCD Type 2)。

3.1 DWD 层:明细日志打宽(维度退化)

将订单明细维表进行广播 Join(Map端聚合,避免Shuffle):

代码语言:javascript
复制
-- 设置广播阈值
SET spark.sql.autoBroadcastJoinThreshold = 104857600; -- 100MB
SET spark.sql.adaptive.enabled = true;
SET spark.sql.adaptive.skewJoin.enabled = true; -- 腾讯云EMR自适应执行解决数据倾斜

INSERT OVERWRITE TABLE dwd_order_detail PARTITION(dt='${bizdate}')
SELECT 
    o.order_id,
    o.user_id,
    u.user_name,
    u.user_level,
    p.product_name,
    p.category_id,
    o.amount,
    o.order_status,
    unix_timestamp(o.pay_time) - unix_timestamp(o.create_time) as pay_duration_sec
FROM ods_orders o
JOIN ods_users u ON o.user_id = u.user_id AND u.dt='${bizdate}'
JOIN ods_products p ON o.product_id = p.product_id AND p.dt='${bizdate}'
WHERE o.dt='${bizdate}';

3.2 DWS 层:多维度汇总(预聚合)

面向业务主题(用户、商品、渠道)构建宽表,使用 GROUPING SETS 实现多维度上卷:

代码语言:javascript
复制
INSERT OVERWRITE TABLE dws_user_category_agg PARTITION(dt='${bizdate}')
SELECT 
    user_id,
    category_id,
    SUM(amount) as gmv,
    COUNT(DISTINCT order_id) as order_cnt,
    AVG(amount) as avg_order_value,
    GROUPING__ID as group_level -- 用于区分汇总粒度
FROM dwd_order_detail
WHERE dt='${bizdate}'
GROUP BY user_id, category_id
GROUPING SETS ((user_id, category_id), (user_id), (category_id), ());

3.3 ADS 层:同步至腾讯云数据库 MySQL(供业务查询)

使用 Spark JDBC 并行写入,注意批次控制以防锁表

代码语言:javascript
复制
# PySpark 代码
jdbc_url = "jdbc:mysql://your-cdb.mysql.tencentcdb.com:3306/ads_db?useSSL=false"
properties = {"user": "ads_writer", "password": "***", "driver": "com.mysql.jdbc.Driver"}

# 覆盖写入当天的运营大盘数据
df_ads = spark.sql("SELECT * FROM ads_dashboard WHERE dt='2026-08-15'")
df_ads.write.jdbc(url=jdbc_url, table="dashboard_daily", mode="overwrite", properties=properties)

4. 实时计算层(第四层):Flink 双流Join与状态过期

除了离线T+1,平台还需支持实时大屏(5秒刷新)。我们在EMR上部署Flink On YARN,使用双流Join关联订单流与支付流,并设置合理的 State TTL 防止状态无限膨胀。

代码语言:javascript
复制
// Flink Java/Scala 核心逻辑伪代码
DataStream<OrderEvent> orderStream = ...;
DataStream<PayEvent> payStream = ...;

orderStream.join(payStream)
    .where(OrderEvent::getOrderId)
    .equalTo(PayEvent::getOrderId)
    .window(JoinWindows.of(Time.seconds(10))) // 10秒间隔
    .apply(new JoinFunction<OrderEvent, PayEvent, RichOrder>() {
        @Override
        public RichOrder join(OrderEvent order, PayEvent pay) {
            // 实时计算支付转化时延
            return new RichOrder(order, pay.getPayTime());
        }
    })
    .map(rich -> {
        // 输出到 Redis (腾讯云 CRS) 供大屏查询
        // 使用 Jedis Cluster 批量写入 Pipeline
    });

// 设置 State TTL (防止状态堆积)
ValueStateDescriptor<T> descriptor = new ValueStateDescriptor<>("order-state", TypeInformation.of(...));
descriptor.enableTimeToLive(StateTtlConfig
    .newBuilder(Time.hours(24))
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
    .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
    .build());

5. 调度层(第五层):DolphinScheduler 跨层递归依赖

企业级调度不能只依赖Crontab。我们使用 DolphinScheduler 定义工作流(Workflow),利用其依赖节点(Dependent Node) 自动解析上游任务是否成功,完美解决跨层数据滞后问题。

5.1 定义复杂依赖(使用 Python API 或 UI 配置)

代码语言:javascript
复制
# 伪配置文件,实际UI操作
项目: 电商数仓
工作流: DWD_DWS_ADS_Flow
  - 任务A (ODS导入) : Shell脚本,执行sqoop或spark load
  - 任务B (DWD层ETL) : Spark SQL,依赖A完成
  - 任务C (DWS层聚合) : Spark SQL,依赖B完成
  - 任务D (ADS导出MySQL) : Spark JDBC,依赖C完成
  - 任务E (数据质量校验) : Python脚本,依赖B和C完成 (行数校验)
  - 依赖配置: 任务B的`上游依赖`设置为任务A,且任务A必须状态为success。
  - 失败重试: 重试3次,间隔5分钟,发送企微告警。

5.2 补数据与回刷(解决历史数据修正)

当业务逻辑变更时,需要回刷近30天数据。DolphinScheduler 支持 “补数” 功能,动态传入 {bizdate} 参数,并行启动多个工作流实例:

代码语言:javascript
复制
# 补数脚本 (通过API触发)
curl -X POST http://dolphin-api/projects/1/executors/start \
  -d '{"processDefinitionCode": 123, "scheduleTime": "2026-07-15 00:00:00..2026-08-15 00:00:00", "parallelism": 5}'

核心技术点:必须确保同一分区(dt)的数据幂等(Insert Overwrite),避免重复补数导致数据翻倍。


6. 治理层:元数据血缘与DQC质量闭环

全能型大数据开发必须具备数据治理意识。我们利用腾讯云 DataHub(或Atlas) 采集血缘,并自研基于Great Expectations的质量校验。

6.1 数据质量校验(DQC)核心代码

在ADS层导出前,必须执行质量检查,否则阻断下游:

代码语言:javascript
复制
# PySpark DQC 校验脚本
def dqc_check(table_name, date):
    df = spark.sql(f"SELECT * FROM {table_name} WHERE dt='{date}'")
    total_count = df.count()
    
    # 规则1:波动率不能超过20%(对比7天均值)
    avg_7d = spark.sql(f"SELECT AVG(cnt) FROM (SELECT COUNT(*) as cnt FROM {table_name} WHERE dt BETWEEN date_sub('{date}',7) AND date_sub('{date}',1) GROUP BY dt) tmp").collect()[0][0]
    if abs(total_count - avg_7d) / avg_7d > 0.2:
        raise Exception(f"DQC 失败: 表{table_name} 日活波动超过20%,当前{total_count},历史均值{avg_7d}")
    
    # 规则2:关键字段空值率 < 5%
    null_rate = df.filter(df.order_id.isNull()).count() / total_count
    if null_rate > 0.05:
        raise Exception(f"DQC 失败: order_id 空值率 {null_rate} > 5%")
    
    print("Quality check passed.")

6.2 血缘自动解析(获取影响分析)

通过解析 spark.sql 执行计划,自动获取输入输出表:

代码语言:javascript
复制
# 利用Spark的ListenerBus 或 直接解析 SQL 的 Abstract Syntax Tree (AST)
# 此处以解析逻辑为例,实际可在EMR配置中开启血缘日志写入ES
explain_str = spark.sql("EXPLAIN EXTENDED SELECT ...").collect()[0][0]
# 正则提取 TableScan 和 InsertIntoTable
# 将血缘关系 (source->target) 上报至腾讯云元数据中心

7. 最终SLA保障与监控大盘

我们通过多层架构设计,最终将离线数仓产出时间从早上的10:00提前至凌晨6:00,实时数据延迟控制在5秒内。

  • EMR监控:配置CLS日志告警,当Executor GC时间超过10s或Task失败率>5%时触发钉钉机器人。
  • COS存储优化:设置生命周期策略,30天前的数据自动沉降到归档存储,降低成本60%。

8. 总结:多层次的“多”究竟在哪里?

很多开发者只关注写SQL,而忽略了架构的多层次性

  1. 数据链路多:实时+离线双轨并行,互为备份。
  2. 计算引擎多:Flink做流,Spark做批,各取所长。
  3. 治理手段多:从调度依赖、重试机制到DQC质量门禁,形成闭环。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 五层架构实战:基于腾讯云EMR与Flink CDC构建企业级实时离线一体化大数据平台
    • 1. 多层次架构全景设计(第一层:总体架构)
    • 2. 采集与存储层:Flink CDC 实时入湖(COS)
      • 2.1 Flink CDC 启动脚本(含断点续传)
      • 2.2 解决小文件堆积(关键代码)
    • 3. 离线计算层(第三层):Spark SQL 四层模型(ODS->DWD->DWS->ADS)
      • 3.1 DWD 层:明细日志打宽(维度退化)
      • 3.2 DWS 层:多维度汇总(预聚合)
      • 3.3 ADS 层:同步至腾讯云数据库 MySQL(供业务查询)
    • 4. 实时计算层(第四层):Flink 双流Join与状态过期
    • 5. 调度层(第五层):DolphinScheduler 跨层递归依赖
      • 5.1 定义复杂依赖(使用 Python API 或 UI 配置)
      • 5.2 补数据与回刷(解决历史数据修正)
    • 6. 治理层:元数据血缘与DQC质量闭环
      • 6.1 数据质量校验(DQC)核心代码
      • 6.2 血缘自动解析(获取影响分析)
    • 7. 最终SLA保障与监控大盘
    • 8. 总结:多层次的“多”究竟在哪里?
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档