
# scala样例代码
#### 开发说明
支持对接CloudTable的OpenTSDB和MRS的OpenTSDB。
- 前提条件 在DLI管理控制台上已完成创建跨源连接。具体操作请参考《[数据湖探索用户指南](https://support.huaweicloud.com/usermanual-dli/dli_01_0426.html)》。
  ![](https://support.huaweicloud.com/devg-dli/public_sys-resources/note_3.0-zh-cn.png)
  认证用的password硬编码到代码中或者明文存储都有很大的安全风险，建议在配置文件或者环境变量中密文存放，使用时解密，确保安全。
  
- 构造依赖信息，创建SparkSession
  1. 导入依赖。
     涉及到mvn依赖
     ```
     <dependency>
       <groupId>org.apache.spark</groupId>
       <artifactId>spark-sql_2.11</artifactId>
       <version>2.3.2</version>
     </dependency>
     ```
     import相关依赖包
     ```
     import scala.collection.mutable
     import org.apache.spark.sql.{Row, SparkSession}
     import org.apache.spark.rdd.RDD
     import org.apache.spark.sql.types._
     ```
     
  
  2. 创建会话。
     ```
     val sparkSession = SparkSession.builder().getOrCreate()
     ```
     
  
  3. 创建DLI关联跨源访问 OpenTSDB的关联表。
     ```
     sparkSession.sql("create table opentsdb_test using opentsdb options(
     'Host'='opentsdb-3xcl8dir15m58z3.cloudtable.com:4242',
             'metric'='ctopentsdb',
     'tags'='city,location')")
     ```
      表1创建表参数 
     | 参数     | 说明                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                   |
     |:---|:---|
     | host   | OpenTSDB连接地址。 - 访问CloudTable OpenTSDB，填写OpenTSDB链接地址，具体可以登录CloudTable控制台，单击"集群模式 \> *集群名称*"，在集群信息获取OpenTSDB链接地址。  - 访问MRS OpenTSDB，若使用增强型跨源连接，填写OpenTSDB所在节点IP与端口，格式为"IP:PORT"，OpenTSDB存在多个节点时，用分号隔开，获取方式请参考"图 MRS集群OpenTSDB IP信息"和"图 MRS集群OpenTSDB 端口信息"。若使用经典型跨源，填写经典型跨源返回的连接地址，管理控制台操作请参考《数据湖探索用户指南》。   |
     | metric | 所创建的dli表对应的OpenTSDB中的指标名称。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                           |
     | tags   | metric对应的标签，用于归类、过滤、快速检索等操作，可以是1到8个，以"，"分隔，包括对应metric下的所有tagk的值。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                     |
        
     
   
- 通过SQL API访问
  1. 插入数据
     ```
     sparkSession.sql("insert into opentsdb_test values('futian', 'abc', '1970-01-02 18:17:36', 30.0)")
     ```
     
  
  2. 查询数据
     ```
     sparkSession.sql("select * from opentsdb_test").show()
     ```
     返回结果：
     ![](https://support.huaweicloud.com/devg-dli/zh-cn_image_0223997429.png)
     
   
- 通过DataFrame API访问
  1. 构造schema
     ```
     val attrTag1Location = new StructField("location", StringType)
     val attrTag2Name = new StructField("name", StringType)
     val attrTimestamp = new StructField("timestamp", LongType)
     val attrValue = new StructField("value", DoubleType)
     val attrs = Array(attrTag1Location, attrTag2Name, attrTimestamp, attrValue)
     ```
     
  
  2. 根据schema的类型构造数据
     ```
     val mutableRow: Seq[Any] = Seq("aaa", "abc", 123456L, 30.0)
     val rddData: RDD[Row] = sparkSession.sparkContext.parallelize(Array(Row.fromSeq(mutableRow)), 1)
     ```
     
  
  3. 导入数据到OpenTSDB
     ```
     sparkSession.createDataFrame(rddData, new StructType(attrs)).write.insertInto("opentsdb_test")
     ```
     
  
  4. 读取OpenTSDB上的数据
     ```
     val map = new mutable.HashMap[String, String]()
     map("metric") = "ctopentsdb"
     map("tags") = "city,location"
     map("Host") = "opentsdb-3xcl8dir15m58z3.cloudtable.com:4242"
     sparkSession.read.format("opentsdb").options(map.toMap).load().show()
     ```
     返回结果：
     ![](https://support.huaweicloud.com/devg-dli/zh-cn_image_0223997490.png)
     
   
- 提交Spark作业
  1. 将写好的代码生成jar包，上传至OBS桶中。
  
  2. 在Spark作业编辑器中选择对应的Module模块并执行Spark作业。
     ![](https://support.huaweicloud.com/devg-dli/public_sys-resources/note_3.0-zh-cn.png)
     - 如果选择spark版本为2.3.2（即将下线）或2.4.5提交作业时，需要指定Module模块，名称为：sys.datasource.opentsdb。
     
     - 如果选择Spark版本为3.1.1及以上版本时，无需选择Module模块， 需在 'Spark参数（--conf)' 配置 spark.driver.extraClassPath=/usr/share/extension/dli/spark-jar/datasource/opentsdb/\*
       spark.executor.extraClassPath=/usr/share/extension/dli/spark-jar/datasource/opentsdb/\*
       
     
     - 通过控制台提交作业请参考《[数据湖探索用户指南](https://support.huaweicloud.com/usermanual-dli/dli_01_0384.html)》中的"选择依赖资源参数说明"。
     
     - 通过API提交作业请参考《数据湖探索API参考》\>《[创建批处理作业](https://support.huaweicloud.com/api-dli/dli_02_0124.html)》中"表2-请求参数说明"关于"modules"参数的说明。
       
   
 
#### 完整示例代码
- Maven依赖
  ```
  <dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-sql_2.11</artifactId>
    <version>2.3.2</version>
  </dependency>
  ```
  
- 通过SQL API访问
  ```
  import org.apache.spark.sql.SparkSession
  object Test_OpenTSDB_CT {
    def main(args: Array[String]): Unit = {
      // Create a SparkSession session.
      val sparkSession = SparkSession.builder().getOrCreate()
      // Create a data table for DLI association OpenTSDB
      sparkSession.sql("create table opentsdb_test using opentsdb options(
  'Host'='opentsdb-3xcl8dir15m58z3.cloudtable.com:4242',
  'metric'='ctopentsdb',
  'tags'='city,location')")
      //*****************************SQL module***********************************
      sparkSession.sql("insert into opentsdb_test values('futian', 'abc', '1970-01-02 18:17:36', 30.0)")
      sparkSession.sql("select * from opentsdb_test").show()
      sparkSession.close()
    }
  }
  ```
  
- 通过DataFrame API访问
  ```
  import scala.collection.mutable
  import org.apache.spark.sql.{Row, SparkSession}
  import org.apache.spark.rdd.RDD
  import org.apache.spark.sql.types._
  object Test_OpenTSDB_CT {
    def main(args: Array[String]): Unit = {
      // Create a SparkSession session.
      val sparkSession = SparkSession.builder().getOrCreate()
      // Create a data table for DLI association OpenTSDB
      sparkSession.sql("create table opentsdb_test using opentsdb options(
  'Host'='opentsdb-3xcl8dir15m58z3.cloudtable.com:4242',
  'metric'='ctopentsdb',
  'tags'='city,location')")
      //*****************************DataFrame model***********************************
      // Setting schema
      val attrTag1Location = new StructField("location", StringType)
      val attrTag2Name = new StructField("name", StringType)
      val attrTimestamp = new StructField("timestamp", LongType)
      val attrValue = new StructField("value", DoubleType)
      val attrs = Array(attrTag1Location, attrTag2Name, attrTimestamp,attrValue)
      // Populate data according to the type of schema
      val mutableRow: Seq[Any] = Seq("aaa", "abc", 123456L, 30.0)
      val rddData: RDD[Row] = sparkSession.sparkContext.parallelize(Array(Row.fromSeq(mutableRow)), 1)
      //Import the constructed data into OpenTSDB
      sparkSession.createDataFrame(rddData, new StructType(attrs)).write.insertInto("opentsdb_test")
      //Read data on OpenTSDB
      val map = new mutable.HashMap[String, String]()
      map("metric") = "ctopentsdb"
      map("tags") = "city,location"
      map("Host") = "opentsdb-3xcl8dir15m58z3.cloudtable.com:4242"
      sparkSession.read.format("opentsdb").options(map.toMap).load().show()
      sparkSession.close()
    }
  }
  ```
  
 
