一、从一次加班说起

之前我带过一个小项目,是做订单实时看板。老板要求每秒钟能看到全国订单的总金额和下单笔数。一开始,我照着网上教程上了Flink,从消息队列里读数据,用滚动窗口做聚合,再把结果写进Redis。Flink确实能跑,但用起来真的很折腾。环境要搭,作业要打包,任务挂了我们还得半夜爬起来看日志。最烦的是,业务想加一个维度,比如按商品分类统计,我得改代码、重新编译、再发版,前后没有半天搞不定。后来我换了个思路——既然数据都在时序数据库里,为什么不让数据库自己算呢?于是我开始用TDengine的流式计算。一条SQL就能把实时聚合吃进去,结果表自动维护,前端直接查。那之后,整个人轻松多了。这篇文章我就用最直白的语言,讲讲怎么用TDengine替代Flink,做实时聚合的正确姿势。

二、为什么说TDengine能干这件事

流计算听起来很高大上,其实核心就一句话:数据不停来,你得不停按时间窗口算结果。Flink是专门干这个的流处理引擎,但它是独立系统,不管存储。你要做实时聚合,得让Flink从Kafka拉数据,算完再存到Redis或MySQL里,链条很漫长。TDengine不一样,它本身就是时序数据库,把存储和计算合在了一起。你只要定义好一个流计算任务,后面来的数据它自动帮你算,结果也自动存起来。你永远不用关心"作业怎么提交""状态怎么保存"这些问题。

2.1 Flink做实时聚合的常规姿势

Flink要开窗口,比如每隔一分钟汇总一次。这还没完,数据可能乱序到达,你得给数据打水位线,不然窗口可能提前关掉,丢数据。算完之后,还要把结果往外发,通常要先序列化,再写到某个外部系统。这一整套流程,没有一定的编程功底很难快速上手。对于很多简单业务,真的有点牛刀杀鸡。

2.2 TDengine的姿势有多简单

TDengine把窗口意识做成了SQL关键字INTERVAL。你写一条CREATE STREAM,告诉它"从哪张表读数据、按多长时间的窗口聚合、结果写到哪张表",剩下全是它的活。什么水位线、状态后端、checkpoint,你统统不用管。它自己负责数据到达触发聚合、窗口自动翻转、结果自动落库。更省事的是,你查询结果直接用SQL,不用再对接其他数据库。

三、实战:用TDengine实现订单实时聚合

我们现在就开始动手。下面所有示例,技术栈统一用Python(taospy客户端),也就是说,所有操作都通过Python调用TDengine的SQL接口。无论你是刚入门的开发者,还是老手,照着做就能跑通。

3.1 连接TDengine

先建立连接。假设TDengine已经在你本地电脑上跑起来了,地址是localhost,端口6030,用户名和密码保持默认。

# 技术栈:Python(依赖 taospy,安装命令:pip install taospy)
import taos
import time

# 建立到TDengine的连接
conn = taos.connect(
    host="localhost",   # 数据库所在主机
    port=6030,          # TDengine默认端口
    user="root",        # 用户名
    password="taosdata" # 默认密码
)
# 创建游标,所有SQL都通过它来执行
cursor = conn.cursor()

3.2 建库和订单表

我们要做订单实时统计,首先需要一个数据库和一张订单表。订单表里至少要有时间戳、订单号、商品ID和金额。为了演示方便,我设计得比较简单,你业务里可以随意加字段。

# 建库:如果testdb不存在就创建,副本数设为1(单机测试足够)
cursor.execute("CREATE DATABASE IF NOT EXISTS testdb REPLICA 1")

# 切换到testdb这个库
cursor.execute("USE testdb")

# 创建订单表
# ts       :订单产生的时间,精确到毫秒
# order_id :订单编号,用字符串存
# item_id  :商品ID,用字符串存
# amount   :订单金额,浮点数
cursor.execute("""
    CREATE TABLE IF NOT EXISTS orders (
        ts TIMESTAMP,
        order_id VARCHAR(64),
        item_id VARCHAR(32),
        amount DOUBLE
    )
""")

这里要注意,TDengine对时间戳有严格要求,它是时序数据库,每张表都默认有主键时间戳。我们在建表时把ts定义成了时间戳列,这就是它的主轴。

3.3 定义流计算任务

这是核心一步。我想做的是:按一秒一个窗口,实时统计订单数量和总金额。在TDengine里,这个需求用一条CREATE STREAM就能表达。

# 创建流计算任务
# 含义:从orders表连续读取新数据,每1秒滚动聚合一次
# 聚合结果写入order_stats表,这张表由系统自动创建
cursor.execute("""
    CREATE STREAM IF NOT EXISTS order_stats_stream INTO order_stats AS
    SELECT
        _wstart       AS win_start,    -- 窗口的开始时间
        COUNT(*)      AS order_count,  -- 该窗口内的订单数量
        SUM(amount)   AS total_amount  -- 该窗口内的订单总金额
    FROM orders
    INTERVAL(1s)
""")

