Files
obsidian_sycn/知识管理/编程经验/TDengine/Spark与TDengine链接经验总结.md
2026-07-31 15:28:08 +08:00

8.2 KiB
Raw Permalink Blame History

一、官方JDBC使用

TDengine官方提供了JDBC,Spark读取和写入均可以直接使用

1.1 依赖问题

  1. 需要引用taos-jdbcdriver-2.0.42.jar这个包,版本上强烈推荐2.0.42,其他版本会有各种问题;

  2. taos-jdbc的依赖与南京spark集群存在冲突,需要手动引入依赖,方法为:

    1. 下载正确的依赖,并上传hadoop
    2. 在提交spark任务的时候,手动制定依赖:
--conf spark.driver.extraClassPath=guava-30.1.1-jre.jar:failureaccess-1.0.1.jar \
--conf spark.executor.extraClassPath=guava-30.1.1-jre.jar:failureaccess-1.0.1.jar 
--jars path-to-jar/guava-30.1.1-jre.jar,path-to-jar/failureaccess-1.0.1.jar 
<dependency>  
  <groupId>com.taosdata.jdbc</groupId>  
  <artifactId>taos-jdbcdriver</artifactId>  
  <version>2.0.42</version>  
</dependency>
  1. 对于scala可以在pom中引入依赖,并将所有依赖打包到jar包中,但是python做不到,所以对于pyspark程序,最好引用taos-jdbcdriver-2.0.42-dist.jar
--conf spark.driver.extraClassPath=guava-30.1.1-jre.jar:failureaccess-1.0.1.jar:taos-jdbcdriver-2.0.42-dist.jar 
--conf spark.executor.extraClassPath=guava-30.1.1-jre.jar:failureaccess-1.0.1.jar:taos-jdbcdriver-2.0.42-dist.jar 
--jars path-to-jar/guava-30.1.1-jre.jar,path-to-jar/failureaccess-1.0.1.jar,path-to-jar/taos-jdbcdriver-2.0.42-dist.jar

1.2 读取

可以使用下面的方法读取TDegine

  1. 读取整个表
info = spark.read\  
    .format("jdbc")\  
    .option("driver", "com.taosdata.jdbc.rs.RestfulDriver")\  
    .option("url", "jdbc:TAOS-RS://ip:port/db?user=user&password=password")\  
     #这里是设置网络超时的,如果一次性读取的信息比较多可以把这里设置大一些,单位是ms,默认为5000
    .option("httpSocketTimeout", "100000") \ 
    .option("dbtable", "test.test")\  
    .load() 
  1. 读取部分内容
info = spark.read\  
   .format("jdbc")\  
   .option("driver", "com.taosdata.jdbc.rs.RestfulDriver")\  
   .option("url", "jdbc:TAOS-RS://ip:port/db?user=user&password=password")\  
    #这里是设置网络超时的,如果一次性读取的信息比较多可以把这里设置大一些,单位是ms,默认为5000
   .option("httpSocketTimeout", "100000") \ 
   .option("query", "select * from test.test limit 100")\  
   .load() 

1.3 存储

可以使用下面的方法向TDegine保存数据

  toHbaseSave
    .write
    .format("jdbc")
    .option("url", url)
    .option("driver", driver)
    .mode(SaveMode.Append)
    .option("dbtable", "test.test")
    .save()

二、使用自定义数据源写入数据

对于spark来说,直接使用官方的jdbc进行连接只能实现向普通表中写入数据,无法利用TDengine超级表的特性,为了能向超级表中写入数据,需要自定义数据源。由于只需要写入特性,因此下面的内容中只实现了写入的逻辑,没有实现读取的逻辑。

  1. extends DataSourceV2 with WriteSupport 并重写createWriter方法,创建自定义的数据源;
class TDSourceV2 extends DataSourceV2  with WriteSupport with Serializable {  
  
