TDengine Go Connector

driver-go 是 TDengine 的官方 Go 语言连接器,实现了 Go 语言 database/sql 包的接口。Go 开发人员可以通过它开发存取 TDengine 集群数据的应用软件。

driver-go 提供两种建立连接的方式。一种是原生连接,它通过 TDengine 客户端驱动程序(taosc)原生连接 TDengine 运行实例,支持数据写入、查询、订阅、schemaless 接口和参数绑定接口等功能。另外一种是 REST 连接,它通过 taosAdapter 提供的 REST 接口连接 TDengine 运行实例。REST 连接实现的功能特性集合和原生连接有少量不同。

本文介绍如何安装 driver-go,并通过 driver-go 连接 TDengine 集群、进行数据查询、数据写入等基本操作。

driver-go 的源码托管在 GitHub

支持的平台

原生连接支持的平台和 TDengine 客户端驱动支持的平台一致。 REST 连接支持所有能运行 Go 的平台。

版本支持

请参考版本支持列表

支持的功能特性

原生连接

“原生连接”指连接器通过 TDengine 客户端驱动(taosc)直接与 TDengine 运行实例建立的连接。支持的功能特性有:

  • 普通查询
  • 连续查询
  • 订阅
  • schemaless 接口
  • 参数绑定接口

REST 连接

“REST 连接”指连接器通过 taosAdapter 组件提供的 REST API 与 TDengine 运行实例建立的连接。支持的功能特性有:

  • 普通查询
  • 连续查询

安装步骤

安装前准备

  • 安装 Go 开发环境(Go 1.14 及以上,GCC 4.8.5 及以上)
  • 如果使用原生连接器,请安装 TDengine 客户端驱动,具体步骤请参考安装客户端驱动

配置好环境变量,检查命令:

  • go env
  • gcc -v

使用 go get 安装

go get -u github.com/taosdata/driver-go/v3@latest

使用 go mod 管理

  1. 使用 go mod 命令初始化项目:

    1. go mod init taos-demo
  2. 引入 taosSql :

    1. import (
    2. "database/sql"
    3. _ "github.com/taosdata/driver-go/v3/taosSql"
    4. )
  3. 使用 go mod tidy 更新依赖包:

    1. go mod tidy
  4. 使用 go run taos-demo 运行程序或使用 go build 命令编译出二进制文件。

    1. go run taos-demo
    2. go build

建立连接

数据源名称(DSN)

数据源名称具有通用格式,例如 PEAR DB,但没有类型前缀(方括号表示可选):

  1. [username[:password]@][protocol[(address)]]/[dbname][?param1=value1&...&paramN=valueN]

完整形式的 DSN:

  1. username:password@protocol(address)/dbname?param=value

使用连接器进行连接

  • 原生连接
  • REST 连接
  • WebSocket 连接

taosSql 通过 cgo 实现了 Go 的 database/sql/driver 接口。只需要引入驱动就可以使用 database/sql 的接口。

使用 taosSql 作为 driverName 并且使用一个正确的 DSN 作为 dataSourceName,DSN 支持的参数:

  • configPath 指定 taos.cfg 目录

示例:

  1. package main
  2. import (
  3. "database/sql"
  4. "fmt"
  5. _ "github.com/taosdata/driver-go/v3/taosSql"
  6. )
  7. func main() {
  8. var taosUri = "root:taosdata@tcp(localhost:6030)/"
  9. taos, err := sql.Open("taosSql", taosUri)
  10. if err != nil {
  11. fmt.Println("failed to connect TDengine, err:", err)
  12. return
  13. }
  14. }

taosRestful 通过 http client 实现了 Go 的 database/sql/driver 接口。只需要引入驱动就可以使用database/sql的接口。

使用 taosRestful 作为 driverName 并且使用一个正确的 DSN 作为 dataSourceName,DSN 支持的参数:

  • disableCompression 是否接受压缩数据,默认为 true 不接受压缩数据,如果传输数据使用 gzip 压缩设置为 false。
  • readBufferSize 读取数据的缓存区大小默认为 4K(4096),当查询结果数据量多时可以适当调大该值。

示例:

  1. package main
  2. import (
  3. "database/sql"
  4. "fmt"
  5. _ "github.com/taosdata/driver-go/v3/taosRestful"
  6. )
  7. func main() {
  8. var taosUri = "root:taosdata@http(localhost:6041)/"
  9. taos, err := sql.Open("taosRestful", taosUri)
  10. if err != nil {
  11. fmt.Println("failed to connect TDengine, err:", err)
  12. return
  13. }
  14. }

