【系列:TDengine 工业物联网实战:从零搭起可运行系统 · 第 7 篇】

作者:jufeng1307日期:2026/8/20

第 4 篇我们搭好了写入管线,背压、攒批、并发 worker 都就位了。可管线再高效,也得有数据往里喂。今天回到源头:数据从哪来?两条路——模拟器直接造,采集器从三源收,最终都汇入同一条 pipeline。这篇就讲四件事:物理模型、确定性 seed、归一化漏斗、坏数据去向。读完你能回答:模拟器的物理模型为什么这么设计?三源各适合什么场景?坏数据去哪了?

两条路径,一个入口

先看全局。整个系统的数据入口分成两条路,但最终汇聚到同一个 WritePipeline——就是第 4 篇讲背压、攒批、并发的那条管线。

路径一:模拟器直连。 Simulator 里的每个设备直接调用 device.generate() 产出 Record,然后 pipeline.submit(record) 提交。这条路径不经过 normalizer,因为模拟器产出的数据天生就是标准模型——它自己就是"标准答案"。

路径二:采集器三源。 HTTP、MQTT、SCADA CSV 三个数据源接收外部数据,先经过 service.accept_json / accept_dict 的校验和归一化,再 pipeline.submit 提交。这条路径上,normalize_payload 是必经之路——外部设备厂商各异,字段名五花八门,必须统一。

两条路的定位不一样:模拟器服务开发调试、演示、跑基准;采集器服务生产接入。但进了 pipeline 之后,待遇完全一样——背压、批量、并发,一个不少。

模拟器三场景:不是随机数,是物理模型

模拟器不是简单地 random.random() 糊弄数据。三个场景各有各的物理逻辑,这是整个系统能"以假乱真"的关键。

环境传感器:日周期正弦

SensorDevice 模拟的是环境监测站。核心逻辑是日周期正弦——一天 86400 秒,对应 2π 相位:

1daily_phase = timestamp.timestamp() / 86_400 * math.tau
2temperature = self.baseline_temperature + math.sin(daily_phase) * 4
3temperature += self._random.gauss(0, 0.25)
4humidity = self.baseline_humidity - math.sin(daily_phase) * 8
5humidity += self._random.gauss(0, 0.8)
6

基线 24°C、55%RH,温度振幅 ±4°C,湿度振幅 ±8%RH。注意温度和湿度是反相的——温度加 sin,湿度减 sin,白天温度高湿度低,夜里反过来,此消彼长。

噪声用高斯分布:温度 σ=0.25,湿度 σ=0.8。湿度做了钳位 min(100, max(0, ...)),防止超过物理极限。

工业电机:负载耦合 + 长期退化

IndustrialDevice 是这套模拟器里最有意思的部分。它模拟的是退化模型——设备会随着运行时间慢慢老化:

1self.degradation = min(1.0, self.degradation + self._random.uniform(0, 0.00001))
2load = 0.65 + 0.25 * math.sin(self.sequence / 90)
3load += self._random.gauss(0, 0.02)
4rotational_speed = round(1450 * load + self._random.gauss(0, 5))
5vibration = 0.02 + load * 0.04 + self.degradation * 0.5
6

关键在负载耦合:转速 = 1450 × load,振动 = 0.02 + load×0.04 + 退化×0.5,温度 = 30 + load×35 + vibration×12。所有指标都随负载联动,就像真实电机——负载上去,转速、振动、温度一起变。

长期退化是点睛之笔:每轮 degradation 增加 0~1e-5,封顶 1.0。退化步长很小,但架不住积累——按平均 5e-6 每轮估算,跑 5 万条左右振动就能跨过 0.18 进入预警,十万条上下跨过 0.35 报故障。status_code 从 1(正常)→ 2(预警)→ 3(故障)逐级跳变。这就是喂给第 9 篇告警引擎的"慢性病"数据。

车辆轨迹:积分 + 超速报警

