
# 典型场景示例：使用Spark Jar作业读取和查询OBS数据
#### 操作场景
DLI完全兼容开源的[Apache Spark](https://spark.apache.org/)，支持用户开发应用程序代码来进行作业数据的导入、查询以及分析处理。本示例从编写Spark程序代码读取和查询OBS数据、编译打包到提交Spark Jar作业等完整的操作步骤说明来帮助您在DLI上进行作业开发。
#### 环境准备
在进行Spark Jar作业开发前，请准备以下开发环境。
表1Spark Jar作业开发环境 
| 准备项                | 说明                                           |
|:---|:---|
| 操作系统               | Windows系统，支持Windows7以上版本。                    |
| 安装JDK              | JDK使用1.8版本。                                  |
| 安装和配置IntelliJ IDEA | IntelliJ IDEA为进行应用开发的工具，版本要求使用2019.1或其他兼容版本。 |
| 安装Maven            | 开发环境的基本配置。用于项目管理，贯穿软件开发生命周期。                 |
   
#### 开发流程
DLI进行Spark Jar作业开发流程参考如下：
图1Spark Jar作业开发流程   
![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001251908699.png "点击放大")
表2开发流程说明 
| 序号 | 阶段                | 操作界面                                                                    | 说明                                                                                              |
|:---|:---|:---|:---|
| 1  | 创建DLI通用队列         | DLI控制台                                                                  | 创建作业运行的DLI队列。                                                                                   |
| 2  | 上传数据到OBS桶         | OBS控制台                                                                  | 将测试数据上传到OBS桶下。                                                                                  |
| 3  | 新建Maven工程，配置pom文件 | IntelliJ IDEA                                                           |  参考样例代码说明，编写程序代码读取OBS数据。  |
| 4  | 编写程序代码            | IntelliJ IDEA                                                           |  参考样例代码说明，编写程序代码读取OBS数据。  |
| 5  | 调试，编译代码并导出Jar包    | IntelliJ IDEA                                                           |  参考样例代码说明，编写程序代码读取OBS数据。  |
| 6  | 上传Jar包到OBS和DLI    | OBS控制台 DLI控制台 | 将生成的Spark Jar包文件上传到OBS目录下和DLI程序包中。                                                              |
| 7  | 创建Spark Jar作业     | DLI控制台                                                                  | 在DLI控制台创建Spark Jar作业并提交运行作业。                                                                    |
| 8  | 查看作业运行结果          | DLI控制台                                                                  | 查看作业运行状态和作业运行日志。                                                                                |
   
 #### 步骤1：创建DLI通用队列
提交Spark作业需要先创建队列，本例创建名为"sparktest"的通用队列。
1. 登录DLI管理控制台。
2. 在左侧导航栏单击"资源管理 \> 弹性资源池"，可进入弹性资源池管理页面。
3. 在弹性资源池管理界面，单击界面右上角的"购买弹性资源池"。
4. 在"购买弹性资源池"界面，填写具体的弹性资源池参数。
   本例在华东-上海二区域购买按需计费的弹性资源池。相关参数说明如[表3]所示。
    表3参数说明 
   | 参数名称 | 参数说明                                                             | 配置样例              |
   |:---|:---|:---|
   | 计费模式 | 选择弹性资源池计费模式。                                                     | 按需计费              |
   | 区域   | 选择弹性资源池所在区域。                                                     | 华东-上海二            |
   | 项目   | 每个区域默认对应一个项目，由系统预置。                                              | 系统默认项目            |
   | 名称   | 弹性资源池名称。                                                         | dli_resource_pool |
   | 规格   | 选择弹性资源池规格。                                                       | 标准版               |
   | CU范围 | 弹性资源池最大最小CU范围。                                                   | 64-64             |
   | 网段   | 规划弹性资源池所属的网段。如需使用DLI增强型跨源，弹性资源池网段与数据源网段不能重合。**弹性资源池网段设置后不支持更改**。 | 172.16.0.0/19     |
   | 企业项目 | 选择对应的企业项目。                                                       | default           |
      
   
5. 参数填写完成后，单击"立即购买"，在界面上确认当前配置是否正确。
6. 单击"提交"完成弹性资源池的创建。
7. 在弹性资源池的列表页，选择要操作的弹性资源池，单击操作列的"添加队列"。
8. 配置队列的基础配置，具体参数信息如下。
   表4弹性资源池添加队列基础配置 
   | 参数名称 | 参数说明                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                            | 配置样例                                                                                                                |
   |:---|:---|:---|
   | 名称   | 弹性资源池添加的队列名称。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                   | dli_queue_01                                                                                                        |
   | 类型   | 选择创建的队列类型。 - 执行SQL作业请选择SQL队列。  - 执行Flink或Spark作业请选择通用队列。   | SQL作业场景请选择"SQL队列"。 其他场景请选择"通用队列"。 |
   | 执行引擎 | SQL队列可以选择队列引擎为Spark。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                            | Spark                                                                                                               |
   | 企业项目 | 选择对应的企业项目。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                      | default                                                                                                             |
      
   
9. 单击"下一步"，配置队列的扩缩容策略。 单击"新增"，可以添加不同优先级、时间段、"最小CU"和"最大CU"扩缩容策略。
   本例配置的扩缩容策略如[图2]所示。
   图2添加队列时配置扩缩容策略   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000002000460541.png "点击放大")
   表5扩缩容策略参数说明 
   | 参数名称 | 参数说明                                                                                                                           | 配置样例  |
   |:---|:---|:---|
   | 优先级  | 当前弹性资源池中的优先级数字越大表示优先级越高。本例设置一条扩缩容策略，默认优先级为1。                                                                                   | 1     |
   | 时间段  | 首条扩缩容策略是默认策略，不能删除和修改时间段配置。 即设置00-24点的扩缩容策略。 | 00-24 |
   | 最小CU | 设置扩缩容策略支持的最小CU数。                                                                                                               | 16    |
   | 最大CU | 当前扩缩容策略支持的最大CU数。                                                                                                               | 64    |
      
   