taosWS 通过 WebSocket 实现了 Go 的 database/sql/driver 接口。只需要引入驱动(driver-go 最低版本 3.0.2)就可以使用database/sql的接口。

使用 taosWS 作为 driverName 并且使用一个正确的 DSN 作为 dataSourceName,DSN 支持的参数:

  • writeTimeout 通过 WebSocket 发送数据的超时时间。
  • readTimeout 通过 WebSocket 接收响应数据的超时时间。

示例:

  1. package main
  2. import (
  3. "database/sql"
  4. "fmt"
  5. _ "github.com/taosdata/driver-go/v3/taosWS"
  6. )
  7. func main() {
  8. var taosUri = "root:taosdata@ws(localhost:6041)/"
  9. taos, err := sql.Open("taosWS", taosUri)
  10. if err != nil {
  11. fmt.Println("failed to connect TDengine, err:", err)
  12. return
  13. }
  14. }

使用示例

写入数据

SQL 写入

  1. package main
  2. import (
  3. "database/sql"
  4. "fmt"
  5. "log"
  6. _ "github.com/taosdata/driver-go/v3/taosRestful"
  7. )
  8. func createStable(taos *sql.DB) {
  9. _, err := taos.Exec("CREATE DATABASE power")
  10. if err != nil {
  11. log.Fatalln("failed to create database, err:", err)
  12. }
  13. _, err = taos.Exec("CREATE STABLE power.meters (ts TIMESTAMP, current FLOAT, voltage INT, phase FLOAT) TAGS (location BINARY(64), groupId INT)")
  14. if err != nil {
  15. log.Fatalln("failed to create stable, err:", err)
  16. }
  17. }
  18. func insertData(taos *sql.DB) {
  19. sql := `INSERT INTO power.d1001 USING power.meters TAGS('California.SanFrancisco', 2) VALUES ('2018-10-03 14:38:05.000', 10.30000, 219, 0.31000) ('2018-10-03 14:38:15.000', 12.60000, 218, 0.33000) ('2018-10-03 14:38:16.800', 12.30000, 221, 0.31000)
  20. power.d1002 USING power.meters TAGS('California.SanFrancisco', 3) VALUES ('2018-10-03 14:38:16.650', 10.30000, 218, 0.25000)
  21. power.d1003 USING power.meters TAGS('California.LosAngeles', 2) VALUES ('2018-10-03 14:38:05.500', 11.80000, 221, 0.28000) ('2018-10-03 14:38:16.600', 13.40000, 223, 0.29000)
  22. power.d1004 USING power.meters TAGS('California.LosAngeles', 3) VALUES ('2018-10-03 14:38:05.000', 10.80000, 223, 0.29000) ('2018-10-03 14:38:06.500', 11.50000, 221, 0.35000)`
  23. result, err := taos.Exec(sql)
  24. if err != nil {
  25. log.Fatalln("failed to insert, err:", err)
  26. }
  27. rowsAffected, err := result.RowsAffected()
  28. if err != nil {
  29. log.Fatalln("failed to get affected rows, err:", err)
  30. }
  31. fmt.Println("RowsAffected", rowsAffected)
  32. }
  33. func main() {
  34. var taosDSN = "root:taosdata@http(localhost:6041)/"
  35. taos, err := sql.Open("taosRestful", taosDSN)
  36. if err != nil {
  37. log.Fatalln("failed to connect TDengine, err:", err)
  38. }
  39. defer taos.Close()
  40. createStable(taos)
  41. insertData(taos)
  42. }

查看源码

InfluxDB 行协议写入

  1. package main
  2. import (
  3. "fmt"
  4. "log"
  5. "github.com/taosdata/driver-go/v3/af"
  6. )
  7. func prepareDatabase(conn *af.Connector) {
  8. _, err := conn.Exec("CREATE DATABASE test")
  9. if err != nil {
  10. panic(err)
  11. }
  12. _, err = conn.Exec("USE test")
  13. if err != nil {
  14. panic(err)
  15. }
  16. }
  17. func main() {
  18. conn, err := af.Open("localhost", "root", "taosdata", "", 6030)
  19. if err != nil {
  20. fmt.Println("fail to connect, err:", err)
  21. }
  22. defer conn.Close()
  23. prepareDatabase(conn)
  24. var lines = []string{
  25. "meters,location=California.LosAngeles,groupid=2 current=11.8,voltage=221,phase=0.28 1648432611249",
  26. "meters,location=California.LosAngeles,groupid=2 current=13.4,voltage=223,phase=0.29 1648432611250",
  27. "meters,location=California.LosAngeles,groupid=3 current=10.8,voltage=223,phase=0.29 1648432611249",
  28. "meters,location=California.LosAngeles,groupid=3 current=11.3,voltage=221,phase=0.35 1648432611250",
  29. }
  30. err = conn.InfluxDBInsertLines(lines, "ms")
  31. if err != nil {
  32. log.Fatalln("insert error:", err)
  33. }
  34. }