VehicleDevice 模拟的是 GPS 轨迹,从上海 (121.4737, 31.2304) 出发:

1speed = max(0, 45 + 25 * math.sin(self.sequence / 120) + self._random.gauss(0, 3))
2self.direction = (self.direction + self._random.gauss(0, 1.5)) % 360
3distance_km = speed / 3600
4radians = math.radians(self.direction)
5self.latitude += math.cos(radians) * distance_km / 111.0
6longitude_scale = max(0.01, math.cos(math.radians(self.latitude)))
7self.longitude += math.sin(radians) * distance_km / (111.0 * longitude_scale)
8

这是轨迹积分:每 tick 前进 speed/3600 公里(1 秒一条数据,一小时 3600 秒),按方向角分解成经纬度增量。1 度纬度 ≈ 111km,经度按 cos(纬度) 缩放——高纬度地区经度圈更小,这是真实的地理投影。

速度 45±25 km/h,方向随机游走。超 100 km/h 就置 alarm_code=1。产出的是 VehicleTrack,走独立的超级表,和第 2 篇的 vehicle_track 对应。

确定性 seed:为什么重跑必须逐字节一致

模拟器里每个设备都持有独立的 random.Random(self.seed),而 factory 传的是 seed + index。默认 seed 是 20260808

为什么要确定性?两个原因。

第一,bug 可复现。 线上出问题,同一 seed 重跑一遍,数据逐字节一致。你能在本地复现、调试、修复,而不是"这次跑出来的数据和上次不一样,不知道是不是修好了"。

第二,基准可对比。 第 6 篇我们讲公平基准——不同方案对比,数据必须完全一致才有意义。如果每次跑数据都随机,你没法判断性能差异是方案带来的还是数据波动带来的。

factory 里还有一层细节:factory_no = index % 8 + 1(8 个工厂)、workshop_no = index % 20 + 1(20 个车间)、region 四选一。设备 ID 是 device-{index:07d},车辆有 fleet-{index%50:03d} 和车型/省份维度。这些维度组合起来,就是 TDengine 超级表里的 tag 体系。

运行器:不漂移的节拍器

模拟器跑起来靠 runner.py 的定时调度。核心是 next_tick 机制:

1interval = 1 / self.settings.simulator_rate
2next_tick = monotonic()
3while not self.stop_event.is_set():
4    for device in devices:
5        record = device.generate()
6        await self.pipeline.submit(record)
7        self.generated += 1
8    next_tick += interval
9    delay = next_tick - monotonic()
10    if delay > 0:
11        await asyncio.sleep(delay)
12    else:
13        next_tick = monotonic()
14

next_tick绝对时间推进,每轮生成耗时不计入间隔。就算某轮生成慢了,下一轮也会立即补发,不累积漂移。

分片并发:设备按 index % shard_count 轮转分到多个 producer 协程,充分利用 asyncio 并发。SIGINT/SIGTERM 触发优雅停机——先 stop_event 取消 producers,再 pipeline.stop() 把缓冲区数据排空。

CLI 默认参数:100 台设备、每秒每设备 1 条、mixed 场景、WebSocket 传输、batch_size=1000、4 个 worker。跑完输出 JSON 统计:{records, elapsed_seconds, records_per_second}

采集器:失败不抛异常

采集器服务和模拟器完全不同——它面对的是不可控的外部世界。设备厂商的固件可能发垃圾数据,网络可能半路截断,JSON 可能格式错误。所以 service.py 的设计哲学是:绝不因单条坏数据崩溃

1async def accept_json(self, value, *, source):
2    try:
3        raw = json.loads(value)
4    except (json.JSONDecodeError, UnicodeDecodeError) as error:
5        await self.dead_letters.write(str(value), str(error), source)
6        RECORDS_FAILED.labels("collector", "invalid_json").inc()
7        return False
8    if not isinstance(raw, dict):
9        await self.dead_letters.write(raw, "payload must be an object", source)
10        RECORDS_FAILED.labels("collector", "invalid_shape").inc()
11        return False
12    return await self.accept_dict(raw, source=source)
13

