当前位置:首页>python>TDengine Python 连接器 — taospy 与 taos-ws-py 完整指南

TDengine Python 连接器 — taospy 与 taos-ws-py 完整指南

  • 2026-10-11 06:16:12
TDengine Python 连接器 — taospy 与 taos-ws-py 完整指南

分类:11.客户端与连接器 | 篇章:04 Python 连接器详情

TDengine 提供 taospy(Native)和 taos-ws-py(WebSocket)两个 Python 包,均实现 PEP 249 DB-API 规范。本文涵盖安装、基本用法、STMT2、Schemaless、订阅、Pandas/SQLAlchemy 集成。

包速查表

包
协议
端口
依赖
taospy
Native
6030
libtaos
taos-ws-py
WebSocket
6041
无
taospy (RestClient)
REST
6041
无

详细解析

1. 安装

# Native (需要 libtaos)pip install taospy# WebSocket (推荐)pip install taos-ws-py# 两个都装(按需选用)pip install taospy taos-ws-py

2. 基本用法

# Nativeimport taosconn = taos.connect(    host="localhost",    port=6030,    user="root",    password="taosdata",    database="test")# WebSocketimport taoswsconn = taosws.connect("taosws://root:taosdata@localhost:6041/test")# 查询cursor = conn.cursor()cursor.execute("CREATE DATABASE IF NOT EXISTS test")cursor.execute("USE test")cursor.execute("""  CREATE STABLE IF NOT EXISTS meters (    ts TIMESTAMP, current FLOAT, voltage INT  ) TAGS (location VARCHAR(32))""")# 插入cursor.execute("""  INSERT INTO d001 USING meters TAGS('Beijing')   VALUES (NOW, 25.3, 220)""")# 查询cursor.execute("SELECT * FROM meters LIMIT 10")for row in cursor:    print(row)cursor.close()conn.close()

3. 参数化查询

# 标准 DB-APIcursor.execute(    "SELECT * FROM meters WHERE ts > ? AND current > ?",    (datetime.now() - timedelta(hours=1), 100))# 批量插入cursor.executemany(    "INSERT INTO d001 VALUES (?, ?, ?)",    [(ts1, 25.0, 220), (ts2, 26.0, 221), ...])

4. STMT2 高性能写入

# WebSocket STMT2import taoswsconn = taosws.connect("taosws://root:taosdata@localhost:6041/test")stmt = conn.statement2(    "INSERT INTO ? USING meters TAGS(?) VALUES(?, ?, ?)")# 准备多个子表的数据table_names = ["d001", "d002"]tags = [["Beijing"], ["Shanghai"]]cols = [    [        [int(time.time()*1000), int(time.time()*1000)+1000],  # ts        [25.3, 26.1],                                          # current        [220, 221]                                             # voltage    ],    [        [int(time.time()*1000), int(time.time()*1000)+1000],        [27.5, 28.2],        [218, 219]    ]]stmt.bind_param(table_names, tags, cols)stmt.execute()stmt.close()

5. Schemaless 写入

import taosfrom taos import SmlProtocol, SmlPrecisionconn = taos.connect(host="localhost", database="test")# InfluxDB Linelines = [    "meters,location=Beijing current=25.3,voltage=220i 1717488000000",    "meters,location=Shanghai current=26.1,voltage=219i 1717488001000"]conn.schemaless_insert(    lines,     SmlProtocol.LINE_PROTOCOL,    SmlPrecision.MILLI_SECONDS)# OpenTSDB JSONimport jsondata = [{    "metric": "meters.current",    "timestamp": 1717488000,    "value": 25.3,    "tags": {"location": "Beijing"}}]conn.schemaless_insert(    [json.dumps(data)],    SmlProtocol.JSON_PROTOCOL,    SmlPrecision.SECONDS)

6. TMQ 订阅

from taos.tmq import Consumerconsumer = Consumer({    "td.connect.ip": "localhost",    "td.connect.user": "root",    "td.connect.pass": "taosdata",    "group.id": "my_group",    "auto.offset.reset": "earliest",    "enable.auto.commit": "true"})consumer.subscribe(["my_topic"])try:    while True:        msg = consumer.poll(1)  # 1 秒超时        if msg is None:            continue        for block in msg:            for row in block:                print(row)        consumer.commit()except KeyboardInterrupt:    passfinally:    consumer.close()# WebSocket 版本from taosws import TmqConsumerconsumer = TmqConsumer(    "taosws://root:taosdata@localhost:6041/",    group_id="my_group",    auto_offset_reset="earliest")consumer.subscribe(["my_topic"])