查看源码

OpenTSDB Telnet 行协议写入

  1. package main
  2. import (
  3. "log"
  4. "github.com/taosdata/driver-go/v3/af"
  5. )
  6. func prepareDatabase(conn *af.Connector) {
  7. _, err := conn.Exec("CREATE DATABASE test")
  8. if err != nil {
  9. panic(err)
  10. }
  11. _, err = conn.Exec("USE test")
  12. if err != nil {
  13. panic(err)
  14. }
  15. }
  16. func main() {
  17. conn, err := af.Open("localhost", "root", "taosdata", "", 6030)
  18. if err != nil {
  19. log.Fatalln("fail to connect, err:", err)
  20. }
  21. defer conn.Close()
  22. prepareDatabase(conn)
  23. var lines = []string{
  24. "meters.current 1648432611249 10.3 location=California.SanFrancisco groupid=2",
  25. "meters.current 1648432611250 12.6 location=California.SanFrancisco groupid=2",
  26. "meters.current 1648432611249 10.8 location=California.LosAngeles groupid=3",
  27. "meters.current 1648432611250 11.3 location=California.LosAngeles groupid=3",
  28. "meters.voltage 1648432611249 219 location=California.SanFrancisco groupid=2",
  29. "meters.voltage 1648432611250 218 location=California.SanFrancisco groupid=2",
  30. "meters.voltage 1648432611249 221 location=California.LosAngeles groupid=3",
  31. "meters.voltage 1648432611250 217 location=California.LosAngeles groupid=3",
  32. }
  33. err = conn.OpenTSDBInsertTelnetLines(lines)
  34. if err != nil {
  35. log.Fatalln("insert error:", err)
  36. }
  37. }

查看源码

OpenTSDB JSON 行协议写入

  1. package main
  2. import (
  3. "log"
  4. "github.com/taosdata/driver-go/v3/af"
  5. )
  6. func prepareDatabase(conn *af.Connector) {
  7. _, err := conn.Exec("CREATE DATABASE test")
  8. if err != nil {
  9. panic(err)
  10. }
  11. _, err = conn.Exec("USE test")
  12. if err != nil {
  13. panic(err)
  14. }
  15. }
  16. func main() {
  17. conn, err := af.Open("localhost", "root", "taosdata", "", 6030)
  18. if err != nil {
  19. log.Fatalln("fail to connect, err:", err)
  20. }
  21. defer conn.Close()
  22. prepareDatabase(conn)
  23. payload := `[{"metric": "meters.current", "timestamp": 1648432611249, "value": 10.3, "tags": {"location": "California.SanFrancisco", "groupid": 2}},
  24. {"metric": "meters.voltage", "timestamp": 1648432611249, "value": 219, "tags": {"location": "California.LosAngeles", "groupid": 1}},
  25. {"metric": "meters.current", "timestamp": 1648432611250, "value": 12.6, "tags": {"location": "California.SanFrancisco", "groupid": 2}},
  26. {"metric": "meters.voltage", "timestamp": 1648432611250, "value": 221, "tags": {"location": "California.LosAngeles", "groupid": 1}}]`
  27. err = conn.OpenTSDBInsertJsonPayload(payload)
  28. if err != nil {
  29. log.Fatalln("insert error:", err)
  30. }
  31. }

查看源码

查询数据

  1. package main
  2. import (
  3. "database/sql"
  4. "log"
  5. "time"
  6. _ "github.com/taosdata/driver-go/v3/taosRestful"
  7. )
  8. func main() {
  9. var taosDSN = "root:taosdata@http(localhost:6041)/power"
  10. taos, err := sql.Open("taosRestful", taosDSN)
  11. if err != nil {
  12. log.Fatalln("failed to connect TDengine, err:", err)
  13. }
  14. defer taos.Close()
  15. rows, err := taos.Query("SELECT ts, current FROM meters LIMIT 2")
  16. if err != nil {
  17. log.Fatalln("failed to select from table, err:", err)
  18. }
  19. defer rows.Close()
  20. for rows.Next() {
  21. var r struct {
  22. ts time.Time
  23. current float32
  24. }
  25. err := rows.Scan(&r.ts, &r.current)
  26. if err != nil {
  27. log.Fatalln("scan error:\n", err)
  28. return
  29. }
  30. log.Println(r.ts, r.current)
  31. }
  32. }

