本文源于已完结的19章企业级大数据训练营终极项目,核心解决数据孤岛、任务血缘混乱与SLA无法保障三大痛点。文章全程代码驱动,涵盖MySQL CDC实时采集、Hive 数仓分层(ODS/DWD/DWS/ADS)、Flink流式ETL、DolphinScheduler跨层依赖调度以及元数据质量校验,拒绝泛泛而谈。
企业级大数据平台绝非工具的堆砌,而是分层复用与职责分离的体现。我们设计的五层架构如下:
层级 | 功能定位 | 腾讯云技术选型 | 核心挑战 |
|---|---|---|---|
采集层 | 业务库日志/CDC实时接入 | Canal + Flink CDC + Kafka | 全量/增量同步无丢失、断点续传 |
存储层 | 原始数据冷备与索引 | COS(对象存储)+ Hive Metastore | 分区治理与小文件合并 |
计算层 | 离线T+1批处理 & 实时流处理 | EMR(Spark SQL)+ Flink on EMR | 数据倾斜根治、状态后端优化 |
调度层 | 复杂任务依赖编排与重试 | DolphinScheduler(工作流) | 跨层级递归依赖、死锁检测 |
治理层 | 数据质量监控与血缘分析 | DataHub(元数据)+ 自研DQC规则 | 行数波动监控、空值率告警 |
企业级平台的第一要务是毫秒级捕捉业务变更。我们使用 Flink CDC 2.4 直接读取 MySQL Binlog,同步至腾讯云 COS 并以 Hudi(或Parquet格式) 存储,保留变更历史(CDC日志)。
// 使用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;实时入湖会产生大量小时级甚至分钟级小文件,拖垮 NameNode。我们在 EMR 侧通过 Spark 定时合并(Coalesce)处理:
// 每天凌晨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 或 直接覆盖)多层次的核心在于数据逐层建模。我们在 EMR 上使用 Spark SQL 执行批处理 ETL,引入拉链表处理缓慢变化维度(SCD Type 2)。
将订单明细维表进行广播 Join(Map端聚合,避免Shuffle):
-- 设置广播阈值
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}';面向业务主题(用户、商品、渠道)构建宽表,使用 GROUPING SETS 实现多维度上卷:
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), ());使用 Spark JDBC 并行写入,注意批次控制以防锁表:
# 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)除了离线T+1,平台还需支持实时大屏(5秒刷新)。我们在EMR上部署Flink On YARN,使用双流Join关联订单流与支付流,并设置合理的 State TTL 防止状态无限膨胀。
// 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());企业级调度不能只依赖Crontab。我们使用 DolphinScheduler 定义工作流(Workflow),利用其依赖节点(Dependent Node) 自动解析上游任务是否成功,完美解决跨层数据滞后问题。
# 伪配置文件,实际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分钟,发送企微告警。当业务逻辑变更时,需要回刷近30天数据。DolphinScheduler 支持 “补数” 功能,动态传入 {bizdate} 参数,并行启动多个工作流实例:
# 补数脚本 (通过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),避免重复补数导致数据翻倍。
全能型大数据开发必须具备数据治理意识。我们利用腾讯云 DataHub(或Atlas) 采集血缘,并自研基于Great Expectations的质量校验。
在ADS层导出前,必须执行质量检查,否则阻断下游:
# 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.")通过解析 spark.sql 执行计划,自动获取输入输出表:
# 利用Spark的ListenerBus 或 直接解析 SQL 的 Abstract Syntax Tree (AST)
# 此处以解析逻辑为例,实际可在EMR配置中开启血缘日志写入ES
explain_str = spark.sql("EXPLAIN EXTENDED SELECT ...").collect()[0][0]
# 正则提取 TableScan 和 InsertIntoTable
# 将血缘关系 (source->target) 上报至腾讯云元数据中心我们通过多层架构设计,最终将离线数仓产出时间从早上的10:00提前至凌晨6:00,实时数据延迟控制在5秒内。
很多开发者只关注写SQL,而忽略了架构的多层次性:
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。