坏数据不抛异常,而是写进 dead-letter 文件 + 增加 RECORDS_FAILED 指标,接口层返回 422。主流程照常跑,坏数据不拖累好数据——生产系统的底线。

正常数据走 accept_dictnormalize_payloadpipeline.submit,成功后 RECORDS_ACCEPTED.labels(source, "telemetry").inc()。每个来源的接受/拒绝数都能在 Prometheus 9108 端口看到。

normalizer:异构设备的归一化漏斗

三源的数据格式千奇百怪,但进了系统必须统一成标准模型。normalize_payload 就是那个漏斗。

第一层:别名映射。 设备厂商字段名各不相同——有的叫 temp,有的叫 temperature_c,有的叫 rh 表示湿度,电机厂商用 amps 表示电流。9 个别名映射全部统一到模型字段:

1ALIASES = {
2    "temp": "temperature", "temperature_c": "temperature", "rh": "humidity",
3    "amps": "current_value", "current": "current_value", "kw": "power",
4    "rpm": "rotational_speed", "flow": "flow_rate", "seq": "sequence_no",
5}
6

第二层:时间解析。 timestampts 字段,支持 datetime 对象、数字(大于 10_000_000_000 判为毫秒)、ISO-8601 字符串(Z 自动转 +00:00)。naive 时间补 UTC——不带时区的时间默认当成 UTC,避免歧义。

第三层:范围校验。 温度必须在 -273.15~2000°C,湿度 0~100%,电压 0~1e6,状态码 -128~127(对应 TDengine TINYINT)。越界或非数字直接抛 NormalizationError,进 dead-letter。

还有两个细节:device_id 必填,缺失直接拒绝;六种 tag(product_key/factory_id/workshop_id/region/device_type)缺失给默认值。tag 文本限制 ≤64 字节——这是 TDengine 表名 tag 的长度上限。

三源逐个看

HTTP:最通用的入口

POST /v1/telemetry,aiohttp 实现。client_max_size=1MB + content_length 前置检查,超限直接 413。成功返回 202 {"accepted": true},归一化失败返回 422 {"accepted": false}

健康检查拆成两个:/health/live 恒 UP(进程活着就算活),/health/readywriter.health_check()(pipeline 能收数据才算就绪)。K8s 里 live 和 ready 分开是标准姿势——进程活着但处理不了请求时,不能继续给它流量。

默认监听 0.0.0.0:8090

MQTT:物联网的事实标准

aiomqtt 客户端,URL 校验 scheme 必须是 mqttmqtts。mqtts 自动建 TLS 默认上下文,默认端口 1883/8883。

亮点是主题通配符factory/+/telemetry+ 匹配任意工厂。一个订阅覆盖所有工厂的遥测数据,新工厂上线不用改代码。

async for message in client.messages 流式消费,每条消息走 accept_json。没装 aiomqtt 会提示 install the 'mqtt' extra——依赖按需安装,不强制。

SCADA CSV:存量系统的桥接

很多工厂的 SCADA 系统只导出 CSV。scada_source.pycsv.DictReader 流式读取,不整文件载入内存——几 GB 的导出文件也能处理。

两个细节:utf-8-sig 编码兼容 Excel 的 BOM 头;每 1000 行 asyncio.sleep(0) 让出事件循环,避免阻塞其他协程。返回 (accepted, rejected) 计数,CLI 输出 accepted=.. rejected=..

DeadLetterSink:坏数据的归宿

坏数据写到 data/dead-letter/dead-letter-YYYY-MM-DD.jsonl,按天分文件。O_APPEND|O_CREAT|0o600 + fsync——追加写、断电不丢、权限 600。每条记录含 received_at / source / reason / payload,方便事后排查。

坏数据完整保留在 dead-letter 里——不污染主数据流,现场也都在,事后能查、能重放、能分析。