查看源码

更多示例程序

使用限制

由于 REST 接口无状态所以 use db 语法不会生效,需要将 db 名称放到 SQL 语句中,如:create table if not exists tb1 (ts timestamp, a int)改为create table if not exists test.tb1 (ts timestamp, a int)否则将报错[0x217] Database not specified or available

也可以将 db 名称放到 DSN 中,将 root:taosdata@http(localhost:6041)/ 改为 root:taosdata@http(localhost:6041)/test。当指定的 db 不存在时执行 create database 语句不会报错,而执行针对该 db 的其他查询或写入操作会报错。

完整示例如下:

  1. package main
  2. import (
  3. "database/sql"
  4. "fmt"
  5. "time"
  6. _ "github.com/taosdata/driver-go/v3/taosRestful"
  7. )
  8. func main() {
  9. var taosDSN = "root:taosdata@http(localhost:6041)/test"
  10. taos, err := sql.Open("taosRestful", taosDSN)
  11. if err != nil {
  12. fmt.Println("failed to connect TDengine, err:", err)
  13. return
  14. }
  15. defer taos.Close()
  16. taos.Exec("create database if not exists test")
  17. taos.Exec("create table if not exists tb1 (ts timestamp, a int)")
  18. _, err = taos.Exec("insert into tb1 values(now, 0)(now+1s,1)(now+2s,2)(now+3s,3)")
  19. if err != nil {
  20. fmt.Println("failed to insert, err:", err)
  21. return
  22. }
  23. rows, err := taos.Query("select * from tb1")
  24. if err != nil {
  25. fmt.Println("failed to select from table, err:", err)
  26. return
  27. }
  28. defer rows.Close()
  29. for rows.Next() {
  30. var r struct {
  31. ts time.Time
  32. a int
  33. }
  34. err := rows.Scan(&r.ts, &r.a)
  35. if err != nil {
  36. fmt.Println("scan error:\n", err)
  37. return
  38. }
  39. fmt.Println(r.ts, r.a)
  40. }
  41. }

常见问题

  1. database/sql 中 stmt(参数绑定)相关接口崩溃

    REST 不支持参数绑定相关接口,建议使用db.Execdb.Query

  2. 使用 use db 语句后执行其他语句报错 [0x217] Database not specified or available

    在 REST 接口中 SQL 语句的执行无上下文关联,使用 use db 语句不会生效,解决办法见上方使用限制章节。

  3. 使用 taosSql 不报错使用 taosRestful 报错 [0x217] Database not specified or available

    因为 REST 接口无状态,使用 use db 语句不会生效,解决办法见上方使用限制章节。

  4. readBufferSize 参数调大后无明显效果

    readBufferSize 调大后会减少获取结果时 syscall 的调用。如果查询结果的数据量不大,修改该参数不会带来明显提升,如果该参数修改过大,瓶颈会在解析 JSON 数据。如果需要优化查询速度,需要根据实际情况调整该值来达到查询效果最优。

  5. disableCompression 参数设置为 false 时查询效率降低

    disableCompression 参数设置为 false 时查询结果会使用 gzip 压缩后传输,拿到数据后要先进行 gzip 解压。

  6. go get 命令无法获取包,或者获取包超时

    设置 Go 代理 go env -w GOPROXY=https://goproxy.cn,direct

常用 API

database/sql API

  • sql.Open(DRIVER_NAME string, dataSourceName string) *DB

    该 API 用来打开 DB,返回一个类型为 *DB 的对象。

Go - 图1info

该 API 成功创建的时候,并没有做权限等检查,只有在真正执行 Query 或者 Exec 的时候才能真正的去创建连接,并同时检查 user/password/host/port 是不是合法。

  • func (db *DB) Exec(query string, args ...interface{}) (Result, error)

    sql.Open 内置的方法,用来执行非查询相关 SQL。

  • func (db *DB) Query(query string, args ...interface{}) (*Rows, error)

    sql.Open 内置的方法,用来执行查询语句。

高级功能(af)API

af 包封装了连接管理、订阅、schemaless、参数绑定等 TDengine 高级功能。

