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

225 lines
8.2 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
## 一、官方JDBC使用
TDengine官方提供了JDBC,Spark读取和写入均可以直接使用
### 1.1 依赖问题
1. 需要引用taos-jdbcdriver-2.0.42.jar这个包,版本上强烈推荐2.0.42,其他版本会有各种问题;
2. taos-jdbc的依赖与南京spark集群存在冲突,需要手动引入依赖,方法为:
1. 下载正确的依赖,并上传hadoop;
2. 在提交spark任务的时候,手动制定依赖:
```shell
--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
```
```pom
<dependency>
<groupId>com.taosdata.jdbc</groupId>
<artifactId>taos-jdbcdriver</artifactId>
<version>2.0.42</version>
</dependency>
```
3. 对于scala可以在pom中引入依赖,并将所有依赖打包到jar包中,但是python做不到,所以对于pyspark程序,最好引用taos-jdbcdriver-2.0.42-dist.jar
```shell
--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. 读取整个表
```Python
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()
```
2. 读取部分内容
```python
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保存数据
```python
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方法,创建自定义的数据源;
```scala
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()))
}
```
2. 继承 DataSourceWriter 重写 createWriterFactory 方法并返回自定义的 DataWriterFactory,重写 commit 方法,用来提交整个事务, 重写 abort 方法,用来做事务回滚;
```scala
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
}
```
3. 继承 DataWriterFactory, 重写 createDataWriter方法返回自定义的 DataWriter
```scala
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)
}
}
```
4. 继承 DataWriter 重写 write 方法实现具体的写入数据库逻辑,重写 commit 方法用来提交事务,重写 abort 方法用来做事务回滚 ;
```scala
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
}
```
5. 调用,在程序中指定自定义的datasource并传入参数
```scala
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](https://www.modb.pro/db/132365)
2. [Spark DataSource V1 & V2 API 一文理解](https://blog.csdn.net/penriver/article/details/115672072)
3. [Spark SQL DataSource V2 学习入门 + 代码模板](http://www.jsledd.cn/2019/04/05/datasourcev2/)
4. [Tdengine v2.6 官方文档](https://docs.taosdata.com/2.6/reference/connector/java/)