一、起因:不想再给每个项目重写一遍数据对接
我们内部有套平台,叫数据中台。名字听着挺大,干的活其实很朴素——把“数据从哪来、怎么加工、给谁用”这三件事,从各个业务项目里抽出来,收敛到一套平台上统一做。
在这之前是什么状态?上游有基础数据平台、区域汇聚平台,还有各个单位散落在本地的一堆文件;下游十几个业务系统,各拉各的数据、各写各的清洗脚本。同一份原始数据,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)大致长这样:
@Component
public class HttpCollectPlugin implements DataPlugin {
@Override public PluginType type() { return PluginType.COLLECT; }
@Override public String code() { return "collect_http_v1"; }
@Override
public PluginResult execute(PluginContext ctx) throws PluginException {
String url = ctx.param("url");
int timeoutMs = ctx.param("timeoutMs", 30_000);
String saveDir = ctx.workingDir();
List<FileRef> 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, hashlib
from datetime import datetime
def collect(params: dict) -> dict:
"""params 由平台在任务里配置:pattern / since / working_dir"""
pattern = params["pattern"] # 例如 /mnt/upstream/*.nc
since = 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):
continue
md5 = 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() {
List<CollectTask> 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, json
from plugins import load_plugin
def 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_tag=method.delivery_tag)
except Exception as e:
# 失败不重回队列,由死信队列接管,最多重试 3 次
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
conn = pika.BlockingConnection(pika.ConnectionParameters("rabbitmq"))
ch = conn.channel()
ch.basic_qos(prefetch_count=8) # 按执行器能力限流
ch.basic_consume(queue="data.ingest.collect", on_message_callback=on_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
);
一个“时空剪裁 + 格式转换”的处理插件(Python):
import xarray as xr
def 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(time=slice(*time_range))
if elements:
ds = ds[[e for e in elements if e in ds.data_vars]]
ds = ds.sel(longitude=slice(bbox[0], bbox[2]),
latitude=slice(bbox[1], bbox[3]))
out = f"{out_dir}/{ctx['name']}.nc"
ds.to_netcdf(out)
outputs.append({"path": out, "format": "nc"})
return outputs
n: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()) {
List<String> 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 推送插件(Python):
import paramiko
def dispatch(inputs: list, ctx: dict) -> dict:
t = paramiko.Transport((ctx["host"], ctx.get("port", 22)))
t.connect(username=ctx["user"], password=ctx["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 Gateway):
spring:
cloud:
gateway:
routes:
- id:>public ResponseEntity<byte[]> 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]) end
if n > tonumber(ARGV[2]) then return 0 end
return 1
对外接口我们控制得比较细:按数据粒度、调用频次、产品加工深度分档,不同客户拿到的权限不一样。既保证基础数据能用,又不至于被无节制地薅。
八、部署:K8s 上跑中间件
集群按用途拆成三块,互不打架:
节点类型 | 大致规格 | 用途 |
控制节点 | 16 核 / 32GB / NVMe SSD | 集群控制面 |
工作节点 | 128 核 / 1TB 内存 / NVMe SSD | 业务服务、采集、处理 |
存储节点 | 64 核 / 256GB / 大容量 SSD + HDD | 数据库、对象存储、文件服务 |
中间件(PostgreSQL、Redis、RabbitMQ、配置中心、调度中心)统一以容器方式部署,用 ConfigMap 管配置:
apiVersion: apps/v1
kind: Deployment
metadata:
name:>apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
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 + 中间件撑住了弹性和稳定,机器坏一台不影响整体;
- 元数据和监控让数据“看得见、追得到、管得住”,这才是平台和一堆脚本的本质区别。
如果你们也在做类似的数据平台,我的建议是:先把元数据模型想清楚,再动手写采集。元数据是整个平台的骨架,骨架歪了,后面补得越勤越痛苦。
代码片段都是脱敏后的示例,思路可以直接借鉴。有踩过类似坑的朋友,欢迎评论区交流。