资讯动态

数据中台落地实录:撑起数据接入、处理、分发与服务

发布时间:2026/10/4 8:40:16 来源:尧图企业网站定制
一、起因不想再给每个项目重写一遍数据对接我们内部有套平台叫数据中台。名字听着挺大干的活其实很朴素——把“数据从哪来、怎么加工、给谁用”这三件事从各个业务项目里抽出来收敛到一套平台上统一做。在这之前是什么状态上游有基础数据平台、区域汇聚平台还有各个单位散落在本地的一堆文件下游十几个业务系统各拉各的数据、各写各的清洗脚本。同一份原始数据A 项目解析一遍B 项目又解析一遍出来的格式还不一样。运维最怕的就是上游某天改了接口下游一串系统跟着挂。所以我们给自己定的目标只有十二个字一次接入、统一处理、按需分发、标准服务。平台跑在 K8s 上中间件主力是 PostgreSQL、Redis、RabbitMQ整条链路拆成接入、处理、分发、服务四段。下面把架构、选型和几段关键代码摊开讲顺带说说踩过的坑。二、整体架构先上一张全局图。整体是分层的最底下是数据源往上依次是接入层、处理层、存储层、服务与分发层右侧一列是贯穿全链路的监控。图 1 数据中台总体架构一句话概括数据流向接入层负责把数据“搬进来”处理层负责把数据“加工成产品”分发层把产品“推出去”服务层让用户“自己来取”。图 2 数据流向接入 → 处理 → 分发 / 服务三、技术选型为什么是这几个家伙选型这块我们内部吵过几轮最后落地如下。K8s容器编排——这是最没争议的一条。平台要同时跑接入、处理、分发、网关、调度、监控一堆服务还要能按需扩容。用 K8s 之后滚动升级、副本自愈、资源隔离这些都不用自己造。集群拆成计算和存储两部分控制节点 3 个起步工作节点和存储节点各自独立避免一个业务把整台机器的 IO 打满。PostgreSQL关系型存储——存元数据、任务定义、用户权限、数据目录、操作日志这类强结构、强一致的数据。选 PG 而不是 MySQL主要是看中它对 JSONB 的支持像插件参数、批次参数这种半结构化的东西直接塞 JSONB 字段查询还能走索引省掉一堆宽表。元数据表在平台里是“导航地图”必须靠谱。Redis缓存与协调——三件事热点接口的结果缓存、定时任务的分布式锁、以及接口限流。都是内存级的小操作但对延迟和原子性要求高。RabbitMQ消息队列——采集和处理是典型的“生产快、消费慢”。上游文件一到就是一批处理又要解析又要算同步做必然堵。用 MQ 把任务投递和执行解耦再配合 prefetch 控制执行器压力是整个平台能稳住的关键。对象存储 / 分布式文件——原始文件、中间文件、产品文件都不小尤其是格点类数据动辄几百 MB 到 GB。这类非结构化数据不适合塞数据库走对象存储或者分布式文件系统数据库里只存路径和校验值。顺便说一句关于消息队列我们一开始也考虑过 Kafka。后来发现我们的场景是“任务级”的投递消息量不算极端但对可靠投递、死信、优先级要求更明确RabbitMQ 更顺手。技术选型这东西合适比时髦重要。四、接入层把多源异构数据拉进来接入这块的核心矛盾是源太多、协议太杂、更新频率还不一样。我们的做法是抽象出“采集插件 任务定义”两层。插件负责“怎么拉”——接口、数据湖挂载、FTP/SFTP/SMB/NFS 各写各的插件任务负责“什么时候拉、拉哪些、拉完放哪”——用 Cron 表达式描述频率。任务表长这样-- 数据源CREATE TABLE t_data_source (id BIGSERIAL PRIMARY KEY,code VARCHAR(64) NOT NULL UNIQUE, -- 数据源编码name VARCHAR(128) NOT NULL,global_params JSONB, -- 全局参数enabled BOOLEAN DEFAULT TRUE,created_at TIMESTAMPTZ DEFAULT now());-- 采集任务CREATE TABLE t_collect_task (id BIGSERIAL PRIMARY KEY,source_id BIGINT NOT NULL REFERENCES t_data_source(id),data_code VARCHAR(64) NOT NULL,cron_expr VARCHAR(64) NOT NULL, -- 采集频率plugin_code VARCHAR(64) NOT NULL, -- 采集插件plugin_params JSONB,executor_id BIGINT, -- 代码执行器file_format VARCHAR(32),time_zone VARCHAR(64) DEFAULT Asia/Shanghai,retain_days INT DEFAULT 30, -- 超过时长自动清理enabled BOOLEAN DEFAULT TRUE,updated_at TIMESTAMPTZ DEFAULT now());CREATE INDEX idx_task_enabled ON t_collect_task(enabled);插件统一走一个 SPI 接口Java 和 Python 都能实现平台按 code 动态装载public interface DataPlugin {/** 插件类型COLLECT / PROCESS / DISPATCH */PluginType type();/** 插件编码全局唯一 */String code();/** 执行入口返回本次产出的文件或数据引用 */PluginResult execute(PluginContext context) throws PluginException;}一个 HTTP 拉取的采集插件Java大致长这样Componentpublic class HttpCollectPlugin implements DataPlugin {Override public PluginType type() { return PluginType.COLLECT; }Override public String code() { return collect_http_v1; }Overridepublic PluginResult execute(PluginContext ctx) throws PluginException {String url ctx.param(url);int timeoutMs ctx.param(timeoutMs, 30_000);String saveDir ctx.workingDir();ListFileRef files new ArrayList();RequestConfig cfg RequestConfig.custom().setConnectTimeout(timeoutMs).setSocketTimeout(timeoutMs).build();try (CloseableHttpClient client HttpClients.custom().setDefaultRequestConfig(cfg).build()) {HttpGet get new HttpGet(url);ctx.headers().forEach(get::setHeader);try (CloseableHttpResponse resp client.execute(get)) {if (resp.getCode() / 100 ! 2) {throw new PluginException(拉取失败, status resp.getCode());}Path target Paths.get(saveDir, ctx.taskCode() _ ctx.timestamp());Files.copy(resp.getEntity().getContent(), target,StandardCopyOption.REPLACE_EXISTING);files.add(FileRef.of(target, DigestUtils.md5Hex(target)));}return PluginResult.of(files);} catch (IOException e) {throw new PluginException(采集异常: e.getMessage(), e);}}}本地挂载目录的采集用 Python 写更省事# -*- coding: utf-8 -*-import os, glob, hashlibfrom datetime import datetimedef collect(params: dict) - dict:params 由平台在任务里配置pattern / since / working_dirpattern params[pattern] # 例如 /mnt/upstream/*.ncsince params.get(since) # 上次水位时间避免重复拉out_dir params[working_dir]files []for f in sorted(glob.glob(pattern)):mtime datetime.fromtimestamp(os.path.getmtime(f))if since and mtime datetime.fromisoformat(since):continuemd5 hashlib.md5(open(f, rb).read()).hexdigest()files.append({path: f,name: os.path.basename(f),mtime: mtime.isoformat(),md5: md5,})return {code: 0, count: len(files), files: files}调度器只管投消息不管干活。这点很重要——调度器和执行器解耦之后采集慢不会把调度线程拖死Scheduled(cron ${ingest.scheduler.cron})public void dispatchDueTasks() {ListCollectTask tasks taskMapper.findDue(Instant.now());for (CollectTask t : tasks) {rabbitTemplate.convertAndSend(data.ingest.exchange,collect. t.getPriority(),CollectMessage.of(t),msg - {msg.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT); // 持久化return msg;});taskMapper.markQueued(t.getId()); // 标记防止重复投递}}消费端Python worker# worker.py —— 消费采集任务并执行插件import pika, jsonfrom plugins import load_plugindef on_message(ch, method, props, body):task json.loads(body)try:plugin load_plugin(task[plugin_code])result plugin.execute(task[params])publish_done(task, result)ch.basic_ack(delivery_tagmethod.delivery_tag)except Exception as e:# 失败不重回队列由死信队列接管最多重试 3 次ch.basic_nack(delivery_tagmethod.delivery_tag, requeueFalse)conn pika.BlockingConnection(pika.ConnectionParameters(rabbitmq))ch conn.channel()ch.basic_qos(prefetch_count8) # 按执行器能力限流ch.basic_consume(queuedata.ingest.collect, on_message_callbackon_message)ch.start_consuming()五、处理层插件化 流程编排处理层的目标是把原始数据变成“能用的产品”。这块分两步先解析再加工。解析要面对一堆格式——二进制格点、NetCDF、JSON、CSV、文本、图像每种都有自己的读法。加工则包括插值、时空剪裁、要素裁剪、格式转换、图像生成这些常见动作。我们把它做成了可编排的流程每个“中间数据”下挂若干处理节点节点分两种模式。1:1——一个输入对一个输出互不影响满足条件就触发n:1——多个输入凑齐一批才触发比如三个来源的数据要合到一起才算完整。流程定义表CREATE TABLE t_process_flow (id BIGSERIAL PRIMARY KEY,mid_code VARCHAR(64) NOT NULL, -- 中间数据编码source_type SMALLINT NOT NULL, -- 1原始数据 2中间数据source_code VARCHAR(64) NOT NULL,plugin_code VARCHAR(64) NOT NULL,plugin_params JSONB,mode VARCHAR(8) NOT NULL DEFAULT 1:1,batch_no VARCHAR(32), -- n:1 时用于聚合match_pattern VARCHAR(256), -- 批次匹配正则enabled BOOLEAN DEFAULT TRUE);一个“时空剪裁 格式转换”的处理插件Pythonimport xarray as xrdef process(inputs: list, ctx: dict) - list:bbox ctx[bbox] # [lon_min, lat_min, lon_max, lat_max]time_range ctx[time_range] # [start, end]elements ctx[elements] # 需要的字段out_dir ctx[working_dir]outputs []for path in inputs:ds xr.open_dataset(path)ds ds.sel(timeslice(*time_range))if elements:ds ds[[e for e in elements if e in ds.data_vars]]ds ds.sel(longitudeslice(bbox[0], bbox[2]),latitudeslice(bbox[1], bbox[3]))out f{out_dir}/{ctx[name]}.ncds.to_netcdf(out)outputs.append({path: out, format: nc})return outputsn:1 的批次聚合我们用 Redis 的 List 攒文件攒够数就触发public void onFileReady(String midCode, FileRef file) {FlowDef flow flowCache.get(midCode);if (1:1.equals(flow.getMode())) {execute(flow, List.of(file));return;}String batchKey batch: midCode : flow.matchBatch(file.getName());redis.opsForList().rightPush(batchKey, file.getPath());Long size redis.opsForList().size(batchKey);if (size ! null size flow.getExpectSize()) {ListString paths redis.opsForList().range(batchKey, 0, -1);execute(flow, paths.stream().map(FileRef::of).toList());redis.delete(batchKey);}}这样设计的好处是新增一类数据处理不用改平台代码写个插件 配条流程就行。用我们同事的话说叫“积木式”扩展。六、分发层多协议主动推送不是所有下游都愿意调接口有些单位就是习惯收文件。所以分发层要支持多种协议FTP、SFTP、SCP、HTTP(S)、S3按客户和场景来选。分发的模型是“三层”客户 → 接收主机 → 区域链路。客户是业务上的概念比如“某某项目”接收主机是客户侧实际收数据的机器一个客户可以有多个区域链路描述网络怎么走——目标主机如果在平台不能直达的网络里就通过代理节点逐跳转发。图 3 数据分发链路一个 SFTP 推送插件Pythonimport paramikodef dispatch(inputs: list, ctx: dict) - dict:t paramiko.Transport((ctx[host], ctx.get(port, 22)))t.connect(usernamectx[user], passwordctx[password])sftp paramiko.SFTPClient.from_transport(t)sent []for f in inputs:name ctx[rename](f) # 按正则重命名remote f{ctx[remote_dir]}/{name}sftp.put(f, remote)sent.append(remote)sftp.close()t.close()return {code: 0, sent: sent}这里有个经验分发最好做成“可重试 幂等”。网络抖动导致半截文件的情况太常见了我们的做法是先传临时名校验 md5 一致后再改名落地下游拿到的一定是完整文件。七、服务层API 网关与缓存服务层是给“主动来取”的用户用的——业务系统、第三方平台、还有现在越来越多的智能体。对外统一走网关做鉴权、限流、路由和协议转换。网关路由配置Spring Cloud Gatewayspring:cloud:gateway:routes:- id:>public ResponseEntitybyte[] getData(String apiKey, String dataCode, String params) {String cacheKey api: dataCode : DigestUtils.md5Hex(params);byte[] cached (byte[]) redis.opsForValue().get(cacheKey);if (cached ! null) {return ResponseEntity.ok().header(X-Cache, HIT).body(cached);}byte[] data dataService.fetch(dataCode, params);redis.opsForValue().set(cacheKey, data, Duration.ofMinutes(10));return ResponseEntity.ok().header(X-Cache, MISS).body(data);}限流用一段 Lua 保证原子性-- KEYS[1]限流key ARGV[1]窗口秒 ARGV[2]阈值local n redis.call(INCR, KEYS[1])if n 1 then redis.call(EXPIRE, KEYS[1], ARGV[1]) endif n tonumber(ARGV[2]) then return 0 endreturn 1对外接口我们控制得比较细按数据粒度、调用频次、产品加工深度分档不同客户拿到的权限不一样。既保证基础数据能用又不至于被无节制地薅。八、部署K8s 上跑中间件集群按用途拆成三块互不打架节点类型大致规格用途控制节点16 核 / 32GB / NVMe SSD集群控制面工作节点128 核 / 1TB 内存 / NVMe SSD业务服务、采集、处理存储节点64 核 / 256GB / 大容量 SSD HDD数据库、对象存储、文件服务中间件PostgreSQL、Redis、RabbitMQ、配置中心、调度中心统一以容器方式部署用 ConfigMap 管配置apiVersion: apps/v1kind: Deploymentmetadata:name:>apiVersion: autoscaling/v2kind: HorizontalPodAutoscalermetadata:name:>public void runOnce(String jobKey, Runnable job) {String token UUID.randomUUID().toString();Boolean ok redis.opsForValue().setIfAbsent(lock:job: jobKey, token, Duration.ofMinutes(5));if (!Boolean.TRUE.equals(ok)) {return; // 别的实例在跑直接退出}try {job.run();} finally {unlock(lock:job: jobKey, token); // Lua 原子删除防误删}}6. 消息堆积要提前设阈值。RabbitMQ 队列一旦积压光看面板没用。我们给每个队列设了长度告警配合消费端的 prefetch 和死信队列异常能第一时间发现。十、小结回头看不复杂无非就是把接入、处理、分发、服务四件事标准化、插件化、自动化。但真做起来价值主要在几个地方插件化让能力可以“即插即用”新数据源、新处理算法、新分发协议都不用动核心代码流程编排把数据加工变成可视化配置业务同学也能参与不用每次拉开发K8s 中间件撑住了弹性和稳定机器坏一台不影响整体元数据和监控让数据“看得见、追得到、管得住”这才是平台和一堆脚本的本质区别。如果你们也在做类似的数据平台我的建议是先把元数据模型想清楚再动手写采集。元数据是整个平台的骨架骨架歪了后面补得越勤越痛苦。代码片段都是脱敏后的示例思路可以直接借鉴。有踩过类似坑的朋友欢迎评论区交流。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