10. 单击"确定"完成添加队列配置。
 
#### 步骤2：上传数据到OBS桶
1. 根据如下数据，创建people.json文件。
   ```
   {"name":"Michael"}
   {"name":"Andy", "age":30}
   {"name":"Justin", "age":19}
   ```
   
2. 进入OBS管理控制台，在"桶列表"下，单击已创建的OBS桶名称，本示例桶名为"dli-test-obs01"。
3. 单击"上传对象"，将people.json文件上传到OBS桶根目录下。
4. 在OBS桶根目录下，单击"新建文件夹"，创建名为"result"的文件夹。
5. 单击"result"的文件夹，在"result"下单击"新建文件夹"，创建名为"parquet"的文件夹。
 
#### 步骤3：新建Maven工程，配置pom依赖
以下通过IntelliJ IDEA 2020.2工具操作演示。
1. 打开IntelliJ IDEA，选择"File \> New \> Project"。
   图3新建Project   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001252187705.png "点击放大") 
2. 选择Maven，Project SDK选择1.8，单击"Next"。
   图4新建Project   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001637557382.png "点击放大") 
3. 定义样例工程名和配置样例工程存储路径，单击"Finish"完成工程创建。
   图5创建工程   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001637398494.png "点击放大")
   如上图所示，本示例创建Maven工程名为：SparkJarObs，Maven工程路径为："D:\\DLITest\\SparkJarObs"。
   
4. 在pom.xml文件中添加如下配置。
   ```
   <dependencies>
           <dependency>
               <groupId>org.apache.spark</groupId>
               <artifactId>spark-sql_2.11</artifactId>
               <version>2.3.2</version>
           </dependency>
   </dependencies>
   ```
   图6修改pom.xml文件   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001252053711.png "点击放大") 
5. 在工程路径的"src \> main \> java"文件夹上鼠标右键，选择"New \> Package"，新建Package和类文件。
   图7新建Package   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001637399398.png "点击放大")
   Package根据需要定义，本示例定义为："com.huawei.dli.demo"，完成后回车。
   在包路径下新建Java Class文件，本示例定义为：SparkDemoObs。
   图8新建Java Class文件   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001637559882.png "点击放大") 
 
#### 步骤4：编写代码
编写SparkDemoObs程序读取OBS桶下的[1]的"people.json"文件，并创建和查询临时表"people"。
完整的样例请参考[完整样例代码参考]，样例代码分段说明如下：
1. 导入依赖的包。
   ```
   import org.apache.spark.sql.Dataset;
   import org.apache.spark.sql.Row;
   import org.apache.spark.sql.SaveMode;
   import org.apache.spark.sql.SparkSession;
   import static org.apache.spark.sql.functions.col;
   ```
   
