简介:这份《中国移动省级数据共享平台功能规范》是面向省级移动通信网络数据平台建设者、数据治理与运维人员的标准化文档,用于指导数据集成、管理与共享系统的规划、建设和运行维护。包内共1个docx文件,约127KB,正文以标准条款加目录结构呈现,含范围、规范性引用文件、术语定义与缩略语等基础章节。内容围绕数据统一采集、统一存储、统一处理、统一共享、统一管理五大模块展开,细分OMC数据采集与其他网管系统接入、存储能力与存储管理、数据处理管道设计与作业管理、数据目录资产与数据订阅、共享数据SLO/SLI管理、数据质量与元数据管理、数据模型管理及数据分级、加密脱敏、安全审计等要求;同时覆盖系统日志、业务监控、用户分权分域与统一认证、系统门户等运营管理内容。已有768人学习,适合数据平台架构师、数据治理工程师及运营商数字化项目人员对照查阅,快速理清平台功能边界与规范落地要点。
1. 从「每个网管都自己采一遍」说起:省级数据共享平台到底规范了什么
一个省公司 O 域往往并行跑着性能、告警、资源、工单、拨测、DPI 好几套网管系统,每套系统都觉得自己缺原始数据,于是各自找 OMC 开北向账号、各自写 MML 脚本、各自解析同一批性能文件,最后建出来的宽表口径还不一样——同一个小区的小时级流量,两个系统能差出 8%。中国移动省级数据共享平台功能规范要解决的就是这个局面:把采集、存储、处理、共享、管理五段能力收拢成一套平台,域内其他系统不再直连底层网元,而是通过平台取数。规范里明确了采集范围、存储能力、处理管道、数据目录资产、SLI/SLO、数据分级这些硬约束,适合省公司数据平台建设方、网管系统对接方,以及被派去啃这份文档做落地方案的人。
2. 数据统一采集与存储:接口矩阵、采集通道与分层落库
规范第 5、6 章把采集和存储拆成两段,但真正落地时这两段耦合很紧——采集通道决定了落库的粒度,落库的分区设计又反过来约束采集的批次。这一章按「接口怎么选—采集怎么写—数据怎么落」的顺序拆。
2.1 采集范围与接口选型矩阵
规范列了 Kafka、JDBC、文件、RESTful、FTP/SFTP、SDTP、MML、Corba、SNMP、Socket 一长串接口,但没人会全用上。选型的核心判断只有两个:数据是「推」过来还是「拉」过来,以及时效要求落在哪个档位。
| 数据类别 | 典型来源 | 推荐接口 | 时效档 |
|---|---|---|---|
| 资源数据(网元、小区、拓扑) | OMC、资源网管 | JDBC / RESTful / 文件 | T+1 |
| 性能数据(KPI/KQI 计数器) | OMC 北向 | FTP/SFTP 文件、SDTP | 15 分钟 |
| 告警数据 | OMC、告警网管 | SNMP Trap、Corba、Kafka | 秒级 |
| 工单数据 | 工单系统 | RESTful / JDBC | 分钟级 |
| DPI、上网日志 | DPI 系统 | Kafka | 准实时 |
| 拨测数据 | 拨测系统 | 文件 / RESTful | 分钟级 |
| 交换侧指令数据 | 程控交换系统 | MML | 按需 |
判断逻辑很直白:大批量、允许分钟级延迟的性能文件走 SFTP 拉取最稳,因为文件本身带生成时间和记录数,天然可做完整性校验;告警这种偶发高频事件必须走推模式,SNMP Trap 或 Kafka 都行,用轮询一定漏;MML 是程控交换时代留下的人机对话语言,它的特点是「有会话状态」,一条指令的执行结果依赖前一条指令的返回,采集器必须维持长连接而不能每条指令新建会话。
2.2 OMC 北向接口采集的实现要点
性能文件采集最容易被忽略的是幂等。同一批文件在重试时会被重复拉取,如果下游直接 append,指标会被放大。常见做法是在采集侧维护一张本地检查点表,用「文件名 + 文件大小 + MD5」三个字段做主键去重。
# omc_perf_pull.py —— OMC 北向性能文件断点续传采集 import hashlib, os, sqlite3, time from pathlib import Path import paramiko CKPT_DB = "/data/omc_ckpt.db" def init_ckpt(): conn = sqlite3.connect(CKPT_DB) conn.execute("""CREATE TABLE IF NOT EXISTS ckpt( fname TEXT PRIMARY KEY, -- 远端文件名,作为唯一键 size INTEGER, -- 文件字节数,用于发现同名不同内容 md5 TEXT, -- 内容指纹,落库前计算 state TEXT, -- pulled / loaded / failed ts INTEGER)""") conn.commit() return conn def md5_of(path, chunk=1 << 20): h = hashlib.md5() with open(path, "rb") as f: for blk in iter(lambda: f.read(chunk), b""): h.update(blk) return h.hexdigest() def pull(host, port, user, keyfile, remote_dir, local_dir, pattern, batch=50): conn = init_ckpt() ssh = paramiko.SSHClient() ssh.set_missing_host_key_policy(paramiko.AutoAddPolicy()) ssh.connect(hostname=host, port=port, username=user, key_filename=keyfile, look_for_keys=False, timeout=30) sftp = ssh.open_sftp() sftp.chdir(remote_dir) files = [f for f in sftp.listdir() if f.startswith(pattern)] for fname in sorted(files)[:batch]: st = sftp.stat(fname) row = conn.execute("SELECT size, state FROM ckpt WHERE fname=?", (fname,)).fetchone() if row and row[0] == st.st_size and row[1] == "loaded": continue # 已完整入库,跳过 local = os.path.join(local_dir, fname) sftp.get(fname, local) m = md5_of(local) conn.execute("INSERT OR REPLACE INTO ckpt VALUES (?,?,?,?,?)", (fname, st.st_size, m, "pulled", int(time.time()))) conn.commit() sftp.close(); ssh.close() if __name__ == "__main__": pull(host="10.x.x.x", port=22, user="north", keyfile="/etc/keys/omc_rsa", remote_dir="/bns/north/perf/20240514", local_dir="/data/landing/perf", pattern="PM_")逻辑上是「先比对、后下载、再登记」:SELECT命中且状态为loaded就跳过,避免重跑时重复落库。参数方面,pattern按各省 OMC 的性能文件命名规则填,一般是「网元类型 + 时间戳」,不要用*全量匹配,否则一个目录几万文件会把listdir拖死;batch控制在 50~200,太大容易在 SFTP 会话超时前传不完;keyfile用密钥而不是密码,规范里对采集账号的安全要求通常禁止口令登录;timeout=30是连接超时,不代表传输超时,大文件传输要单独设sftp.get的超时或用getfo分块读。
落地的经验是:ckpt表状态字段一定要有中间态。直接从pulled跳到loaded,一旦加载失败就丢了线索;保留failed并记录次数,采集管理模块才有东西可看。
告警侧走 Kafka 时是另一种写法:
# alarm_consume.py —— 告警流消费,手动提交位点保证至少一次 from kafka import KafkaConsumer import json consumer = KafkaConsumer( "omc-alarm-topic", bootstrap_servers=["kafka1:9092", "kafka2:9092"], group_id="province-dsp-alarm", auto_offset_reset="earliest", enable_auto_commit=False, # 关掉自动提交 max_poll_records=500, value_deserializer=lambda b: json.loads(b.decode("utf-8")), ) for batch in consumer: # 先落 HBase/ODS,成功后再提交位点 write_ods(batch.value) consumer.commit()enable_auto_commit=False是关键,开了自动提交就等于接受丢数据;max_poll_records要跟下游写入的批处理能力匹配,设太大反而会因为处理超时被踢出消费组,触发 rebalance。
2.3 分层存储与介质映射
规范第 6 章强调存储能力与存储管理,实际落地就是一张分层映射表。省级平台的数据量级通常是 PB 级,不可能用一种介质兜住所有查询模式。
| 分层 | 存储介质 | 典型表 | 保留策略 |
|---|---|---|---|
| ODS 原始层 | HDFS + Hive 外部分区表 | ods_perf_raw | 6~12 个月 |
| 明细层 | Hive 分区表、HBase | dwd_perf_15min | 12 个月 |
| 轻度汇总层 | Hive / MPP | dws_cell_hour | 24 个月 |
| 应用层 | MPP / RMDB | ads_region_kpi | 长期 |
| 维表与缓存 | Redis | dim_ne, dim_cell | 随源更新 |
| 元数据 | MySQL / PostgresDB | meta_table, meta_column | 长期 |
Hive 分区表按「天 + OMC 标识」两级分区最实用,因为绝大多数查询带省份和日期条件:
CREATE EXTERNAL TABLE IF NOT EXISTS ods_perf_raw ( ne_id STRING COMMENT '网元标识', counter_id STRING COMMENT '计数器编码', collect_ts BIGINT COMMENT '采集时间戳(秒)', val DOUBLE COMMENT '计数器值' ) COMMENT 'OMC 性能原始数据' PARTITIONED BY (dt STRING COMMENT '日期 yyyyMMdd', omc_id STRING COMMENT 'OMC 标识') STORED AS ORC TBLPROPERTIES ('orc.compress'='SNAPPY'); -- 按分区挂载,采集完成一个目录挂一个分区 ALTER TABLE ods_perf_raw ADD IF NOT EXISTS PARTITION (dt='20240514', omc_id='OMC-A01') LOCATION '/data/ods/perf/dt=20240514/omc_id=OMC-A01';用EXTERNAL是因为 ODS 数据可能被其他引擎直接读 HDFS 路径,删表不该删数据;ORC + SNAPPY是性能与压缩率的平衡点;分区挂载跟采集任务绑定,采集到一个完整目录就ADD PARTITION一次,这样分区可见性就等于数据完整性信号——查不到分区,说明那批数据没采到。
2.4 数据源信息管理与采集管理
规范把「数据源信息管理」单列一节,原因是采集任务的元信息必须集中。注册一个数据源至少要落这些字段:数据源编码、类型(OMC / 其他网管)、接入协议、连接地址、认证方式、责任人、所属网络域、采集周期、最近水位时间。采集调度器按「数据源编码 + 采集周期」生成实例,任务失败时能反查到责任人,这比全平台一个告警群有效得多。
3. 数据处理管道设计到实例化:批流一体的作业编排
规范第 7 章用「管道设计—管道实例化—管道及作业管理」三步描述处理能力,这套抽象的好处是把「逻辑」和「运行实例」解耦:一个清洗逻辑只设计一次,然后按省份、按网元类型实例化出几十个作业。
3.1 管道设计:从算子链到 DAG
管道在规范语义里是「源 + 处理算子 + 汇」的有向无环图。设计阶段只画逻辑图,不落具体资源;实例化阶段才绑定并行度、资源队列、检查点周期。常见的四类算子:过滤(丢无效计数器)、映射(编码转名称)、聚合(按时间窗汇总)、关联(挂维表补小区归属)。
以性能数据 15 分钟粒度汇总为例,Flink SQL 写得最省事:
-- 源:Kafka 中的性能明细流 CREATE TABLE perf_src ( ne_id STRING, counter_id STRING, val DOUBLE, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '30' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'dsp-perf-detail', 'properties.bootstrap.servers' = 'kafka1:9092,kafka2:9092', 'properties.group.id' = 'dsp-perf-agg', 'scan.startup.mode' = 'group-offsets', 'format' = 'json' ); -- 汇:Hive 汇总表 CREATE TABLE dws_cell_15min ( ne_id STRING, counter_id STRING, win_start TIMESTAMP(3), avg_val DOUBLE, max_val DOUBLE, PRIMARY KEY (ne_id, counter_id, win_start) NOT ENFORCED ) WITH ( 'connector' = 'hive', 'hive-version' = '3.1.2', 'sink.partition-commit.policy.kind' = 'metastore,success-file' ); INSERT INTO dws_cell_15min SELECT ne_id, counter_id, TUMBLE_START(event_time, INTERVAL '15' MINUTE) AS win_start, AVG(val), MAX(val) FROM perf_src GROUP BY ne_id, counter_id, TUMBLE(event_time, INTERVAL '15' MINUTE);WATERMARK ... - INTERVAL '30' SECOND决定了容忍多久的乱序,设太大窗口迟迟不关闭、结果延迟高,设太小迟到数据会被丢弃;scan.startup.mode = group-offsets保证作业重启从消费组位点续跑,而不是从头重刷;sink.partition-commit.policy.kind里的metastore让 Hive 元数据跟文件一起提交,避免出现「文件写了但分区查不到」的中间态。
3.2 管道实例化与参数注入
实例化的本质是模板渲染。把上面的 SQL 抽成带占位符的模板,实例化时注入参数:
{ "pipeline_code": "PERF_AGG_15MIN", "template_id": "TPL_PERF_TUMBLE", "params": { "source_topic": "dsp-perf-detail", "sink_table": "dws_cell_15min", "window_size": "15m", "parallelism": 8, "checkpoint_interval_ms": 60000, "state_backend": "rocksdb", "restart_strategy": "fixed-delay:3:30s" }, "owner": "perf-domain", "slo_ref": "SLO-PERF-15MIN" }参数里parallelism跟 Kafka 分区数对齐最好,分区数 16 而并行度 8,每个算子实例要消费两个分区,吞吐上不去;checkpoint_interval_ms60 秒是流作业的常见起点,状态大就拉长到 180 秒,否则检查点本身成为负担;state_backend选rocksdb的前提是状态超过内存容量,小状态用hashmap更快;restart_strategy的fixed-delay:3:30s表示最多重启 3 次、每次间隔 30 秒,超过就置为失败等人工介入——生产环境不建议用无限制重启,会把真实故障掩盖成「一直在重试」。
3.3 管道及作业管理
实例化之后要能被统一管理,最小可用的能力是四件事:启停、并行度调整、位点查看、血缘追溯。批作业一般交给调度系统按依赖触发,流作业常驻运行、由平台做健康检查。批作业侧一个典型依赖配置:
# 用调度系统的命令行提交一个带依赖的批作业 dolphinscheduler task create \ --project dsp_province \ --name dwd_perf_15min_load \ --type HIVE \ --raw "INSERT OVERWRITE TABLE dwd_perf_15min PARTITION(dt='${bizdate}') SELECT * FROM ods_perf_raw WHERE dt='${bizdate}'" \ --pre dws_cell_15min_agg \ --timeout 3600 \ --retry-times 2--pre声明前置依赖,缺了它下游会在上游分区还没挂载时开跑,跑出空表;--timeout 3600防止任务卡死占住工作流实例;--retry-times 2只对可重试的失败有意义,像权限错误这种重试多少次都一样。
3.4 数据质量核查嵌入处理链路
规范第 9.1.3 节的「数据质量核查分析」如果只在事后跑,问题发现时已经污染了下游。务实做法是把轻量核查卡在管道汇之前,重核查放 T+1 批处理。轻量核查的例子:
-- 15 分钟窗口的完整性核查:计数器数量不应低于基线的 95% SELECT win_start, COUNT(DISTINCT counter_id) AS cnt, COUNT(DISTINCT counter_id) * 1.0 / 900 AS ratio -- 900 为基线计数器数 FROM dws_cell_15min WHERE win_start >= CURRENT_TIMESTAMP - INTERVAL '1' HOUR GROUP BY win_start HAVING COUNT(DISTINCT counter_id) * 1.0 / 900 < 0.95;HAVING里放阈值判断,命中即产出质量告警记录,写进质量结果表供监控模块读取。900这个基线值应该来自元数据里的「应有计数器清单」而不是硬编码,否则网元版本升级后误报会淹掉真告警。
4. 数据共享与统一管理:目录资产、订阅、SLO 与安全分级
规范第 8 章管共享、第 9 章管治理,这两章在实现上共用同一份元数据仓库。数据目录资产是入口,订阅是授权动作,共享是通道,SLO 是对外承诺,安全分级贯穿全流程。
4.1 数据目录资产注册与数据地图
一条目录资产记录要能回答四个问题:这是什么数据、在哪里、谁能看、质量如何。注册时的字段建议包含:资产编码、资产名称、所属主题域、数据分层、物理位置(库表或 API 路径)、更新频率、owner、密级、质量分、版本号。数据地图是把这些资产按主题域和血缘关系铺开,让使用方按「网络域—业务域—资产」三级路径找数据,而不是靠人问。
资产变更、注销、版本管理这三个动作必须有审计留痕。变更表结构时先走新版本注册、双版本并行、下游切换、旧版本注销,跳过并行期直接改,下游一定会挂。
4.2 数据订阅与共享通道
规范把共享通道归为三类:批量文件、数据库接口、数据服务 API。选型看使用方的消费能力。
| 通道 | 协议 | 适用场景 | 关键约束 |
|---|---|---|---|
| 批量文件 | FTP/SFTP、SDTP | 大批量、T+1 对账 | 文件命名、校验文件必须成对 |
| 数据库接口 | JDBC | 对方有库、需要 SQL 自由度 | 只开放视图,不给基表权限 |
| 数据服务 API | RESTful/HTTPS | 应用系统实时取数 | 限流、分页、鉴权令牌 |
订阅本身是一个审批流:使用方在门户提订阅申请,选资产、选字段、填用途、填期限,数据 owner 审批,通过后系统按密级自动决定是否需要脱敏视图。规范强调「分权分域按需订阅与共享」,落到实现就是订阅记录里必须有「域」和「权」两个维度,不能只按角色给全量权限。
4.3 SLI/SLO 体系怎么落地
规范第 8.4 节把 SLI 定义为精细测量的服务水平指标,SLO 是用 SLI 描述的期望状态。落地时先定少量 SLI,再给每个共享资产挂 SLO。
| SLI | 定义 | 采集方式 | 示例 SLO |
|---|---|---|---|
| 数据到达及时率 | 按时到达批次数 / 应到达批次数 | 采集任务水位监控 | ≥ 99% |
| 数据完整率 | 实际记录数 / 期望记录数 | 质量核查结果表 | ≥ 99.5% |
| API 可用性 | 成功响应数 / 总请求数 | 网关访问日志 | ≥ 99.9% |
| API P99 延迟 | 99 分位响应耗时 | 网关直方图指标 | ≤ 800ms |
| 共享任务成功率 | 成功共享任务 / 总任务 | 共享任务表 | ≥ 99% |
用指标查询语言把 SLO 写成可告警的表达式:
-- 查询语言中:API 可用性 = 非 5xx 请求占比(5 分钟窗口) sum(rate(dsp_api_requests_total{code!~"5.."}[5m])) / sum(rate(dsp_api_requests_total[5m]))rate(...[5m])取的是每秒速率,用比值消除流量波动的影响;code!~"5.."用正则排除 5xx,注意不要写成排除 4xx,因为鉴权失败属于调用方问题,把它算进可用性会让 SLO 长期虚低。SLO 定完还要配错误预算:可用性目标 99.9%,一个月允许的不可用时间约 43 分钟,预算烧完就该冻结变更、优先修稳定性,这条纪律比指标本身更重要。
4.4 元数据、数据模型与数据标准管理
元数据分技术元数据和业务元数据。技术元数据靠采集器从 Hive Metastore、数据库系统表自动抽取,业务元数据靠数据标准人工维护。规范第 9.2.1 节的数据标准管理,落到系统里就是一张「标准项」表:标准编码、中文名、英文名、数据类型、取值范围、单位、对应安全级别。新建表时字段必须挂标准项,挂不上就说明这个字段还没有标准,先补标准再建表。
数据模型管理覆盖导入、呈现、查询、导出四件事。导入支持从建表语句或建模工具文件解析,呈现用图形化 ER 图,查询支持按表名和字段名模糊搜,导出支持生成建表语句和字段清单。模型的价值在于影响分析:改一个字段前先查它在哪些模型和下游作业里被引用,避免改完才发现有十几个作业在跑。
4.5 数据分级、脱敏与分权分域
数据分级是安全控制的基准。常见分四档:公开、内部、敏感、机密。分级结果要落到字段级而不是表级,因为一张工单表里可能只有手机号是敏感字段。
脱敏实现上,静态脱敏用于共享出去的落地数据,动态脱敏用于 API 返回:
# 手机号脱敏:保留前 3 位后 4 位,中间打码 def mask_msisdn(v: str) -> str: if not v or len(v) != 11: return "***" return v[:3] + "****" + v[-4:] # 身份证脱敏:只留前 6 位地区码 def mask_idcard(v: str) -> str: return v[:6] + "*" * (len(v) - 6) if v else ""逻辑说明:先做长度校验再截取,避免脏数据导致切片错位把明文暴露出来。参数上要注意 11 位是手机号长度假设,遇到带国家码的号码要先归一化再脱敏。分权分域则在下游查询入口强制拼接域条件,比如地市用户只能查本地市数据,这个过滤必须做在数据服务层而不是靠前端传参,前端传什么参数都不可信。
5. 排错与验证:采集断点、SLO 告警与质量核查的定位手法
出问题时,排查顺序应该固定下来,不然每次都在猜。
第一步看采集水位。查检查点表里每个数据源的最大时间戳:
SELECT fname, size, state, ts FROM ckpt WHERE state <> 'loaded' ORDER BY ts DESC LIMIT 20;有pulled但迟迟不到loaded,说明文件拉下来了、加载环节卡住,去看加载作业日志;有文件根本没进表,说明远端目录里还没生成,去找 OMC 侧确认。这一步能区分「采不到」和「采到了没入库」,是最省时间的一次分流。
第二步看分区连续性。Hive 分区断层是下游空结果的常见原因:
SHOW PARTITIONS ods_perf_raw;如果dt=20240513和dt=20240515之间缺了一天,先别怀疑数据,检查那一天有没有挂过分区。MSCK REPAIR TABLE能把 HDFS 上存在但元数据里没有的分区补回来,这是采集任务和元数据不同步时的常用修复手段。
第三步看质量核查结果表,按时间和主题域过滤,看是单点异常还是全量异常。单点异常通常是某个 OMC 或某类网元的问题,全量异常往往是自己改了逻辑,这时候去比对变更记录比看代码快。
第四步看 SLO 与错误预算。可用性 SLO 触发告警时,先分清楚是流量突增还是真实故障:
-- 按接口维度看请求量和错误率,定位是哪个接口拖垮整体 sum by (path) (rate(dsp_api_requests_total{code=~"5.."}[5m])) / sum by (path) (rate(dsp_api_requests_total[5m]))by (path)让结果按接口路径分组,一眼能看出是某个接口还是全部接口。如果只有一个接口错误率高,问题在接口实现;如果全部接口一起抖,先查共享层依赖的存储或认证服务。
一个常被忽略的技巧是给采集任务加「静默期」判断。网络割接、网元版本升级期间数据本来就会缺,这时候刷出来的质量告警全是噪音,值班人员很快会麻木。把割接窗口写进配置,窗口内的缺失记录到日报但不触发告警,窗口外的缺失才告警——这一条能把误报压掉大半。
本文还有配套的精品资源,点击获取