连接管理

  • af.Open(host, user, pass, db string, port int) (*Connector, error)

    该 API 通过 cgo 创建与 taosd 的连接。

  • func (conn *Connector) Close() error

    关闭与 taosd 的连接。

订阅

  • func NewConsumer(conf *tmq.ConfigMap) (*Consumer, error)

    创建消费者。

  • func (c *Consumer) Subscribe(topic string, rebalanceCb RebalanceCb) error 注意:出于兼容目的保留 rebalanceCb 参数,当前未使用

    订阅单个主题。

  • func (c *Consumer) SubscribeTopics(topics []string, rebalanceCb RebalanceCb) error 注意:出于兼容目的保留 rebalanceCb 参数,当前未使用

    订阅主题。

  • func (c *Consumer) Poll(timeoutMs int) tmq.Event

    轮询消息。

  • func (c *Consumer) Commit() ([]tmq.TopicPartition, error) 注意:出于兼容目的保留 tmq.TopicPartition 参数,当前未使用

    提交消息。

  • func (c *Consumer) Close() error

    关闭连接。

schemaless

  • func (conn *Connector) InfluxDBInsertLines(lines []string, precision string) error

    写入 InfluxDB 行协议。

  • func (conn *Connector) OpenTSDBInsertTelnetLines(lines []string) error

    写入 OpenTDSB telnet 协议数据。

  • func (conn *Connector) OpenTSDBInsertJsonPayload(payload string) error

    写入 OpenTSDB JSON 协议数据。

参数绑定

  • func (conn *Connector) StmtExecute(sql string, params *param.Param) (res driver.Result, err error)

    参数绑定单行插入。

  • func (conn *Connector) InsertStmt() *insertstmt.InsertStmt

    初始化参数。

  • func (stmt *InsertStmt) Prepare(sql string) error

    参数绑定预处理 SQL 语句。

  • func (stmt *InsertStmt) SetTableName(name string) error

    参数绑定设置表名。

  • func (stmt *InsertStmt) SetSubTableName(name string) error

    参数绑定设置子表名。

  • func (stmt *InsertStmt) BindParam(params []*param.Param, bindType *param.ColumnType) error

    参数绑定多行数据。

  • func (stmt *InsertStmt) AddBatch() error

    添加到参数绑定批处理。

  • func (stmt *InsertStmt) Execute() error

    执行参数绑定。

  • func (stmt *InsertStmt) GetAffectedRows() int

    获取参数绑定插入受影响行数。

  • func (stmt *InsertStmt) Close() error

    结束参数绑定。

通过 WebSocket 订阅

  • func NewConsumer(conf *tmq.ConfigMap) (*Consumer, error)

    创建消费者。

  • func (c *Consumer) Subscribe(topic string, rebalanceCb RebalanceCb) error 注意:出于兼容目的保留 rebalanceCb 参数,当前未使用

    订阅单个主题。

  • func (c *Consumer) SubscribeTopics(topics []string, rebalanceCb RebalanceCb) error 注意:出于兼容目的保留 rebalanceCb 参数,当前未使用

    订阅主题。

  • func (c *Consumer) Poll(timeoutMs int) tmq.Event

    轮询消息。

  • func (c *Consumer) Commit() ([]tmq.TopicPartition, error) 注意:出于兼容目的保留 tmq.TopicPartition 参数,当前未使用

    提交消息。

  • func (c *Consumer) Close() error

    关闭连接。

完整订阅示例参见 GitHub 示例文件

通过 WebSocket 进行参数绑定

  • func NewConnector(config *Config) (*Connector, error)

    创建连接。

  • func (c *Connector) Init() (*Stmt, error)

    初始化参数。

  • func (c *Connector) Close() error

    关闭连接。

  • func (s *Stmt) Prepare(sql string) error

    参数绑定预处理 SQL 语句。

  • func (s *Stmt) SetTableName(name string) error

    参数绑定设置表名。

  • func (s *Stmt) SetTags(tags *param.Param, bindType *param.ColumnType) error

    参数绑定设置标签。

  • func (s *Stmt) BindParam(params []*param.Param, bindType *param.ColumnType) error

    参数绑定多行数据。

  • func (s *Stmt) AddBatch() error

    添加到参数绑定批处理。

  • func (s *Stmt) Exec() error

    执行参数绑定。

  • func (s *Stmt) GetAffectedRows() int

    获取参数绑定插入受影响行数。

  • func (s *Stmt) Close() error

    结束参数绑定。

完整参数绑定示例参见 GitHub 示例文件

API 参考

全部 API 见 driver-go 文档