看不懂没关系,我拆开讲。_wstart是窗口开始时间,不用你算,引擎给。COUNT(*)就是数行数,SUM(amount)就是累加金额。INTERVAL(1s)告诉引擎窗口大小是一秒。就这么点东西,一个实时聚合任务就定义好了。怎么样,比Flink那套水位线、窗口函数、输出器简单得多吧。

3.4 模拟实时写入订单数据

流计算任务已经在后台默默运行了。现在我要向orders表插入订单数据,模拟业务系统持续产生订单。为了让结果有明显区别,我故意分两批插入:第一批两笔订单,放在同一个秒窗口里;第二批一笔订单,放在下一个秒窗口里。

# 获取当前时间的毫秒时间戳
now = int(time.time() * 1000)

# 第一批订单:两笔,时间相差100毫秒,会落在同一个秒窗口内
cursor.execute(
    f"INSERT INTO orders (ts, order_id, item_id, amount) VALUES "
    f"({now}, '001', 'A', 10.5), "
    f"({now + 100}, '002', 'B', 20.0)"
)

# 等1.2秒,让第一个一秒窗口自然闭合
time.sleep(1.2)

# 第二批订单:一笔,时间比第一批晚2秒,会落到下一个秒窗口
cursor.execute(
    f"INSERT INTO orders (ts, order_id, item_id, amount) VALUES "
    f"({now + 2000}, '003', 'C', 30.0)"
)

# 再等1.2秒,确保第二个窗口也完成计算
time.sleep(1.2)

注意,我用的时间戳是毫秒级。now是当前时间,now + 100代表过了100毫秒。如果当前时间正好位于某个秒窗口的尾部,这两条订单可能跨到下一个窗口,但概率很小。为了演示稳定,你可以在程序里故意把时间戳设成一个固定值附近,保证落在同一秒。这里我用now只是为了贴近真实场景。

3.5 查询聚合结果

现在,我们来看看流计算帮我们算出了什么结果。直接查询order_stats表就可以了。

# 查询结果表
cursor.execute("SELECT win_start, order_count, total_amount FROM order_stats")

# 逐行打印
for row in cursor.fetchall():
    print(f"窗口开始时间: {row[0]}, 订单数: {row[1]}, 总金额: {row[2]}")

如果你看到类似这样的输出:

窗口开始时间: 2025-04-01 10:00:01, 订单数: 2, 总金额: 30.5
窗口开始时间: 2025-04-01 10:00:02, 订单数: 1, 总金额: 30.0

就说明流计算完全正常。第一个窗口聚合了两笔订单,金额是10.5加20.0等于30.5;第二个窗口聚合了一笔订单,金额是30.0。更妙的是,如果你继续往里插数据,下一次查询时,这张表会自动多出新的窗口结果,不再需要你手动触发任何东西。

3.6 完整脚本参考

上面几步拆开看可能有点碎,这里给你一份完整的Python脚本,直接复制就能用。技术栈依然是Python和taospy,所有逻辑都在一个文件里。

# 技术栈:Python(依赖 taospy)
import taos
import time

# 连接
conn = taos.connect(host="localhost", port=6030, user="root", password="taosdata")
cursor = conn.cursor()

# 建库建表
cursor.execute("CREATE DATABASE IF NOT EXISTS testdb REPLICA 1")
cursor.execute("USE testdb")
cursor.execute("""
    CREATE TABLE IF NOT EXISTS orders (
        ts TIMESTAMP,
        order_id VARCHAR(64),
        item_id VARCHAR(32),
        amount DOUBLE
    )
""")

# 创建流计算:每1秒滚动聚合一次
cursor.execute("""
    CREATE STREAM IF NOT EXISTS order_stats_stream INTO order_stats AS
    SELECT
        _wstart AS win_start,
        COUNT(*) AS order_count,
        SUM(amount) AS total_amount
    FROM orders
    INTERVAL(1s)
""")

# 模拟写入两批订单
now = int(time.time() * 1000)
cursor.execute(
    f"INSERT INTO orders (ts, order_id, item_id, amount) VALUES "
    f"({now}, '001', 'A', 10.5), "
    f"({now + 100}, '002', 'B', 20.0)"
)
time.sleep(1.2)
cursor.execute(
    f"INSERT INTO orders (ts, order_id, item_id, amount) VALUES "
    f"({now + 2000}, '003', 'C', 30.0)"
)
time.sleep(1.2)

# 查询结果
cursor.execute("SELECT win_start, order_count, total_amount FROM order_stats")
for row in cursor.fetchall():
    print(f"窗口开始时间: {row[0]}, 订单数: {row[1]}, 总金额: {row[2]}")

# 关闭连接
cursor.close()
conn.close()

这份脚本虽然短,但已经具备了流计算最关键的三个要素:数据源表、流计算定义、结果表查询。你完全可以在这个基础上扩展自己的业务逻辑。

