提示
预计用时 60 分钟。学习目标:以电商订单分析场景为例,端到端跑通 MySQL → 数据湖仓 → PySpark 加工 → 工作流编排 → 定时调度 → 运维监控 的完整数据工程链路。
场景介绍
您是一名刚加入数据团队的数据工程师,业务方希望每天自动把 MySQL 中的订单流水同步到数据湖仓,并按地区维度聚合产出收入汇总表,供下游 BI 与运营分析使用。
本教程将带您:
1. 在工作空间中注册 MySQL 数据源,并确认数据湖仓所需的计算资源就绪。
2. 通过 离线同步 把 MySQL 中的
amazon_orders 表加载到数据湖仓 Catalog。3. 在 Studio 中用 PySpark Notebook 完成数据清洗与按地区聚合,写回数据湖仓。
4. 把离线同步与 Notebook 串联成一条 DAG 工作流。
5. 配置每天 00:00 的定时调度。
6. 在工作流运行中跟踪执行状态、查看运行日志,并完成一次失败重跑。
完成本教程后,您将掌握 DataBuddy 数据工程师日常最高频的全套操作。
前提条件
在开始本教程前,请确保满足以下条件:
已开通 DataBuddy 服务并完成空间初始化。
当前账号在目标工作空间中拥有 数据工程师 角色(权限详见 成员与权限管理 )。
工作空间已绑定以下三类计算资源(详见 计算资源概述 ):
数据接入型 计算资源(用于离线同步任务)。
作业集群 (用于 Notebook 任务)。
交互式集群 (用于 SQL 调试,按需)。
已准备好一个可访问的 MySQL 实例,库表样例为
wedata_demo_db.amazon_orders(14 个字段,约 20 行示例数据,见 步骤 1 的数据准备说明)。工作空间已绑定数据湖仓 Catalog(默认
DataLakeCatalog),且当前账号具备目标 Schema wedata_demo_db 的建表与读写权限。提示
本教程使用的样例库表
wedata_demo_db.amazon_orders 仅作演示用途。您可以使用任意结构化数据源替换;同步与 PySpark 步骤同样适用。教程数据样例
amazon_orders 表字段示意:字段 | 类型 | 说明 |
order_id | VARCHAR(20) | 订单号 |
user_id | VARCHAR(20) | 用户 ID |
product_id | VARCHAR(20) | 商品 ID |
product_name | VARCHAR(200) | 商品名 |
category | VARCHAR(50) | 商品品类 |
unit_price | DECIMAL(10,2) | 单价 |
quantity | INT | 数量 |
total_amount | DECIMAL(12,2) | 订单金额 |
order_status | VARCHAR(20) | 订单状态:Delivered / Shipped / Cancelled 等 |
payment_method | VARCHAR(30) | 支付方式 |
order_date | VARCHAR(10) | 下单日期( dd/MM/yyyy) |
state_code | VARCHAR(5) | 州简码 |
state | VARCHAR(50) | 州名 |
city | VARCHAR(100) | 城市 |
步骤 1:注册 MySQL 数据源
注册数据源是后续所有同步任务的基础。数据湖仓 Catalog 通常由管理员预先绑定,本步骤只需补齐 MySQL 侧的连接信息。
1. 登录 DataBuddy 控制台,在工作空间列表中点击空间名称,进入工作空间。在左侧导航选择 管理 > 数据源管理 。
2. 在页面右上方单击 创建数据源 ,类型选择 MySQL 。
3. 填写连接信息:
数据源名称 :例如
quickstart_mysql,用于在后续任务中引用。连接串信息 :例如
jdbc:mysql://<host>:3306/wedata_demo_db。账号、密码 :MySQL 实例的访问凭据。
1. 选择一个 数据接入型 计算资源,单击 测试连通性 ,状态显示 连接成功 后单击 创建并接入数据 。
提示
数据湖仓的 Catalog(如
DataLakeCatalog)在工作空间绑定作业集群后由平台自动生成,无需手动新建。若没有该 Catalog,请联系工作空间管理员确认作业集群绑定情况。详见 添加数据源连接 。步骤 2:把 MySQL 表同步到数据湖仓
接下来创建一个离线同步任务,把
amazon_orders 表全量加载到数据湖仓 Catalog 中。1. 在左侧导航选择 数据集成 ,选择 接入任务 ,单击 创建任务 。
2. 填写基础信息:
数据源类型 :MySQL。
接入方式 :离线单表接入 。
任务名称 :例如
quickstart_mysql2dlc(建议加上自己的姓名前缀以避免多人冲突)。计算资源 :已创建平台型 计算资源 。
1. 配置来源端:
数据源类型 :MySQL。
数据源连接 :选择 步骤 1 创建的
quickstart_mysql。来源端数据库 :
wedata_demo_db。来源端数据表 :
amazon_orders。1. 配置目标端:
目标端 Catalog :
DataLakeCatalog。目标端 Schema :
wedata_demo_db(如不存在请提前创建或联系管理员)。目标端 Table :例如
quickstart_amazon_orders。如目标表不存在,单击 创建目标表 ,平台会基于来源端 Schema 生成建表 SQL,必要时手动调整后提交。写入模式 :overwrite (首次全量加载推荐使用 overwrite;增量场景请改为 append 或 upsert)。
1. 进入字段映射步骤,单击 同名映射 ,让来源端字段与目标端同名字段一一对齐。
2. 单击 保存 ,再单击 运行 进行调试执行;待状态变为 成功 后继续下一步。
3. 点击 保存并发布 ,任务进入已发布状态,后续可以在 步骤4 工作流选择当前任务进行工作流的编排。
提示
步骤 3:在 Studio 中开发并调试 PySpark Notebook
数据落到数据湖仓后,用 PySpark 完成清洗、按地区聚合,并把汇总结果写回数据湖仓。
3.1 新建 Notebook 并配置 Kernel
1. 在左侧导航选择 开发工作台 。
2. 单击 + 创建 > Notebook ,文件名例如
quickstart_transform.ipynb。3. 在顶部右侧区域点击“未连接”,搜索相关作业集群资源并进行连接。
4. 等待 Kernel 状态变为该作业集群名,则状态变为 已连接 。
提示
3.2 准备样例数据(可选)
如果您未完成 步骤 2 的离线同步,可以在 Notebook 中新建 SQL 单元格,直接创建目标表并插入样例数据:
-- 在 Notebook 中新建 SQL 单元格(单元格左上角语言选择 SQL),先执行以下建表语句-- 1. 创建 Schema(如果不存在)CREATE SCHEMA IF NOT EXISTS `DataLakeCatalog`.`wedata_demo_db`;-- 2. 创建表(如果不存在)CREATE TABLE IF NOT EXISTS `DataLakeCatalog`.`wedata_demo_db`.`quickstart_amazon_orders` (`order_id` STRING COMMENT '订单号',`user_id` STRING COMMENT '用户ID',`product_id` STRING COMMENT '商品ID',`product_name` STRING COMMENT '商品名',`category` STRING COMMENT '商品品类',`unit_price` DECIMAL(10,2) COMMENT '单价',`quantity` INT COMMENT '数量',`total_amount` DECIMAL(12,2) COMMENT '订单金额',`order_status` STRING COMMENT '订单状态',`payment_method` STRING COMMENT '支付方式',`order_date` STRING COMMENT '下单日期',`state_code` STRING COMMENT '州简码',`state` STRING COMMENT '州名',`city` STRING COMMENT '城市') COMMENT '样例订单表';
建表成功后,再新建一个 SQL 单元格执行插入数据:
-- 3. 插入样例数据INSERT INTO `DataLakeCatalog`.`wedata_demo_db`.`quickstart_amazon_orders`(`order_id`, `user_id`, `product_id`, `product_name`, `category`, `unit_price`, `quantity`, `total_amount`, `order_status`, `payment_method`, `order_date`, `state_code`, `state`, `city`)VALUES('ORD-001','U001','P001','Towel Set','Home',CAST(15.99 AS DECIMAL(10,2)),2,CAST(31.98 AS DECIMAL(12,2)),'Delivered','CreditCard','01/03/2024','CA','California','Los Angeles'),('ORD-002','U002','P002','Bedsheet','Home',CAST(25.50 AS DECIMAL(10,2)),1,CAST(25.50 AS DECIMAL(12,2)),'Shipped','PayPal','02/03/2024','NY','New York','New York City'),('ORD-003','U003','P003','Curtain','Home',CAST(40.00 AS DECIMAL(10,2)),3,CAST(120.00 AS DECIMAL(12,2)),'Delivered','CreditCard','03/03/2024','TX','Texas','Houston'),('ORD-004','U004','P004','Carpet','Home',CAST(120.50 AS DECIMAL(10,2)),1,CAST(120.50 AS DECIMAL(12,2)),'Cancelled','ApplePay','04/03/2024','FL','Florida','Miami'),('ORD-005','U005','P005','Pillow','Home',CAST(12.99 AS DECIMAL(10,2)),4,CAST(51.96 AS DECIMAL(12,2)),'Delivered','CreditCard','05/03/2024','WA','Washington','Seattle'),('ORD-006','U006','P006','Mattress','Home',CAST(299.99 AS DECIMAL(10,2)),1,CAST(299.99 AS DECIMAL(12,2)),'Delivered','DebitCard','06/03/2024','IL','Illinois','Chicago'),('ORD-007','U007','P007','Blanket','Home',CAST(35.00 AS DECIMAL(10,2)),2,CAST(70.00 AS DECIMAL(12,2)),'Shipped','CreditCard','07/03/2024','AZ','Arizona','Phoenix'),('ORD-008','U001','P001','Towel Set','Home',CAST(15.99 AS DECIMAL(10,2)),5,CAST(79.95 AS DECIMAL(12,2)),'Delivered','CreditCard','08/03/2024','CA','California','San Diego'),('ORD-009','U008','P008','Rug','Home',CAST(65.00 AS DECIMAL(10,2)),1,CAST(65.00 AS DECIMAL(12,2)),'Delivered','PayPal','09/03/2024','NV','Nevada','Las Vegas'),('ORD-010','U009','P009','Comforter','Home',CAST(85.00 AS DECIMAL(10,2)),1,CAST(85.00 AS DECIMAL(12,2)),'Shipped','CreditCard','10/03/2024','OR','Oregon','Portland');
执行成功后,可以再执行以下查询验证数据:
-- 4. 验证数据已插入SELECT * FROM `DataLakeCatalog`.`wedata_demo_db`.`quickstart_amazon_orders`;
执行成功后,切换回 PySpark 单元格继续下一步。
3.3 读取数据湖仓表
在第一个 PySpark 单元格中读取 步骤 2 同步进来的表(或 步骤 3.2 中手动创建的表):
df_orders = spark.sql("""SELECT *FROM `DataLakeCatalog`.`wedata_demo_db`.`quickstart_amazon_orders`""")display(df_orders)print(f"总行数: {df_orders.count()}")
数据湖仓中的表使用 三段名
Catalog.Schema.Table 进行引用。3.4 清洗与维度衍生
from pyspark.sql import functions as F# 过滤已取消订单df_valid = df_orders.filter(F.col("order_status") != "Cancelled")# 把 dd/MM/yyyy 文本日期解析为 date,并衍生年月df_valid = (df_valid.withColumn("order_date_parsed", F.to_date("order_date", "dd/MM/yyyy")).withColumn("order_year", F.year("order_date_parsed")).withColumn("order_month", F.month("order_date_parsed")))display(df_valid.select("order_id", "order_date", "order_date_parsed", "order_year", "order_month"))
3.5 按州维度聚合
df_state_revenue = (df_valid.groupBy("state", "state_code").agg(F.count("order_id").alias("order_count"),F.sum("total_amount").alias("total_revenue"),F.avg("total_amount").alias("avg_order_value"),F.countDistinct("user_id").alias("unique_customers"),).orderBy(F.desc("total_revenue")))display(df_state_revenue)
3.6 创建目标表并写回数据湖仓
spark.sql("""CREATE TABLE IF NOT EXISTS `DataLakeCatalog`.`wedata_demo_db`.`quickstart_state_revenue` (state STRING COMMENT '州名',state_code STRING COMMENT '州简码',order_count BIGINT COMMENT '订单数',total_revenue DOUBLE COMMENT '总收入',avg_order_value DOUBLE COMMENT '客单价',unique_customers BIGINT COMMENT '去重用户数') COMMENT '按州维度的收入汇总'""")(df_state_revenue.select("state", "state_code", "order_count", "total_revenue", "avg_order_value", "unique_customers").write.mode("overwrite").insertInto("`DataLakeCatalog`.`wedata_demo_db`.`quickstart_state_revenue`"))display(spark.sql("SELECT * FROM `DataLakeCatalog`.`wedata_demo_db`.`quickstart_state_revenue`"))
确认输出表中有数据后,单击编辑器顶部的 保存 ,生成一个版本快照。
提示
步骤 4:把任务串成一条工作流
工作流(Workflow)是统一调度的最小单元。本步骤把离线同步任务与 Notebook 串联成一条 DAG。
4.1 新建工作流
1. 在左侧导航选择 数据工程 > 工作流 。
2. 单击 创建工作流 ,自动进入工作流编辑器画布。
4.2 在画布上添加两个任务
添加离线数据接入任务 :
1. 在画布上单击 + 添加任务 ,任务类型选择 离线数据同步 。
2. 任务名称 :
sync_mysql2dlc。3. 任务来源 :下拉选择 步骤 2 创建的
quickstart_mysql2dlc。4. 计算资源 :选择数据接入型资源组。
5. 单击 创建任务 。
添加 Notebook 任务 :
1. 继续单击 + 添加任务 ,任务类型选择 Notebook 。
2. 任务名称 :
pyspark_transform。3. 任务来源 :工作空间 。
4. 路径 :单击输入框,选择 步骤 3 中创建的
quickstart_transform.ipynb。5. 计算资源 :选择与 Notebook 调试一致的作业集群;资源模式、资源配置保持默认即可。
6. 单击 创建任务 。
4.3 连接依赖关系
在画布中把鼠标移到
sync_mysql2dlc 节点的下侧锚点,按住左键拖动到 pyspark_transform 节点的上侧锚点,松开后生成一条依赖箭头。最终 DAG 形如:
sync_mysql2dlc ───────▶ pyspark_transform(离线数据接入) (Notebook)
单击
pyspark_transform 节点,确认右侧 上游依赖任务 中已列出 sync_mysql2dlc,运行条件 保持默认的 全部成功 。4.4 调试运行整条工作流
1. 单击画布顶部的 运行 立即触发一次调试运行。
2. 观察 工作流内任务的 执行顺序:
sync_mysql2dlc 先运行,成功后 pyspark_transform 自动开始。3. 单击运行记录后,点击任务节点在 运行日志 和 数据对账 查看明细:行数、Spark Job 状态、错误信息等。
两个任务都显示 成功 即说明 DAG 链路打通。
提示
任务配置项(任务名称、上游依赖、运行条件、参数、监控指标、告警、失败与超时策略)通用规则参考 配置任务 。
步骤 5:配置定时调度
调试通过后,把工作流配置为每天自动运行。
1. 在工作流编辑器右侧侧栏单击 调度配置 > 添加调度 。
2. 填写调度参数:
调度状态 :启动。
触发方式 :定时触发。
调度时区 :根据业务选择(如
(UTC+08:00) Beijing, Chongqing, Hong Kong SAR, Urumqi)。调度时间 :开始时间设为今日,结束时间使用默认
2099-12-31。配置方式 :常规。
调度周期 :天。
执行时间 :单时间点,填写
00:00。1. 确认 调度预览 区域显示的最近 5 次预计运行时间符合预期(第 1 次、第 2 次……第 5 次的日期与时间)。
2. 单击 确定 保存调度配置。
提示
更复杂的调度(如每小时、每周指定星期几、Cron 表达式、文件到达触发、持续运行)参考 调度与触发器 。
完成保存后,工作流即按调度自动产出周期实例。您可以随时回到工作流编辑器修改任务或调度,保存 即生效,下次调度时使用最新代码与配置。
步骤 6:在工作流运行中监控与排障
配置定时调度后,日常运维主要在 工作流运行 模块完成。
6.1 查看运行列表
1. 在左侧导航选择 工作流 > 工作流运行 。
2. 顶部按 工作流名称 筛选
quickstart_orders_etl,可看到本工作流的所有运行实例及状态趋势图。3. 单击某次 运行流 ID 进入运行详情。
6.2 在运行详情页定位问题
运行详情页提供三种视图:
视图 | 适用场景 |
图模式(DAG) | 直观查看每个任务状态,节点颜色:绿色=成功 / 红色=失败 / 蓝色=运行中 / 灰色=等待中 / 黄色=失败重试。 |
列表模式 | 表格化展示所有任务的运行 ID、类型、时长、计算资源、源文件路径。 |
甘特图 | 按时间轴展示每个任务的开始-结束区间,便于发现长尾任务与瓶颈。 |
在 DAG 中单击任意任务节点,进入 任务运行详情 页:
运行结果 :展示任务代码与产出。
运行日志 :错误码、错误信息、完整日志,支持搜索、下载、自动刷新。
任务洞察 :累计 CPU 时、扫描数据量、Shuffle 大小、输出行数等性能指标。
右上角 Spark UI :单击按钮跳转到 Spark 原生 UI 排查任务细节。
6.3 重跑失败任务
当任务因数据延迟、网络抖动或临时配置问题失败时:
1. 在工作流运行详情页中确认失败任务及其根因(参考运行日志)。
2. 修复根因(例如修改 Notebook 代码后在 Studio 中重新 保存 )。
3. 在工作流运行详情页单击 重跑 ,弹窗 DAG 中默认勾选失败节点,可选择是否包含其上游。
4. 单击 运行 触发重跑。重跑使用 任务最新代码、最新配置 。
注意
重跑时若任务配置了告警,告警通道仍会被触发。如需在排障重跑过程中静音告警,请提前在告警配置中开启免打扰。
提示
步骤 7:清理资源(可选)
如不再需要本教程产出的资源,可按以下顺序清理,避免持续计费与占用配额:
1. 暂停调度 :在工作流编辑器 调度配置 中单击 暂停调度 ,停止后续周期实例。
2. 删除运行记录 :在工作流运行列表中筛选本工作流,删除已到达终态的运行记录。
3. 删除工作流 :在工作流列表行尾单击 删除 (必须先终止所有非终态运行)。
4. 删除目标表 (可选):在 Studio 中执行 SQL,或在 Catalog 数据目录中删除:
DataLakeCatalog.wedata_demo_db.quickstart_amazon_ordersDataLakeCatalog.wedata_demo_db.quickstart_state_revenue1. 删除数据源 :在 平台管理 > 数据源管理 中删除
quickstart_mysql(如已无其他任务引用)。警告
删除工作流与运行记录后不可恢复,操作前请确认无业务依赖。
验证结果
完成全部步骤后,您应当确认:
1. 离线同步任务
quickstart_mysql2dlc 已成功执行,目标表 DataLakeCatalog.wedata_demo_db.quickstart_amazon_orders 中有数据,行数与 MySQL 源一致。2. Notebook 文件
quickstart_transform.ipynb 在 Studio 中已 保存版本 ,单元格可独立运行成功。3. 汇总目标表
DataLakeCatalog.wedata_demo_db.quickstart_state_revenue 中按州维度展示了订单数、总收入、客单价、去重用户数等指标。4. 工作流
quickstart_orders_etl 在工作流列表中可见,调度状态为 启动 ,DAG 中两个任务均显示为成功。5. 工作流运行 中至少有一次状态为 成功 的运行实例,任务运行日志、洞察、Spark UI 均可正常访问。
总结
在本教程中,您完成了:
在 平台管理 > 数据源管理 中注册 MySQL 数据源并测试连通性。
在 数据接入 > 任务管理 中创建并跑通离线同步任务(MySQL → 数据湖仓)。
在 Studio 中用 PySpark Notebook 实现数据清洗、地区聚合,并把结果写回数据湖仓。
在 工作流 中把离线数据接入任务与 Notebook 任务串成一条 DAG,配置上下游依赖。
配置每天 00:00 定时调度。
在 工作流运行 中跟踪运行状态、查看日志、识别失败任务并完成一次重跑。
本条端到端链路覆盖了数据工程师日常 80% 以上的高频操作。后续您可以按业务复杂度扩展更多任务类型(SQL、Python File、Ray、For Each、IF ELSE 等)与高级能力(参数传递、告警、质量监控、CI/CD)。