7. Pandas 集成

import pandas as pdimport taoswsconn = taosws.connect("taosws://root:taosdata@localhost:6041/test")# 查询 → DataFramedf = pd.read_sql(    "SELECT _wstart, AVG(current) AS avg_c FROM meters "    "WHERE ts > NOW - 1d INTERVAL(1h)",    conn)print(df.head())# 处理df.set_index('_wstart', inplace=True)df.plot(figsize=(12, 4))# DataFrame → 写入df_to_insert = pd.DataFrame({    'ts': pd.date_range('2026-06-04', periods=1000, freq='S'),    'current': np.random.rand(1000) * 30,    'voltage': np.random.randint(218, 223, 1000)})cursor = conn.cursor()for _, row in df_to_insert.iterrows():    cursor.execute(        f"INSERT INTO d001 VALUES ('{row.ts}', {row.current}, {row.voltage})"    )

8. SQLAlchemy 集成

from sqlalchemy import create_engine, text# Nativeengine = create_engine("taos://root:taosdata@localhost:6030/test")# WebSocketengine = create_engine("taosws://root:taosdata@localhost:6041/test")# RESTengine = create_engine("taosrest://root:taosdata@localhost:6041/test")with engine.connect() as conn:    result = conn.execute(text("SELECT * FROM meters LIMIT 10"))    for row in result:        print(row)# Pandas + SQLAlchemydf = pd.read_sql("SELECT * FROM meters LIMIT 100", engine)

代码示例

完整数据采集脚本

import timeimport randomimport taoswsfrom datetime import datetimedef main():    conn = taosws.connect(        "taosws://root:taosdata@localhost:6041/iot")    cursor = conn.cursor()    # 准备 schema    cursor.execute("""      CREATE DATABASE IF NOT EXISTS iot    """)    cursor.execute("""      CREATE STABLE IF NOT EXISTS sensors (        ts TIMESTAMP, temp FLOAT, humidity FLOAT      ) TAGS (device_id VARCHAR(32), location VARCHAR(32))    """)    # 模拟采集    devices = [        ("d001", "Floor1"),        ("d002", "Floor1"),        ("d003", "Floor2"),    ]    while True:        ts = int(time.time() * 1000)        for dev_id, loc in devices:            temp = 20 + random.random() * 10            humidity = 40 + random.random() * 30            cursor.execute(                f"INSERT INTO {dev_id} USING sensors "                f"TAGS('{dev_id}', '{loc}') "                f"VALUES ({ts}, {temp:.2f}, {humidity:.2f})"            )        time.sleep(1)if __name__ == "__main__":    main()

数据分析

import taoswsimport pandas as pdimport matplotlib.pyplot as pltconn = taosws.connect("taosws://root:taosdata@localhost:6041/iot")# 小时聚合df = pd.read_sql("""  SELECT _wstart, location, AVG(temp) AS avg_temp   FROM sensors   WHERE ts > NOW - 1d   PARTITION BY location   INTERVAL(1h)""", conn)# 数据透视pivot = df.pivot(index='_wstart', columns='location', values='avg_temp')pivot.plot(figsize=(12, 6))plt.title("Hourly Temperature by Location")plt.savefig("temp_trend.png")

性能考量

Python 连接器性能

操作
吞吐
cursor.execute() 单条
几千行/秒
executemany 批量
几万行/秒
Schemaless
几十万行/秒
STMT2 列绑定
上百万行/秒

优化建议

措施
效果
用 STMT2
5~10x 提升
大批量绑定
显著提升
asyncio 异步
并发提升
Pandas to_sql
简洁但通常慢

FAQ

Q1: taospy vs taos-ws-py 怎么选?

  • 本地开发:taos-ws-py(无 libtaos 依赖)
  • 容器化部署:taos-ws-py
  • 极致性能:taospy
  • 浏览器场景:taos-ws-py

Q2: 安装 taospy 报错找不到 libtaos?