四、有哪些需要留心的地方

实时聚合不是说无脑用就行,有几个坑你一定会遇到。

第一,窗口对齐规则。TDengine的窗口是按时间戳对齐的,比如INTERVAL(1s),它以整秒为边界。并不是从你创建流计算那刻开始算,而是从每个整秒开始到下一整秒结束。这意味着如果你的数据时间戳不是正好落在某个窗口范围内,它可能会进入你预期之外的窗口。所以写数据时,时间戳一定要准确。

第二,乱序数据的问题。Flink可以通过水位线来等待乱序数据,TDengine目前更多是"来一条算一条"。如果你的业务对乱序非常敏感,比如订单可能延迟很久才上报,那就要谨慎。TDengine的流计算也提供了一些参数来设置容忍延迟,但默认行为不适合极端乱序场景。

第三,结果表是自动建的,字段类型和长度可能不是你期望的。比如COUNT(*)的结果是BIGINTSUM(amount)的结果是DOUBLE。如果后续你要用这些值做计算,记得先确认类型。另外,结果表名如果你不加库名前缀,它会被创建在当前数据库下面。

第四,流任务不会自动停止。你建了CREATE STREAM,它会一直运行,不停消费新数据。如果只是做测试,不想要这个任务了,要用DROP STREAM order_stats_stream来删掉,否则它会一直占着资源。

第五,别把流计算当成万能的。它适合简单的窗口聚合,比如SUM、COUNT、AVG、MAX、MIN。如果你要做多表复杂关联、事件状态机、自定义窗口,TDengine确实不如Flink灵活。这时候就不该硬用TDengine,老老实实上Flink。

五、适合用在什么场景

用TDengine做实时聚合,最适合的场景就是"时序数据 + 简单窗口统计"。我举几个实际例子。

第一个,订单实时统计。就是我们刚才演示的场景,每天订单量大的业务,按分钟、小时聚合销售额、订单数,给数据看板用。因为结果表是实时更新的,前端只要定时查询,就能看到最新数据。

第二个,物联网设备监控。设备每秒钟上报温度、湿度、电压等指标。你可以在TDengine里建一条流计算,按一分钟或一小时滑动窗口,算出每台设备的平均温度、最高温度、最低温度,一旦超阈值还能触发告警。这比把数据全部拿到流处理引擎里算要省事得多。

第三个,日志实时汇总。比如应用日志每天产生几亿条,你只想知道每个接口每分钟的调用次数和平均响应时间。日志一旦写入TDengine,流计算立刻帮你把结果算好,你只需要查询结果表,不用跑MapReduce。

第四个,车联网轨迹聚合。车辆GPS数据高频上报,你需要按车辆和五分钟窗口统计行驶里程、平均速度。TDengine的流计算和它的时序存储能力结合起来,效率很高。

这些场景都有几个共同点:数据带时间戳、需要按时间窗口聚合、业务逻辑不复杂。只要满足这三点,TDengine的流计算基本可以胜任。

六、技术优缺点总结

优点非常明显。第一,简单,一条SQL搞定实时聚合,没有编程门槛。第二,轻量,不需要额外部署一套流计算引擎,TDengine本身就有存储,结果直接落库。第三,好维护,改聚合逻辑只需要修改CREATE STREAM语句,不像Flink那样要重新打包发布。第四,查询方便,结果可以用标准的SQL查询,和普通表没区别。

缺点也不能忽视。第一,只擅长时序数据,如果数据不是带时间戳的,或者你需要非常复杂的业务事件处理,它就不适合。第二,乱序处理能力比Flink弱,极端乱序场景可能会算错。第三,生态和扩展能力不如Flink,Flink周边有各种连接器、管理平台、监控系统,TDengine作为数据库,还在慢慢完善。第四,如果流计算任务特别多,或者数据量极其巨大,TDengine的吞吐和性能可能不如专门为流处理设计的引擎。

这里我的态度很明确:TDengine的流计算不是来替代Flink生态的,它只是给大量简单场景提供了一个"多快好省"的选择。你完全可以把TDengine和Flink放在一起用:Flink处理复杂逻辑,结果再写进TDengine做存储和实时查询,各干各擅长的。

七、总结

咱们从一次加班说起,聊了用Flink做实时聚合的痛点,然后引出了TDengine流计算这个更轻的做法。接着我用一个订单实时统计的例子,手把手带你建了库、建了表、创建了流计算、模拟了数据写入,最后看到了聚合结果。说真的,当我在生产环境把Flink作业替换成TDengine流计算后,最大的感受就是省心。以前要盯作业状态、要处理反压,现在只要看数据库有没有数据进来,结果自己就出来了。如果你也正被实时聚合搞得焦头烂额,不妨先看看自己的业务是否属于简单时序聚合。如果是,那么用TDengine,可能就是你要找的"正确姿势"。