尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

省级数据共享平台落地:OMC采集、批流一体处理与SLI/SLO治理

省级数据共享平台落地:OMC采集、批流一体处理与SLI/SLO治理 简介这份《中国移动省级数据共享平台功能规范》是面向省级移动通信网络数据平台建设者、数据治理与运维人员的标准化文档用于指导数据集成、管理与共享系统的规划、建设和运行维护。包内共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 / 文件T1性能数据KPI/KQI 计数器OMC 北向FTP/SFTP 文件、SDTP15 分钟告警数据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, chunk1 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, batch50): conn init_ckpt() ssh paramiko.SSHClient() ssh.set_missing_host_key_policy(paramiko.AutoAddPolicy()) ssh.connect(hostnamehost, portport, usernameuser, key_filenamekeyfile, look_for_keysFalse, timeout30) 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(host10.x.x.x, port22, usernorth, keyfile/etc/keys/omc_rsa, remote_dir/bns/north/perf/20240514, local_dir/data/landing/perf, patternPM_)逻辑上是「先比对、后下载、再登记」SELECT命中且状态为loaded就跳过避免重跑时重复落库。参数方面pattern按各省 OMC 的性能文件命名规则填一般是「网元类型 时间戳」不要用*全量匹配否则一个目录几万文件会把listdir拖死batch控制在 50200太大容易在 SFTP 会话超时前传不完keyfile用密钥而不是密码规范里对采集账号的安全要求通常禁止口令登录timeout30是连接超时不代表传输超时大文件传输要单独设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_idprovince-dsp-alarm, auto_offset_resetearliest, enable_auto_commitFalse, # 关掉自动提交 max_poll_records500, value_deserializerlambda b: json.loads(b.decode(utf-8)), ) for batch in consumer: # 先落 HBase/ODS成功后再提交位点 write_ods(batch.value) consumer.commit()enable_auto_commitFalse是关键开了自动提交就等于接受丢数据max_poll_records要跟下游写入的批处理能力匹配设太大反而会因为处理超时被踢出消费组触发 rebalance。2.3 分层存储与介质映射规范第 6 章强调存储能力与存储管理实际落地就是一张分层映射表。省级平台的数据量级通常是 PB 级不可能用一种介质兜住所有查询模式。分层存储介质典型表保留策略ODS 原始层HDFS Hive 外部分区表ods_perf_raw612 个月明细层Hive 分区表、HBasedwd_perf_15min12 个月轻度汇总层Hive / MPPdws_cell_hour24 个月应用层MPP / RMDBads_region_kpi长期维表与缓存Redisdim_ne, dim_cell随源更新元数据MySQL / PostgresDBmeta_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.compressSNAPPY); -- 按分区挂载采集完成一个目录挂一个分区 ALTER TABLE ods_perf_raw ADD IF NOT EXISTS PARTITION (dt20240514, omc_idOMC-A01) LOCATION /data/ods/perf/dt20240514/omc_idOMC-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 节的「数据质量核查分析」如果只在事后跑问题发现时已经污染了下游。务实做法是把轻量核查卡在管道汇之前重核查放 T1 批处理。轻量核查的例子-- 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大批量、T1 对账文件命名、校验文件必须成对数据库接口JDBC对方有库、需要 SQL 自由度只开放视图不给基表权限数据服务 APIRESTful/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;如果dt20240513和dt20240515之间缺了一天先别怀疑数据检查那一天有没有挂过分区。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)让结果按接口路径分组一眼能看出是某个接口还是全部接口。如果只有一个接口错误率高问题在接口实现如果全部接口一起抖先查共享层依赖的存储或认证服务。一个常被忽略的技巧是给采集任务加「静默期」判断。网络割接、网元版本升级期间数据本来就会缺这时候刷出来的质量告警全是噪音值班人员很快会麻木。把割接窗口写进配置窗口内的缺失记录到日报但不触发告警窗口外的缺失才告警——这一条能把误报压掉大半。本文还有配套的精品资源点击获取
返回列表