汇总:谁用哪条路?

两条路径,一张表说清楚:

路径适用场景数据形态是否过 normalizer
模拟器直连开发调试、演示、性能基准标准 Record(Telemetry / VehicleTrack)
采集器三源生产环境接入真实设备异构 JSON / MQTT 消息 / CSV

模拟器是开发者的"自来水"——开箱即用,数据标准,确定性可复现。采集器是生产环境的"引水渠"——面对不可控的外部世界,用归一化和 dead-letter 保底。

两条路最终都汇入同一条 pipeline,从第 4 篇的背压、攒批、并发开始,一路走向存储和查询。

下一篇,我们进入第 8 篇:Java 多数据源 API——PostgreSQL 管设备档案,TDengine 管时序数据。入口有了,数据有了,该聊聊怎么把数据用起来了。


你现在的系统里,数据入口是模拟器还是真实采集?归一化遇到最头疼的脏数据是什么样的?欢迎留言聊聊。

觉得有用?点个关注,持续获取优质内容。


【系列:TDengine 工业物联网实战:从零搭起可运行系统 · 第 7 篇】》 是转载文章,点击查看原文


相关推荐


uni-app 生命周期深度解析(iOS / Android / 鸿蒙 / Vue3 四端对照)
90后晨仔2026/8/7

📌 本文定位: 面向同时具备 iOS、Android、鸿蒙原生开发经验,正在转向 uni-app 的资深工程师。所有技术点均基于 uni-app 官方文档 及 Vue3 官方文档整理。 一、先搞懂一个 JS 语法问题:为什么说"与 data/methods 平级的叫生命周期"? 很多原生工程师初学 uni-app 时会对这句话感到困惑。我们从 JavaScript 语言层面彻底讲清楚。 1.1 Options API 的本质是一个"配置对象" 在 Vue2/Vue3 的 Options A


GPT 5.6 的真正变革:从"最强模型"到"最合适模型",AI 工程范式正在重置
浮生望2026/7/29

摘要: GPT 5.6以Sol、Terra、Luna三层架构重新定义AI使用方式。跑分不再决定竞争力,真正的差距在于把什么任务分给什么模型,用更少Token和更少返工完成更高质量交付。AI正从顾问变成协作者。 每次新模型发布,人们都在问同一个问题,但那个问题已经过时了 GPT 5.6 发布后,讨论一如既往地集中在几个老问题上:跑分涨了多少?代码能力有没有超过 Claude?谁又成了"最强模型"? 这些问题很重要,但如果你只盯着它们,就错过了 GPT 5.6 真正改变游戏规则的地方。 Open


影像创作领域的GitHub!深度体验了这个国产AI,我卸了3个AI视频工具!
程序员X小鹿2026/7/21

大家好,我是X小鹿,今天分享一个不一样的 AI 视频工具。 以前做一条 AI 短片,从剧本创作、脚本编写、到角色定妆、场景道具生成、再到分镜生成、视频生成,以及最后的剪辑成片,每一步都需要用到单独的工具。 生成一条短片,需要在多个工具间来回切换。 后来也陆续出现了一些 AI 视频工具,把上面的很多流程都集成在了一个平台上,用起来确实比之前方便了不少。 最近用了 AI 视频创作平台 updream ,发现它虽然也是将很多功能集成在了一起,但和之前接触的一些 AI 视频工具又不太一样。 体验下来,


图像的分辨率
元來2026/7/12

一、什么是分辨率 图像分辨率表示: 图像能够表示多少空间细节,或者图像在水平和垂直方向上包含多少个像素。 日常最常见的分辨率写法是: 宽度 × 高度 例如: 1920 × 1080 4096 × 4096 6144 × 6144 其中: 1920表示图像水平方向有1920个像素;1080表示图像垂直方向有1080个像素。 因此一张1920 × 1080的图像,像素总数为:1920×1080=2,073,600 即大约207万个像素,也可以称为约200万像素图像。 二、


