
简介这份资源是Mars实时数据库的完整C#项目源码压缩包面向希望深入理解实时数据库内部实现、提升C#工程能力的开发者与学习者。项目围绕数据采集、存储与分析三大环节展开涵盖数据结构、查询处理、事务管理、并发控制与索引优化等核心模块适合具备一定C#与.NET基础、想研究数据库系统设计的中高级开发者参考。压缩包共1222个文件约28.55MB以705个cs源码文件为主体辅以xaml与razor界面文件、csproj工程配置、proto协议定义、json与cfg配置、png与ico资源及少量sqlite数据文件目录结构清晰便于按模块检索阅读。目前已有125人学习下载。通过研读源码读者可了解C#异步编程在数据采集中的应用、内存数据库与缓存策略的取舍、ADO.NET与LINQ在查询分析中的实践以及事务隔离与并发控制的落地方式是一份兼具学习与借鉴价值的实时数据库实现范例。1. Mars 数据库到底在解决什么问题从数据采集到实时分析的一条链路如果你手头有几十台设备在持续上报温度、压力、电流同时业务侧又要求「刚过去 5 秒的平均值」能立刻查到传统做法通常是采集程序写一份、时序库存一份、离线任务再算一份三套东西各管一段中间靠定时任务搬运。Mars 数据库这个标题指向的正是把这三段合成一条链路采集端把数据推进来存储层按时间组织分析层直接在库内出结果。它适合做设备联网、传感器汇聚、实时看板这类场景的工程师尤其是被「采集归采集、分析归分析」割裂折磨过的人。这一章先把这条链路的边界讲清楚后面几章再落到具体怎么搭、参数怎么调、哪里容易翻车。2. 采集层怎么接把设备数据稳定送进 Mars 数据库2.1 先想清楚采集的三种接入形态设备数据进库绕不开三种形态选错了后面全是补丁。第一种是主动推送设备或网关通过 MQTT、HTTP 把数据发出来采集程序只做接收和转发第二种是轮询拉取针对 Modbus、OPC UA 这类不支持主动上报的工业设备采集程序按周期去读寄存器第三种是文件落地设备先把数据写成 CSV 或二进制文件再由采集程序扫描入库。Mars 数据库作为实时库对这三种形态的承接方式不一样推送型适合走消息队列缓冲轮询型适合在采集程序里做批量攒批文件型适合做增量扫描加断点记录。我一般会先确认设备的协议和上报频率再决定用哪种。如果设备每秒上报一次以上优先推送加队列如果设备是 PLC 且只支持 Modbus就老老实实轮询但轮询周期不要低于设备的最小刷新时间否则读到的全是重复值。这里有个容易被忽略的点采集程序的时间戳一定要用设备侧时间或网关侧时间不要用入库时间否则网络抖动会让数据在时间轴上错位后面做窗口聚合时全是坑。2.2 用 Python 写一个最小采集入库脚本下面这段代码演示的是推送型场景下采集程序从 MQTT 拿到消息后批量写入 Mars 数据库的最小逻辑。不同部署方式下客户端库名字可能不同这里用常见的连接和写入接口示意重点是攒批和异常重试的结构。import json import time import paho.mqtt.client as mqtt # 假设 Mars 提供了 Python 客户端连接参数按实际部署填写 from mars_client import MarsClient BATCH_SIZE 200 # 攒够 200 条再写一次减少网络往返 FLUSH_INTERVAL 1.0 # 即使不够 200 条最多等 1 秒也要写 buffer [] client MarsClient(host127.0.0.1, port8300, databasedevice_metrics) client.connect() def on_message(mqtt_client, userdata, msg): payload json.loads(msg.payload.decode(utf-8)) # 关键时间戳用设备上报的 ts不用本地 time.time() row { device_id: payload[device_id], metric: payload[metric], value: float(payload[value]), ts: int(payload[ts]) } buffer.append(row) if len(buffer) BATCH_SIZE: flush() def flush(): global buffer if not buffer: return try: client.write_points(metrics, buffer) buffer [] except Exception as e: # 写失败不要丢数据保留 buffer 等下次重试 print(write failed, keep buffer:, e) def on_connect(mqtt_client, userdata, flags, rc): mqtt_client.subscribe(device//metrics) mqtt_client mqtt.Client() mqtt_client.on_connect on_connect mqtt_client.on_message on_message mqtt_client.connect(127.0.0.1, 1883, 60) # 主循环里做定时 flush保证低频设备的数据不会一直卡在 buffer while True: mqtt_client.loop(timeout0.1) if buffer and time.time() - last_flush FLUSH_INTERVAL: flush() last_flush time.time()逻辑说明on_message只负责解析和入 buffer不直接写库避免每条消息都触发一次网络请求flush做批量写入失败时保留 buffer 而不是清空这是采集程序不丢数据的关键。参数上BATCH_SIZE和FLUSH_INTERVAL要一起调批量太大延迟高太小写入压力大一般从 200 条 / 1 秒起步再根据设备数量和写入吞吐调整。ts字段必须来自设备侧这是后面做时间窗口分析的前提。2.3 轮询型设备的采集节奏怎么定轮询型设备没有推送能力采集程序要自己控制节奏。常见做法是给每类设备配一个采集周期周期值取设备刷新时间的 1 到 2 倍。比如某注塑机数据采集模块的寄存器 500ms 刷新一次采集周期设 500ms 到 1s 都合理设 100ms 只会读到重复值还白白增加负载。轮询程序里要记录每个设备的最后采集时间避免因为某台设备超时拖慢整个循环。如果设备数量多用线程池或异步 IO 并发拉取但并发数不要超过网关或串口服务器的承载能力否则会出现大面积超时。3. 存储层怎么设计Mars 数据库的表结构和时间分区3.1 表结构宽表还是窄表先看查询模式Mars 数据库存储设备数据表结构设计直接决定后面查询顺不顺手。窄表是一行一个测点字段是 device_id、metric、value、ts宽表是一行一台设备多个测点字段是 device_id、ts、temp、pressure、current。窄表写入灵活新增测点不用改表结构适合测点种类多、变化频繁的场景宽表查询方便一次能取到一台设备多个测点的对齐数据适合测点固定、需要做多指标关联分析的场景。我的经验是如果业务侧经常要「查某台设备某段时间的所有指标」用宽表如果经常要「查某个指标在所有设备上的分布」用窄表。两者也可以混用原始数据进窄表再建物化视图或定时任务聚合成宽表供分析用。Mars 数据库如果支持标签索引把 device_id、metric 建成标签查询时按标签过滤会比全表扫描快很多。3.2 时间分区按天还是按小时看数据量和保留周期实时数据库的数据是随时间无限增长的不做分区查询会越来越慢。常见做法是按时间做分区分区粒度看两个因素每天的数据量和需要保留多久。每天写入量在千万级以内按天分区够用如果每天上亿条按小时分区更合适查询时能更快裁剪掉无关分区。保留周期也要提前定比如原始数据保留 30 天聚合数据保留 1 年过期分区定期删除或归档到对象存储。下面是一个按天分区的建表示意具体语法按 Mars 数据库实际支持的方式调整-- 按天分区保留最近 30 天原始数据 CREATE TABLE metrics ( ts TIMESTAMP NOT NULL, device_id VARCHAR(64) NOT NULL, metric VARCHAR(64) NOT NULL, value DOUBLE, PRIMARY KEY (device_id, metric, ts) ) PARTITION BY RANGE (ts) INTERVAL 1 day RETENTION 30 days;逻辑说明PARTITION BY RANGE按时间范围切分数据INTERVAL 1 day表示每天一个分区RETENTION控制自动清理。参数上分区粒度不要小于查询最常用的时间窗口否则跨分区查询会变多保留周期要跟业务和合规要求对齐删之前确认有没有下游任务还在读。3.3 写入批量与压缩别让存储成为瓶颈Mars 数据库写入时批量提交比单条提交吞吐高一个数量级但批量太大会增加内存占用和失败重试成本。一般建议单批 500 到 2000 条具体看单条数据大小和客户端内存。如果数据是数值型且变化平缓开启列式压缩能显著降低存储占用压缩算法选默认的就行不要为了追求压缩率选高 CPU 开销的算法否则写入会变慢。另外写入时尽量按时间顺序追加乱序写入会触发更多的分区合并操作长期看影响性能。4. 分析层怎么用在 Mars 数据库里做实时聚合与查询4.1 窗口聚合实时看板背后的查询怎么写实时看板要的是「最近 5 分钟每 10 秒的平均值」这类结果用窗口聚合直接在库内算比把原始数据拉到应用层再算快得多。下面是一个按 10 秒窗口聚合的查询示意SELECT time_bucket(10 seconds, ts) AS window_start, device_id, AVG(value) AS avg_value, MAX(value) AS max_value FROM metrics WHERE metric temperature AND ts NOW() - INTERVAL 5 minutes GROUP BY window_start, device_id ORDER BY window_start DESC;逻辑说明time_bucket把时间切成固定窗口WHERE先按时间范围裁剪分区GROUP BY按窗口和设备分组。参数上窗口大小要跟看板刷新频率匹配看板 10 秒刷一次窗口就设 10 秒如果窗口太小数据点少平均值波动大可以适当放大窗口或加滑动窗口。查询里一定要带时间范围条件否则会扫全部分区。4.2 降采样与物化视图让历史查询不拖垮实时写入原始数据查久了会慢常见做法是建降采样任务把秒级数据聚合成分钟级、小时级查询历史趋势时直接查聚合表。Mars 数据库如果支持物化视图可以定义好聚合规则让库自动维护如果不支持就用定时任务定期跑聚合 SQL 写入另一张表。降采样表的保留周期可以比原始表长因为数据量小很多。注意聚合任务要错开写入高峰避免和采集写入抢资源。4.3 查询参数怎么调并发、超时和缓存分析查询和采集写入跑在同一个库上时要限制分析查询的资源占用。常见做法是给分析查询设单独的连接池和超时时间超时时间根据查询复杂度设简单聚合 5 到 10 秒复杂关联 30 秒到 1 分钟。如果库支持查询缓存对看板这类重复查询开启缓存能明显降低负载但缓存过期时间不要设太长否则看板数据会滞后。并发数上分析查询的并发不要超过库的连接上限留一部分给写入。5. 避坑与排查Mars 数据库落地时最容易翻车的五个地方5.1 现象采集程序跑几天后内存暴涨原因buffer 只增不减写入失败后没有上限控制或者 MQTT 消息积压导致 buffer 越堆越大。解决给 buffer 设最大长度超过后要么丢弃最旧数据并告警要么落盘做临时缓冲写入失败要区分可重试和不可重试错误不可重试的直接记日志丢弃避免无限重试。5.2 现象查询最近数据很快查一周前数据特别慢原因时间分区没建对或者查询条件没带时间范围导致全分区扫描。解决确认分区粒度和查询窗口匹配查询 SQL 里强制加时间范围条件必要时在应用层做参数校验不允许不带时间范围的查询打到库上。5.3 现象设备时间戳和入库时间不一致聚合结果对不上原因采集程序用了本地时间入库或者设备时间没同步。解决统一用设备侧时间戳采集程序只做透传如果设备时间不可信在网关侧做一次时间校准并记录校准偏移量方便排查。5.4 现象批量写入偶尔失败重试后数据重复原因批量写入不是幂等的失败后整批重试已经写入的部分又写了一遍。解决写入时带上唯一键或版本号让库做去重或者把批量拆小失败后只重试失败的那部分。Mars 数据库如果支持主键去重建表时把 device_id、metric、ts 设成主键。5.5 现象降采样任务跑完后实时查询变慢原因降采样任务和实时查询抢 IO 和 CPU或者聚合表没建索引。解决降采样任务安排在写入低谷期限制任务并发聚合表按查询模式建索引常用过滤字段放前面。6. 进阶技巧用 Mars 数据库做设备健康度实时评分设备健康度评分是个典型的实时分析场景采集层拿到振动、温度、电流存储层按时间存好分析层用滑动窗口算特征再按规则或模型打分。我一般会这样做先用 1 分钟滑动窗口算每个设备的振动均值和标准差再用 5 分钟窗口算温度变化率然后把特征拼成一行用一个简单的评分函数算出 0 到 100 的分值写入一张评分表供看板查询。-- 每分钟算一次设备健康度特征 INSERT INTO device_health_features SELECT time_bucket(1 minute, ts) AS window_start, device_id, AVG(CASE WHEN metric vibration THEN value END) AS vib_avg, STDDEV(CASE WHEN metric vibration THEN value END) AS vib_std, AVG(CASE WHEN metric temperature THEN value END) AS temp_avg, MAX(CASE WHEN metric current THEN value END) AS current_max FROM metrics WHERE ts NOW() - INTERVAL 2 minutes GROUP BY window_start, device_id;逻辑说明用CASE WHEN把窄表里的不同指标转成宽表特征time_bucket控制窗口WHERE只扫最近两分钟的数据保证实时性。参数上窗口大小和任务执行频率要匹配1 分钟窗口就每分钟跑一次特征字段按评分模型需要增减。评分函数可以先用阈值规则比如振动标准差超过阈值扣分温度变化率超过阈值扣分跑一段时间后再考虑换成模型。验证方法上我会先拿历史数据回放确认评分曲线和设备实际状态对得上再上线实时任务。上线后重点看两个指标评分更新延迟和评分跳变频率。延迟高说明任务跑得慢要优化查询或加资源跳变频繁说明窗口太小或阈值太敏感要调窗口和阈值。这套东西跑顺之后设备异常往往能在评分下降后几分钟内被发现比等人工巡检快得多。我自己踩过的最大坑是早期没做时间戳统一设备侧时间、网关时间、入库时间混着用评分曲线看着正常一出问题就对不上原始数据查了两天才定位到是时间错位。从那以后采集链路里时间戳只认一个来源其他环节一律透传。希望帮到你。本文还有配套的精品资源点击获取