第 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 之后,待遇完全一样——背压、批量、并发,一个不少。

信息图:两条数据路径,模拟器直连 pipeline,三源经 normalizer 后进同一 pipeline

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

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

环境传感器:日周期正弦

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

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

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

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

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

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

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

关键在负载耦合:转速 = 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) 出发:

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

这是轨迹积分:每 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 机制:

interval = 1 / self.settings.simulator_rate
next_tick = monotonic()
while not self.stop_event.is_set():
    for device in devices:
        record = device.generate()
        await self.pipeline.submit(record)
        self.generated += 1
    next_tick += interval
    delay = next_tick - monotonic()
    if delay > 0:
        await asyncio.sleep(delay)
    else:
        next_tick = monotonic()

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 的设计哲学是:绝不因单条坏数据崩溃

async def accept_json(self, value, *, source):
    try:
        raw = json.loads(value)
    except (json.JSONDecodeError, UnicodeDecodeError) as error:
        await self.dead_letters.write(str(value), str(error), source)
        RECORDS_FAILED.labels("collector", "invalid_json").inc()
        return False
    if not isinstance(raw, dict):
        await self.dead_letters.write(raw, "payload must be an object", source)
        RECORDS_FAILED.labels("collector", "invalid_shape").inc()
        return False
    return await self.accept_dict(raw, source=source)

坏数据不抛异常,而是写进 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 个别名映射全部统一到模型字段:

ALIASES = {
    "temp": "temperature", "temperature_c": "temperature", "rh": "humidity",
    "amps": "current_value", "current": "current_value", "kw": "power",
    "rpm": "rotational_speed", "flow": "flow_rate", "seq": "sequence_no",
}

第二层:时间解析。 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 的长度上限。

信息图:归一化漏斗,别名映射、时间解析、范围校验,坏数据进 dead-letter

三源逐个看

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 里——不污染主数据流,现场也都在,事后能查、能重放、能分析。

氛围:三源汇聚到 collector 服务的拓扑示意图,数据流清晰

汇总:谁用哪条路?

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

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

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

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

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


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

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

更多推荐