2. 配置访问凭证。 在Spark参数(--conf)参数中配置访问DEW获取凭证。
   - 方式1：使用委托授权：使用Spark 3.3.1及以上版本、Flink 1.15版本的引擎执行作业时，需要您先在IAM页面创建相关委托，并在配置作业时添加新建的委托信息。 委托权限示例请参考[创建DLI自定义委托权限](https://support.huaweicloud.com/usermanual-dli/dli_01_0616.html)和[常见场景的委托权限策略](https://support.huaweicloud.com/usermanual-dli/dli_01_0617.html)。
     
   
   - 方式2：使用DEW授权：
     - 已为授予IAM用户所需的IAM和LakeFormation权限. 具体请参考[授权使用LakeFormation资源](https://support.huaweicloud.com/usermanual-dli/dli_01_0629.html)的操作步骤。
       
     
     - 已在DEW服务创建通用凭证，并存入凭据值。具体操作请参考[创建通用凭据](https://support.huaweicloud.com/usermanual-dew/dew_01_9993.html)。
     
     - 已创建DLI访问DEW的委托并完成委托授权。该委托需具备以下权限：
       - DEW中的查询凭据的版本与凭据值ShowSecretVersion接口权限，csms:secretVersion:get。
       
       - DEW中的查询凭据的版本列表ListSecretVersions接口权限，csms:secretVersion:list。
       
       - DEW解密凭据的权限，kms:dek:decrypt。
       
       
       [Spark Jar 使用DEW获取访问凭证读写OBS](https://support.huaweicloud.com/usermanual-dli/dli_09_0215.html#dli_09_0215)
       
      
    
3. 读取OBS桶中的"people.json"文件数据。
   其中"dli-test-obs01"为演示的OBS桶名，请根据实际的OBS桶名替换。
   ```
   Dataset<Row> df = spark.read().json("obs://dli-test-obs01/people.json");
   df.printSchema();
   ```
   
4. 通过创建临时表"people"读取文件数据。
   ```
   df.createOrReplaceTempView("people");
   ```
   
5. 查询表"people"数据。
   ```
   Dataset<Row> sqlDF = spark.sql("SELECT * FROM people");
   sqlDF.show();
   ```
   
6. 将表"people"数据以parquet格式输出到OBS桶的"result/parquet"目录下。
   ```
   sqlDF.write().mode(SaveMode.Overwrite).parquet("obs://dli-test-obs01/result/parquet");
   spark.read().parquet("obs://dli-test-obs01/result/parquet").show();
   ```
   
7. 关闭SparkSession会话spark。
   ```
   spark.stop();
   ```
   
 
#### 步骤5：调试、编译代码并导出Jar包
1. 双击IntelliJ IDEA工具右侧的"Maven"，参考下图分别双击"clean"、"compile"对代码进行编译。
   编译成功后，双击"package"对代码进行打包。
   图9编译打包   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001637454354.png "点击放大")
   打包成功后，生成的Jar包会放到target目录下，以备后用。本示例将会生成到："D:\\DLITest\\SparkJarObs\\target"下名为"SparkJarObs-1.0-SNAPSHOT.jar"。
   图10导出jar包   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001762766518.png "点击放大") 
 
 #### 步骤6：上传Jar包到OBS和DLI下
- **Spark 3.3及以上版本：**
  仅支持在创建Spark作业时，配置"应用程序"，从OBS选择作业所需的Jar包。
  1. 登录OBS控制台，将生成的Jar包文件上传到OBS路径下。
  
  2. 登录DLI控制台，选择"作业管理 \> Spark作业"。
  
  3. 单击操作列"编辑"。
  
  4. 编辑"应用程序"，选择[1]上传的OBS地址。
     图11配置应用程序   
     ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001842116258.png "点击放大") 
   
- **Spark 3.3以下版本：**
  分别上传Jar包到OBS和DLI下。
  1. 登录OBS控制台，将生成的Jar包文件上传到OBS路径下。
  
  2. 将Jar包文件上传到DLI的程序包管理中，方便后续统一管理。
     1. 登录DLI管理控制台，单击"数据管理 \> 程序包管理"。
     
     2. 在"程序包管理"页面，单击右上角的"创建程序包"。
     
     3. 在"创建程序包"对话框，配置以下参数。
        1. 包类型：选择"JAR"。
        
        2. OBS路径：程序包所在的OBS路径。
        
        3. 分组设置和组名称根据情况选择设置，方便后续识别和管理程序包。
         
     
     4. 单击"确定"，完成创建程序包。
        图12创建程序包   
        ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001633300054.png "点击放大") 
      
   
 
#### 步骤7：创建Spark Jar作业
1. 登录DLI控制台，单击"作业管理 \> Spark作业"。
2. 在"Spark作业"管理界面，单击"创建作业"。
3. 在作业创建界面，配置对应作业运行参数。具体说明如下：
   - 所属队列：选择已创建的DLI通用队列。例如当前选择[步骤1：创建DLI通用队列]创建的通用队列"sparktest"。
   
   - 在下拉列表中选择支持的Spark版本，推荐使用最新版本。
   
   - 作业名称（--name）：自定义Spark Jar作业运行的名称。当前定义为：SparkTestObs。
   
   - 应用程序：选择[步骤6：上传Jar包到OBS和DLI下]中上传到DLI程序包。例如当前选择为："SparkJarObs-1.0-SNAPSHOT.jar"。
   
   - 主类：格式为：程序包名+类名。例如当前为：com.huawei.dli.demo.SparkDemoObs。
   
   
   其他参数可暂不选择。
   了解更多Spark Jar作业提交说明可以参考[创建Spark作业](https://support.huaweicloud.com/usermanual-dli/dli_01_0384.html)。
   
4. 单击"执行"，提交该Spark Jar作业。在Spark作业管理界面显示已提交的作业运行状态。
   图13作业运行状态   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001638204318.png "点击放大") 
 
