分类:11.客户端与连接器 | 篇章:04 Python 连接器
详情
TDengine 提供 taospy(Native)和 taos-ws-py(WebSocket)两个 Python 包,均实现 PEP 249 DB-API 规范。本文涵盖安装、基本用法、STMT2、Schemaless、订阅、Pandas/SQLAlchemy 集成。
包速查表
详细解析
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 连接器性能
优化建议
FAQ
Q1: taospy vs taos-ws-py 怎么选?
- 本地开发:taos-ws-py(无 libtaos 依赖)
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
参考
系统构架篇
数据模型
- 04-《TDengine Tag 设计哲学与 Schema 变更机制》
存储引擎
- 02-《TDengine MemTable 深度解析》
- 05-《TDengine Commit 与 Flush 机制 》
- 06-《TDengine Compaction 合并策略 》
- 09-《TDengine Cache 与 Last 查询加速》
查询引擎
- 02-《TDengine SQL 解析与词法分析》
- 03-《TDengine 语义分析与 AST 重写》
- 11-《TDengine EXPLAIN 与查询优化》
数据写入
数据订阅
- 02-《TDengine 订阅 vs Kafka》
预聚合
- 02-《TDengine TSMA — 时间维度的物化聚合视图》
索引
SQL 语句
- 08-《TDengine SQL 与标准 SQL 差异》
客户端与连接器
关于 TDengine
TDengine 专为物联网IoT平台、工业大数据平台设计。其中,TDengine TSDB 是一款高性能、分布式的时序数据库(Time Series Database),同时它还带有内建的缓存、流式计算、数据订阅等系统功能;TDengine IDMP 是一款AI原生工业数据管理平台,它通过树状层次结构建立数据目录,对数据进行标准化、情景化,并通过 AI 提供实时分析、可视化、事件管理与报警等功能。