帮你快速理解、总结文档立即下载

数据工程师快速上手

最近更新时间:2026-09-11 20:57:30
我的收藏
提示
预计用时 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. 配置目标端:
目标端 CatalogDataLakeCatalog
目标端 Schemawedata_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 状态变为该作业集群名,则状态变为 已连接
提示
Kernel 详细配置、内核复用、运行时环境管理参考 Notebook 运行

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`"))
确认输出表中有数据后,单击编辑器顶部的 保存 ,生成一个版本快照。
提示
Studio 默认开启自动保存草稿,但 只有手动保存生成的版本才会被工作流任务引用 。每次修改 Notebook 后请记得单击 保存 。详见 Notebook IDE 基础操作

步骤 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_orders
DataLakeCatalog.wedata_demo_db.quickstart_state_revenue
1. 删除数据源 :在 平台管理 > 数据源管理 中删除 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)。