#### 步骤8：查看作业运行结果
1. 在Spark作业管理界面显示已提交的作业运行状态。初始状态显示为"启动中"。
2. 如果作业运行成功则作业状态显示为"已成功"，单击"操作"列"更多"下的"Driver日志"，显示当前作业运行的日志。
   图14driver日志   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001686404953.png "点击放大")
   图15"Driver日志"中的作业执行日志   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001251907299.png "点击放大") 
3. 如果作业运行成功，本示例进入OBS桶下的"result/parquet"目录，查看已生成预期的parquet文件。
4. 如果作业运行失败，单击"操作"列"更多"下的"Driver日志"，显示具体的报错日志信息，根据报错信息定位问题原因。
   例如，如下截图信息因为创建Spark Jar作业时主类名没有包含包路径，报找不到类名"SparkDemoObs"。
   图16报错信息   
   ![](https://support.huaweicloud.com/usermanual-dli/zh-cn_image_0000001686339805.png "点击放大")
   可以在"操作"列，单击"编辑"，修改"主类"参数为正确的：com.huawei.dli.demo.SparkDemoObs，单击"执行"重新运行该作业即可。
   
 
#### 后续指引
- 如果您想通过Spark Jar作业访问其他数据源，请参考《[使用Spark作业跨源访问数据源](https://support.huaweicloud.com/devg-dli/dli_09_0020.html)》。
- 如果您想通过Spark Jar作业在DLI创建数据库和表，请参考《[使用Spark作业访问DLI元数据](https://support.huaweicloud.com/devg-dli/dli_09_0176.html)》。
 
 #### 完整样例代码参考
![](https://support.huaweicloud.com/usermanual-dli/public_sys-resources/note_3.0-zh-cn.png)
认证用的access.key和secret.key硬编码到代码中或者明文存储都有很大的安全风险，建议在配置文件或者环境变量中密文存放，使用时解密，确保安全。
```
package com.huawei.dli.demo;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SaveMode;
import org.apache.spark.sql.SparkSession;
import static org.apache.spark.sql.functions.col;
public class SparkDemoObs {
    public static void main(String[] args) {
        SparkSession spark = SparkSession
                .builder()
                .config("spark.hadoop.fs.obs.access.key", "xxx")
                .config("spark.hadoop.fs.obs.secret.key", "yyy")
                .appName("java_spark_demo")
                .getOrCreate();
        // can also be used --conf to set the ak sk when submit the app
        // test json data:
        // {"name":"Michael"}
        // {"name":"Andy", "age":30}
        // {"name":"Justin", "age":19}
        Dataset<Row> df = spark.read().json("obs://dli-test-obs01/people.json");
        df.printSchema();
        // root
        // |-- age: long (nullable = true)
        // |-- name: string (nullable = true)
        // Displays the content of the DataFrame to stdout
        df.show();
        // +----+-------+
        // | age|   name|
        // +----+-------+
        // |null|Michael|
        // |  30|   Andy|
        // |  19| Justin|
        // +----+-------+
        // Select only the "name" column
        df.select("name").show();
        // +-------+
        // |   name|
        // +-------+
        // |Michael|
        // |   Andy|
        // | Justin|
        // +-------+
        // Select people older than 21
        df.filter(col("age").gt(21)).show();
        // +---+----+
        // |age|name|
        // +---+----+
        // | 30|Andy|
        // +---+----+
        // Count people by age
        df.groupBy("age").count().show();
        // +----+-----+
        // | age|count|
        // +----+-----+
        // |  19|    1|
        // |null|    1|
        // |  30|    1|
        // +----+-----+
        // Register the DataFrame as a SQL temporary view
        df.createOrReplaceTempView("people");
        Dataset<Row> sqlDF = spark.sql("SELECT * FROM people");
        sqlDF.show();
        // +----+-------+
        // | age|   name|
        // +----+-------+
        // |null|Michael|
        // |  30|   Andy|
        // |  19| Justin|
        // +----+-------+
        sqlDF.write().mode(SaveMode.Overwrite).parquet("obs://dli-test-obs01/result/parquet");
        spark.read().parquet("obs://dli-test-obs01/result/parquet").show();
        spark.stop();
    }
}
```