  override def createWriter(
                             jobId: String,  
                             structType: StructType,  
                             saveMode: SaveMode,  
                             dataSourceOptions: DataSourceOptions): Optional[DataSourceWriter] = {  
    Optional.of(new TDSourceWriter(  
	//地址
      dataSourceOptions.get("url").get(),
    //用户名  
      dataSourceOptions.get("user").get(),  
	//密码
      dataSourceOptions.get("password").get(),  
	//数据库名称
      dataSourceOptions.get("db").get(),  
	//超级表名称
      dataSourceOptions.get("stable").get()))  
  }
  1. 继承 DataSourceWriter 重写 createWriterFactory 方法并返回自定义的 DataWriterFactory,重写 commit 方法,用来提交整个事务, 重写 abort 方法,用来做事务回滚;
class TDSourceWriter(  
                      url: String,  
                      user: String,  
                      password: String,  
                      db: String,  
                      stable: String) extends DataSourceWriter with Serializable{  
                      
  override def createWriterFactory(): DataWriterFactory[InternalRow] = {  
    new TDWriterFactory(url, user, password, db, stable)  
  }  
  // 2.6版本的TDengine不支持事务
  override def commit(writerCommitMessages: Array[WriterCommitMessage]): Unit = Unit  
  // 2.6版本的TDengine不支持事务
  override def abort(writerCommitMessages: Array[WriterCommitMessage]): Unit = Unit  
}
  1. 继承 DataWriterFactory, 重写 createDataWriter方法返回自定义的 DataWriter
class TDWriterFactory(  
                       url: String,  
                       user: String,  
                       password: String,  
                       db: String,  
                       stable: String) extends DataWriterFactory[InternalRow] with Serializable {  
  override def createDataWriter(  
                                 partitionId: Int,  
                                 taskId: Long,  
                                 epochId: Long): DataWriter[InternalRow] = {  
    new TDDataWriter(url, user, password, db, stable)  
  }  
}
  1. 继承 DataWriter 重写 write 方法实现具体的写入数据库逻辑,重写 commit 方法用来提交事务,重写 abort 方法用来做事务回滚 ;
class TDDataWriter(  
                    url: String,  
                    user: String,  
                    password: String,  
                    db: String,  
                    stable: String) extends DataWriter[InternalRow] with Serializable {  
  
  private val logger = LoggerFactory.getLogger(this.getClass)  
  private var conn: Connection = null  
  private var stmt: Statement = null  
  
  override def write(record: InternalRow): Unit = {  
  
	//在这里编写将数据插入数据库的逻辑
	
    Class.forName("com.taosdata.jdbc.rs.RestfulDriver")  
    val jdbcUrl = s"jdbc:TAOS-RS://$url/$db?user=$user&password=$password"  
    
	// 可以通过record.getxxx获取对应列的内容
	// 如table_name = record.getString(0) 
	
    val variable1 = record.getString(0)  
    val variable2 = record.getInt(1)
    val variable3 = record.getFloat(2)
        
	// 构建插入数据库的sql语句
	// 请注意,只有原生连接方式才支持参数绑定的插入方式,
	// 由于要在spark集群上跑,这里用的是Rest的连接方式,
	// 所以只能手动构建插入数据库的sql语句
	// 关于原生连接、Rest连接和参数绑定插入,请见官方文档
	
    val sql = "insert to ..."
  
    logger.info(sql)  
  
    conn = DriverManager.getConnection(jdbcUrl)  
    stmt = conn.createStatement()  
    stmt.execute(s"use $db")  
  
    try{  
      stmt.execute(sql)  
      conn.commit()  
    }catch {  
      case e: Exception => e.printStackTrace()  
    }finally {  
      conn.close()  
    }  
  }  

  // Tdengine2.6版本不支持事务,因此此处并无实际需要做的事情
  // 创建WriterCommitMessage类,绝对不能传null,网上有些代码是错误的
  object WriteSucceeded extends WriterCommitMessage  
  override def commit(): WriterCommitMessage = WriteSucceeded 
   
  // Tdengine2.6版本不支持事务,因此此处并无实际需要做的事情
  override def abort(): Unit = Unit  
}
  1. 调用,在程序中指定自定义的datasource并传入参数
toTdSave  
  .write  
  .format("org.example.TDSourceV2")  
  .mode(SaveMode.Append)  
  // 地址
  .option("url", url)  
  // 数据库名
  .option("db", db)  
  // 用户名
  .option("user", user)  
  // 超级表名
  .option("stable", stable)  
  // 密码
  .option("password", password)  
  .save()

参考文献:

  1. 暑期2021项目经验分享:实现Spark对接openGauss
  2. Spark DataSource V1 & V2 API 一文理解
  3. Spark SQL DataSource V2 学习入门 + 代码模板
  4. Tdengine v2.6 官方文档