需先安装 TDengine 客户端:

# 下载 taosToolswget ... && tar xzf ...sudo ./install.sh -e no

Q3: 异步支持?

taos-ws-py 提供异步接口:

import taoswsasync def query():    conn = await taosws.connect_async("...")    cursor = await conn.cursor_async()    await cursor.execute("SELECT ...")

Q4: 时间戳类型?

返回 datetime(Native)或 int(WS,毫秒)。注意类型差异。

Q5: 中文乱码?

数据库 PRECISION 与 LOCALE 设置正确:

conn = taos.connect(..., config="/etc/taos")# taos.cfg 中设 locale en_US.UTF-8

参考

系统构架篇

  • 01-《TDengine 整体架构全景》
  • 02-《集群拓扑深度解析》
  • 03-《MNode 内部机制深度解析》
  • 04-《RPC 通信层深度解析》
  • 05-《VNode 生命周期》
  • 06-《RAFT 共识协议》
  • 07-《端到端的消息流》

数据模型

  • 01-《数据库创建与参数详解》
  • 02-《超级表/子表/普通表》
  • 03-《支持数据类型深度解析》
  • 04-《TDengine Tag 设计哲学与 Schema 变更机制》
  • 05-《TDengine 虚拟表实现原理》

存储引擎

  • 01-《TDengine 存储引擎概览》
  • 02-《TDengine MemTable 深度解析》
  • 03-《TDengine WAL 预写日志机制》
  • 04-《TDengine 数据文件格式》
  • 05-《TDengine Commit 与 Flush 机制 》
  • 06-《TDengine Compaction 合并策略 》
  • 07-《TDengine 数据保留与 TTL》
  • 08-《TDengine 压缩编码机制》
  • 09-《TDengine Cache 与 Last 查询加速》
  • 10-《TDengine 逻辑计划生成》

查询引擎

  • 01-《TDengine 查询引擎概览》
  • 02-《TDengine SQL 解析与词法分析》
  • 03-《TDengine 语义分析与 AST 重写》
  • 04-《TDengine 逻辑计划生成》
  • 05-《TDengine 物理计划生成》
  • 06-《TDengine 扫描算子》
  • 07-《TDengine 聚合算子》
  • 08-《TDengine 连接算子》
  • 09-《TDengine 排序、填充与投影》
  • 10-《TDengine 分布式查询执行》
  • 11-《TDengine EXPLAIN 与查询优化》

数据写入

  • 01-《TDengine SQL INSERT》
  • 02-《TDengine 无模式写入》
  • 03-《TDengine STMT 写入》
  • 04-《TDengine 写入内部流程》
  • 05-《TDengine 数据更新删除》

数据订阅

  • 01-《TDengine 数据订阅》
  • 02-《TDengine 订阅 vs Kafka》
  • 03-《TDengine TMQ 消费流程》
  • 04-《TDengine 内部机制》
  • 05-《TDengine TMQ 最佳实践》

预聚合

  • 01-《TDengine RSMA》
  • 02-《TDengine TSMA — 时间维度的物化聚合视图》
  • 03-《TDengine SMA 内部实现》

索引

  • 01-《TDengine Tag 索引》
  • 02-《TDengine SMA 索引》

SQL 语句

  • 01-《TDengine DDL》
  • 02-《TDengine DML SELECT》
  • 03-《TDengine DML 函数完整参考》
  • 04-《TDengine JOIN 完整语法》
  • 05-《TDengine 窗口完整语法》
  • 06-《TDengine 操作符与表达式》
  • 07-《TDengine 系统表》
  • 08-《TDengine SQL 与标准 SQL 差异》

客户端与连接器

  • 01-《TDengine 的连接方式》
  • 02-《TDengine C/C++ 连接器》
  • 03-《TDengine java 连接器》

关于 TDengine

TDengine 专为物联网IoT平台、工业大数据平台设计。其中,TDengine TSDB 是一款高性能、分布式的时序数据库(Time Series Database),同时它还带有内建的缓存、流式计算、数据订阅等系统功能;TDengine IDMP 是一款AI原生工业数据管理平台,它通过树状层次结构建立数据目录,对数据进行标准化、情景化,并通过 AI 提供实时分析、可视化、事件管理与报警等功能。

最新文章

随机文章