EMR Serverless 提供 Spark、Ray、Presto、Doris、StarRocks 等多引擎全托管服务,支持通过控制台快速创建并提交各类引擎作业,本文提供三种常用作业 Demo,跟随操作指引,您可以完整体验从创建作业、提交运行到查看结果的全流程,帮助您快速上手。
基本概念
作业(Job):作业是一个完整的数据计算任务,当用户需要在 EMR Serverless 中提交一个数据计算任务时,可在 EMR 中提交作业,作业提交后,系统会根据配置自动分配计算资源、执行计算、输出结果,执行完成后自动释放计算资源。 队列(Que):是 EMR Serverless 提供的一种计算资源池,包含 CPU、内存等基础计算资源,是运行作业的必要条件,用户提交作业后,Spark、Ray 等计算组件在队列内调度并执行作业,平台自动为作业调配计算资源。
EMR 为您提供了不同类型的计算资源作为底层计算资源支持,您可按需选用:标准计算资源、内存增强型 CPU 计算资源、计算增强型 CPU 计算资源、GPU 加速计算资源(购买 GPU 资源时,您需先联系销售锁定库存)。 EMR 为您提供了多种队列类型:公共队列、独占队列-包年包月(可按需增购弹性资源)、独占队列-按时计费,便于您灵活地基于不同的业务场景选择合适的资源类型、灵活选择计费模式,提高资源利用率降低成本。
前提条件
首次使用 EMR Serverless 前,请跟随准备工作 指引,完成火山引擎账号注册、实名认证、及 EMR 服务开通。EMR 主账号权限较大,建议根据最小权限原则,创建子用户并授予权限。 队列是 EMR Serverless 作业运行所依赖的资源池,使用 EMR Serverless 时您需要准备好对应的队列资源,本次入门实践使用 公共队列。说明
公共队列为开通 EMR Serverless 后为您自动创建好的公用队列资源,实际使用时,您可以按需创建其他队列资源,操作详情请参见创建队列 。
运行作业时,通常您需要将作业代码包上传到 TOS(Torch Object Storage,对象存储)中,后续直接读取 TOS 中的作业运行,因此您需要准备可访问 TOS 服务的 AK/SK。由于主账号秘钥权限很大,而本次快速入门中,仅需访问 TOS 服务的密钥,请创建仅具备 TOS 操作权限的子账号,并使用其密钥。详细操作请参见IAM访问密钥 。说明
PySpark、 RayJob 作业的代码文件需提前上传至 TOS(对象存储)存储桶,作业运行时系统将从指定路径读取文件。
Spark SQL类型直接从控制台SQL编辑器提交作业,无需准备 TOS 资源,可跳过准备工作。
操作步骤
本文提供以下几种作业 Demo 作为示例,为您介绍如何快速上手操作 EMR Serverless。
RayJob:基于 Ray 分布式框架的 Python 作业。 PySpark: Python 编写的 Spark 作业。 Spark SQL:直接在SQL 编辑器中编写 SQL 语句,进行数据查询与处理。
准备工作:准备作业 Demo 本次入门操作为您准备了以下目标作业的示例 Demo 文件,您也可以使用自己的作业代码文件,将作业代码文件上传到 TOS,便于后续在 EMR Serverless 中创建作业时直接引用。
获取作业 Demo。
作业类型
示例 Demo 文件
示例作业代码
RayJob
import ray
import time
# 初始化Ray
ray.init()
# 定义一个远程函数
@ray.remote
def square(x):
time.sleep(1) # 模拟一个耗时的计算
return x * x
# 创建一个任务列表
numbers = [1, 2, 3, 4, 5]
# 使用Ray并行计算平方值
results = ray.get([square.remote(num) for num in numbers])
# 输出结果
print("Squares:", results)
# 关闭Ray
ray.shutdown()
PySpark
py_job.py 代码:
"""
主入口脚本:Spark 作业测试
"""
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql import types as T
from pyspark.sql import SQLContext
from models import Bean
# Spark 作业名称
job_name = 'py_job_test'
# 初始化 SparkSession
spark = SparkSession.builder.appName(job_name).getOrCreate()
# 创建 Bean 对象,传入计数值,验证 Spark 作业中能否正确导入和使用自定义 Python 模块/类
bean = Bean(20)
print(bean.count)
# 执行 SQL 查询,测试 Spark SQL 引擎是否正常工作
spark.sql("select 10000 as test_value").show(10)
# 停止 Spark 会话
spark.stop()
models.py 代码:
class Bean(object):
count = 0
def __init__(self, count):
self.count = count
创建 TOS 存储桶并上传作业 Demo 到 TOS Bucket。
登录对象存储控制台 ,点击控制台左侧 桶列表 ,进入存储桶列表界面后点击左上角 创建桶, 在 创建存储桶 界面,填入配置参数。必填参数说明见下表,其余参数保持默认即可,详细参数说明请参见创建存储桶 。
参数
说明
示例
名称
存储桶的名称。创建后不可更改。
abc
区域
存储桶所归属的地域。需要与 EMR Serverless 队列在统一区域,否则无法直接互相访问。
华北2(北京)
桶策略
设置存储桶级别的访问控制策略(Bucket Policy),建议设置为私有。
私有
将 Demo 上传到存储桶。点击创建好的存储桶名称,进入存储桶详情页,点击 文件列表>上传文件 ,选择 Demo 文件,点击上传 ,上传完成后,在 文件列表 界面,复制文件夹路径,作为备用参数,路径示例:tos://abc/doctest。
进入作业入口 登录 EMR Serverless 控制台 。
在顶部菜单栏中,根据实际场景,下拉选择地域和项目空间:
地域:创建的队列及相应资源将会部署在对应的地域内,一旦创建不能修改。 项目:默认显示默认项目。详细信息请参见项目配置 。 根据目标作业类型,点击 SQL编辑器 按钮,进入 Spark SQL 提交页面;点击创建作业 按钮,进入 RayJob、PySpark 创建作业页面。
方式一:在总览界面
方式二:在左侧导航栏
单击 作业中心 > 作业管理或作业实例,点击右上角创建作业 或SQL编辑器 。
在左侧导航栏选择 Serverless,单击队列名进入队列详情页,点击右上角创建作业 或SQL编辑器 。
创建并运行作业
RayJob 在创建作业页面配置作业信息,本次 Demo 作业必填参数说明见下表。信息配置完成后,点击右下角 创建并运行 ,即可提交 RayJob 并在队列资源中运行。
参数
配置说明
作业名称
自定义填写本次作业的名称,创建后不可更改。
作业类型
选择创建的作业的类型。本示例选择 RayJob。
执行资源
执行作业使用的 EMR 队列资源。
资源类型下拉框:选择 EMR Serverless 队列/项目下拉框:本示例选择队列默认项目 default。 镜像类型
本次快速入门示例保持默认即可。
EMR Serverless 的每个作业运行在容器里,镜像是容器的"操作系统 + 依赖环境"包。
内置镜像:火山引擎提供的预制镜像,适配各种作业场景。 火山镜像仓库:用户自定义构建并推送到镜像仓库(VCR)的镜像。 镜像
点击下拉框,选择任意镜像即可。
任务主文件
作业的核心代码文件,存放在 TOS 存储桶,平台启动作业时会从 TOS 拉取该文件并提交入口命令执行。本次示例点击下拉框,选择在准备工作中已上传的作业,例如:tos://abc/doctest/ray_job.zip。
入口命令
作业容器启动后实际执行的 Shell 命令,通常为 python xxx.py。如有参数,一并加在该命令中即可。
本示例填写 python ray_job.py。
PySpark 在创建作业页面配置作业信息,本次 Demo 作业必填参数说明见下表。信息配置完成后,点击右下角 创建并运行 ,即可提交 RayJob 并在队列资源中运行。
参数
配置操作
作业名称
本次作业的名称,创建后不可更改。
作业类型
创建的作业的类型。本示例选择 PySpark。
执行资源
执行作业使用的EMR资源。
资源类型下拉框:选择 EMR Serverless 队列/项目下拉框:本示例选择队列默认项目 default。 开发模式
本次示例选择UI模式即可。
UI:通过可视化界面提交作业内容。 JSON:通过 JSON 方式提交作业内容。 Python 文件
作业的核心代码文件,存放在 TOS 存储桶,平台启动作业时会从 TOS 拉取该文件并提交入口命令执行。本次示例点击下拉框,选择在准备工作中已上传的作业,例如tos://abc/doctest/py_job.py。
依赖 Python 资源
任务文件执行时依赖的 models 文件,存放在 TOS 存储桶。本次示例点击下拉框,选择在准备工作中已上传的 models,例如 tos://abc/doctest/models.py。
Spark 参数
Spark 组件的环境变量。 此处需填写用户的 AK/SK,用于 TOS 服务鉴权,填写格式:
serverless.emr.access.key=AKL***** serverless.emr.secret.key=WkR*****
Spark SQL 进入 SQL 编辑器 页面后,在工作台提交 SQL 语句并点击 运行 :
说明
EMR SQL 编辑器默认使用 Spark SQL 执行引擎,支持从工作台提交 SQL 语句执行数据查询, 元数据服务由 Las 提供。
在数据库中创建一个表:
说明
首次进入 SQL 编辑器页面时,系统将为用户创建一个默认的 Hive 类型 Catalog ,并创建一个名为 default 的默认数据库,元数据服务由火山引擎 LAS(AI数据湖服务)提供。
CREATE TABLE default.testtb
(
`id` INT,
`name` STRING,
`age` INT
)
为数据表写入数据:
INSERT INTO default.testtb VALUES(1, 'testyw',26), (2, '张三', 26), (3, '李四', 23)
查询表数据:
SELECT * FROM default.testtb
查看作业运行结果 查看作业入口:
在控制台总览界面,单击作业中心 > 作业管理 或作业实例 ,可查看PySpark、 RayJob 作业;Spark SQL作业仅支持从 作业中心 >作业实例 页面查看。支持通过 实例名称/ID 、作业名称/ID 、开始时间 、提交人 、资源类型 筛选作业。 查看作业详细信息:
点击作业名称,进入作业详情页面,可查看系统日志、Driver日志、查询结果等信息。
RayJob 请在 作业详情>运行日志 页面查看 Demo 作业最终输出:
PySpark 请在 作业详情>Driver日志 页面,选择文件名为 stdout ,查看 Demo 作业执行结果。
Spark SQL 请在 作业详情>查询结果 页面或 SQL编辑器 工作台下方 查询结果 处,查看查询作业执行结果。
作业详情>查询结果 页面:SQL编辑器 工作台下方 查询结果 处:
更多操作
通常您也可以和 DataLeap 联合使用,大数据研发治理套件(DataLeap)是火山引擎自研的一站式大数据中台解决方案,集实时&离线数据集成、数据开发、智能运维、数据治理、资产管理能力于一身。在 DataLeap 中对接 EMR 引擎后,即可在 DataLeap 中创建 EMR Serverless 作业、查询数据等,与 DataLeap 联动的快速入门操作可参见 DataLeap on EMR Serverless Spark 快速入门 。 创建 EMR Serverless 队列时,您需选择元数据服务,当前您可选择 火山引擎 LAS Catalog 服务作为 EMR Serverless 队列资源的元数据服务,LAS Catalog 的介绍请参见 LAS 元数据 。