OpenCV-Python实战(31)——实时面部情绪检测与识别系统
盼小辉丶2026/7/4

OpenCV-Python实战(31)——实时面部情绪检测与识别系统 0. 前言1. 规划应用程序1. 人脸检测1.1 基于 Haar 的级联分类器1.2 预训练的级联分类器1.3 使用预训练的级联分类器1.4 FaceDetector 类 2. 收集数据2.1 构建训练数据集2.2 运行应用程序2.3 实现数据收集器 GUI 3. 面部情绪识别3.1 处理数据集3.2 多层感知机3.3 构建 MLP用于面部表情识别 4. 整合所有内容小结系列链接 0. 前言 我


2026年6月远程控制软件横评:UU远程、ToDesk、向日葵全方位对比
凤年徐2026/6/26

2026远程控制软件横评:UU远程、ToDesk、向日葵全方位对比 远程控制早已不是“应急连一下电脑”那么简单了。开发者用它连服务器改代码,设计师用它调家里的渲染机,留学生用它操作国内的网盘和银行App,甚至游戏玩家也用它“云挂机”。市面上主流的三款工具——UU远程、ToDesk、向日葵——各有拥趸,但到底谁在哪个场景下更顺手?这篇横评不吹不黑,直接把你最关心的七个核心能力摆在一起,逐项对比。 一、终端能力:谁能让开发者真正扔掉SSH客户端? 对于开发者、运维和AI训练者来说,远程命令行是最


JavaScript 函数性能优化:配置驱动 + 按需计算实战指南
m0_733915432026/6/17

JavaScript 函数优化是前端性能调优的核心环节。当业务逻辑中存在多层条件分支与重复计算时,代码不仅执行效率低下,还难以维护和扩展。本文将基于真实业务场景,提供一套配置驱动的优化方案,通过规则表集中管理 + 按需计算策略,帮助你提升函数执行效率 80% 以上。 为什么传统条件分支会拖累性能? 理解问题根源是有效优化的前提。传统 if-else 嵌套方案存在三大性能瓶颈: 问题类型 具体表现 性能影响 重复计算


PySide6 + Qt Designer + PyCharm 完整开发流程
资深流水灯工程师2026/6/10

PyCharm 对 Python 桌面开发有更完善的支持,包括智能代码补全、断点调试、集成终端、版本控制等功能,结合 Qt Designer 的可视化 UI 设计,是工业级上位机开发的首选组合。以下是完全适配 PyCharm 的标准化开发流程。 一、环境准备与 PyCharm 配置 1. 创建项目并配置虚拟环境(必做) PyCharm 强烈推荐使用虚拟环境隔离项目依赖,避免版本冲突: 打开 PyCharm → 新建项目 (New Project)选择项目位置,命名为test_equipm


claude-code下载安装与使用
veminhe2026/6/2

1、官方网站 Claude Code by Anthropic | AI Coding Agent, Terminal, IDE 2、github的项目地址 GitHub - anthropics/claude-code: Claude Code is an agentic coding tool that lives in your terminal, understands your codebase, and helps you code faster by executing ro


从社区路标到生态基石:Dave Verwer 的新篇章 -- 肘子的 Swift 周报 #137
东坡肘子2026/5/26

从社区路标到生态基石:Dave Verwer 的新篇章 Dave Verwer 在 iOS Dev Weekly 第 751 期宣布,这份已经持续近 15 年的周报将交由新的团队继续运营,而他自己接下来会全职投入 Swift Package Index。我的博客在早期获得关注,也曾得益于 iOS Dev Weekly 的推荐;而我在周报中坚持撰写每期周评,同样在很大程度上受到 Dave Verwer 的启发。对于很多 Apple 平台开发者来说,iOS Dev Weekly 早已不只是一份链接合

首页编辑器站点地图

本站内容在 CC BY-SA 4.0 协议下发布

Copyright © 2026 聚合阅读