ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

两台服务器搭建轻量级数据平台:批处理、流处理与AI推理实践

两台服务器搭建轻量级数据平台:批处理、流处理与AI推理实践 手头有两台机架式服务器机箱标签上分别印着两组序列号ANV32AA1WDK66和R7KA8T2LFLCAC。团队日常的数据任务杂乱到让人头疼——早上的销售报表 ETL、白天的 Nginx 日志实时清洗、晚上的模型批量推理还有各种临时性的 CSV 转 Parquet、数据库迁表、爬虫结果去重。以前这些活儿全挤在一台老服务器上每天早上 CPU 都能冲到 100%谁跑个重任务其他人就得等。后来我把这两台带有明确序列号标识的节点重新规划了一遍用它们搭了一个轻量级的通用数据处理平台。折腾了大概三天把批处理、流处理、AI 推理和日常文件搬运全部跑通了。这篇文章就把完整的思路、命令、配置和踩坑过程记录下来给同样手上有几台闲置机器、想自己搭一套“啥活儿都能接”的数据处理环境的朋友做个参考。1. 为什么先别急着部署把“序列号”当数据任务的指挥地图1.1 先确认机器身份ANV32AA1WDK66 与 R7KA8T2LFLCAC 到底是什么很多人在一堆机器上部署服务时有个坏习惯直接用 IP 或者“那台旧机器”“新买的大机器”来指代节点。等机器多了驱动装错、服务起错位置、日志刷错主机都是常事。我拿到这两台机器后做的第一件事不是装环境而是采集硬件资产信息。在每台机器上执行dmidecode -s system-serial-number lscpu | grep Model name free -h lspci | grep -E VGA|NVIDIA|RAID df -h第一台ANV32AA1WDK66的输出确认了它的身份8 核 CPU、64GB 内存、一块 NVIDIA 加速卡用于模型推理、一块 2TB NVMe 硬盘。这台机器就是典型的“计算偏科生”——CPU 算力不错但核心数不算多有大内存和 GPU适合跑批处理聚合和 AI 推理。第二台R7KA8T2LFLCAC是 16 核 CPU、128GB 内存、四块 NVMe SSD做了 RAID0和一张万兆网卡。这台机器一看就是“I/O 偏科生”CPU 多、磁盘吞吐猛最适合跑 Kafka、Flink 这类流处理任务以及高频日志写入。我把这两台机的信息登记到了资产管理表里并且约定以后所有运维指令里只认序列号不认 IP。这样做的原因后面会讲。1.2 序列号如何决定数据任务分配序列号本身不产生算力但它是我做任务调度的“指挥地图”。因为两台机器的硬件配置差异非常大我直接按序列号把任务做了规划ANV32AA1WDK66负责离线批量任务、ETL、AI 推理任务。这类任务通常需要大量内存和计算能力但不像流任务那样需要持续高并发。R7KA8T2LFLCAC负责实时流处理、日志采集、消息队列中转。这类任务对磁盘读写和网络吞吐要求极高但对单核性能要求相对低。后面的监控告警、日志采集器也都会以序列号作为主机标识上报。这样一旦出了问题我不用猜是哪台机器直接看序列号。1.3 网络拓扑与基础初始化两台机器都装的是 Ubuntu 22.04 LTS。内网静态 IP 分配如下节点序列号内网 IP职责节点AANV32AA1WDK66192.168.10.10批处理 / AI 推理节点BR7KA8T2LFLCAC192.168.10.11流处理 / 消息队列基础初始化只需要几步但每一步都值得做创建统一用户dataops用于跑所有数据任务配置 SSH 免密登录方便后面写调度脚本关闭本机防火墙因为在可信内网中自己管理端口更灵活开启 NTP 时间同步流处理和日志任务对时间一致性要求很高。sudo apt update sudo apt upgrade -y sudo useradd -m -s /bin/bash dataops sudo usermod -aG sudo dataops ssh-keygen -t ed25519 ssh-copy-id dataops192.168.10.10 ssh-copy-id dataops192.168.10.11至此两台机器的身份问题解决。“谁的序列号是谁”“哪台机器负责哪类任务”都清晰了接下来的部署才能有条不紊。2. 两台机器四类活儿批处理、流处理、AI推理、文件搬运怎么分工2.1 数据任务的四种典型形态“无论任务如何”这句话听起来很夸张但实际团队里的数据任务确实跑不出四种类型任务类型特点典型例子资源瓶颈批处理 ETL数据量大、定时跑、允许延迟销售报表聚合、离线特征计算CPU、内存流处理数据持续到达、要求低延迟日志清洗、实时指标统计磁盘、网络、堆内存AI 推理模型加载后批量预测图片分类、文本向量化GPU、CPU文件搬运与格式转换简单但琐碎、IO密集CSV转Parquet、跨机同步磁盘、网络把这四类任务和两台机器的序列号一对应分工自然就出来了。但这里有个微妙的问题谁也不能保证某类任务永远只在一台机器上跑。所以我后面的调度设计并不是把任务写死在某个 IP 上而是机器根据自己的“序列号标识”选择角色任务通过统一入口分派。2.2 架构选型不上 K8s 时的轻量方案小团队两台机器没必要上一套 Kubernetes——控制面本身就要吃资源而且运维复杂度会立刻把效率吃掉。我选择的是Docker Compose systemd 定时器的组合方案部署快、问题少、好排查。在ANV32AA1WDK66上部署了MinicondaPython 数据环境Onnx Runtime 推理服务定时批处理容器在R7KA8T2LFLCAC上部署了Kafka单节点Flink单机模式日志清洗作业容器每台机器上的容器都加了com.xxx.node.serial的标签方便按序列号过滤和区分服务docker run --name etl-worker \ -l node.serialANV32AA1WDK66 \ -v /data/etl:/data \ -d myrepo/etl-worker:latest2.3 容器化还是裸机的取舍有些朋友可能会问数据任务直接裸机跑不好吗为什么非要套一层容器我的真实感受是容器化的最大价值不是隔离而是让环境可以随时重建。流处理作业的依赖版本经常变动直接在裸机上装 Python 包翻车了要满系统找依赖残骸。用容器把作业环境打包好改一个参数重新docker compose up -d就行省心。代价是容器会带来 5% 到 10% 的网络和磁盘性能损耗但对于这两台机器的配置来说完全不是问题。真正性能敏感的磁盘写入路径我会在流处理那节专门讲怎么绕开 Docker 的默认限制。3. 在 ANV32AA1WDK66 上跑通批任务与模型推理的完整记录3.1 部署 Python 数据环境在这台机器上我选择了 Miniconda 来管理数据科学相关依赖因为numpy、pandas、pyarrow这些包的二进制依赖太多系统自带的 Python 环境装起来很麻烦。wget https://repo.anaconda.com/miniconda/Miniconda3-latest-Linux-x86_64.sh bash Miniconda3-latest-Linux-x86_64.sh -b -p /opt/miniconda3 conda create -n data python3.12 -y conda activate data pip install pandas pyarrow duckdb polars onnxruntime opencv-python-headless这里有一个小教训不要直接用pip install onnxruntime-gpu刚开始我以为 GPU 版会自动匹配 CUDA结果装完一运行发现它把 CPU 和 GPU 的库全混在一起。正确的做法是先nvidia-smi确认 CUDA 版本再装对应的onnxruntime-gpu1.17.0这类明确版本。3.2 批处理实战用 DuckDB 跑销售订单聚合我选了一个 2GB 的模拟订单数据来做测试表结构是常见的订单明细订单号、用户ID、商品ID、金额、订单时间。要算每个商品每天的销售额。如果写 Pandas 硬扛内存会飞到 30GB 以上速度还慢。我改用 DuckDB直接 SQL 聚合速度快到有点意外CREATE TABLE orders AS SELECT * FROM read_csv_auto(/data/orders.csv, headertrue); SELECT product_id, strftime(order_time, %Y-%m-%d) AS day, SUM(amount) FROM orders GROUP BY 1, 2;2GB 数据、约 1200 万行DuckDB 全部跑完只需要 18 秒内存峰值不到 6GB。这个速度对日常 ETL 完全够用。不过我建议如果你的任务里还有复杂的多表 JOIN并且业务量会持续涨建议还是用 Spark 的 Standalone 模式来跑。DuckDB 适合单机数据规模可控的场景胜在部署零成本。3.3 AI推理实战ONNX Runtime 批量图片分类手头有一个图像分类模型导出成了 ONNX 格式任务是给一张目录下大约 5 万张图片做推理。直接用 Python 循环跑 OpenCV 解码 ONNX Runtime 推理我实测下来速度大约是 220 张/秒。为了保证吞吐再往上提我做了两件优化把OMP_NUM_THREADS和torch的线程数全部关掉避免 PyTorch 与 ONNX Runtime 抢占 CPU 线程用ThreadPoolExecutor做了 4 线程并发但每个线程单独创建一个推理 session。import cv2, onnxruntime as ort, os from concurrent.futures import ThreadPoolExecutor sess ort.InferenceSession(/models/model.onnx, providers[CUDAExecutionProvider]) def infer(path): img cv2.imread(path) img cv2.resize(img, (224, 224)) inp img.transpose(2, 0, 1)[None].astype(float32) out sess.run(None, {input: inp}) return path, out[0].argmax() with ThreadPoolExecutor(max_workers4) as ex: for p in os.listdir(/data/images): ex.submit(infer, os.path.join(/data/images, p))优化后吞吐提升到 340 张/秒单张图片平均耗时控制在 3ms 左右。这套配置在这台机器上跑得很稳GPU 利用率大约 70%。3.4 指标记录与任务持久化跑完任务后我把耗时、内存峰值、结果行数都写进了一个本地 SQLite 库方便后续回溯。三天的运行数据显示ANV32AA1WDK66上同时跑两个批任务CPU 平均 70%内存 40GB 左右GPU 利用率 60%-80%没有任何任务互相拖垮的情况。4. R7KA8T2LFLCAC 变身流处理节点的折腾过程4.1 为什么选 Kafka Flink流处理的最出名方案无非 Kafka Flink、Redis Streams、Pulsar。我选择 Kafka Flink 的原因有两条第一团队对“回放”能力有硬需求。日志数据源偶尔会抖动Kafka 可以把消息保留三天等数据源恢复了再重新消费不会丢数据。Redis Streams 虽然简单但内存成本高且消费组管理远不如 Kafka 成熟。第二Flink 的窗口计算能力是刚需。我们需要在滚动窗口内统计接口错误率Flink 的TumblingEventTimeWindow写起来非常顺手。单机版本的 Kafka 其实完全够用只要把分区数拉开。我给了 Kafka 8GB 堆内存Flink TaskManager 同样分配 8GB。4.2 部署 Kafka 单节点 Flink 单机模式Kafka 直接下载官方二进制包解压后修改config/server.propertiesbroker.id1 listenersPLAINTEXT://0.0.0.0:9092 log.dirs/data/kafka-logs num.partitions12 offsets.topic.replication.factor1 log.retention.hours72 zookeeper.connect192.168.10.11:2181num.partitions这里我直接设到了 12。之前用过 3后面你会看到这个教训。Flink 选择 Standalone 模式cd /opt/flink ./bin/start-cluster.sh默认的 Flink 配置只有 1GB 内存必须改conf/flink-conf.yamljobmanager.memory.process.size: 2048m taskmanager.memory.process.size: 8192m taskmanager.numberOfTaskSlots: 4 parallelism.default: 4如果不改这两个内存配置作业跑起来会疯狂 GC。这是流任务最容易忽视的坑。4.3 日志清洗任务从 Nginx Access Log 到结构化指标清明任务的核心逻辑很简单消费 Kafka 里的原始日志行解析 JSON或者正则提取按接口、状态码做 1 分钟窗口聚合最后写入 ClickHouse。我直接用 Flink 的 DataStream API 写了个 60 行的 Java 作业核心逻辑大概是这样DataStreamString raw env.addSource(new FlinkKafkaConsumer( raw_logs, new SimpleStringSchema(), props)); raw.map(line - parseToPojo(line)) .keyBy(e - e.getApiPath()) .window(TumblingProcessingTimeWindows.of(Time.minutes(1))) .aggregate(new CountAggregate()) .addSink(new ClickHouseSink());实测单机吞吐从 Kafka 消费到写 ClickHouse平均每秒处理 4.8 万条日志。瓶颈出在parseToPojo因为大量高并发下 Gson 对象创建太频繁。我把 Gson 换成了手写的字符串 split Long.parseLong吞吐提升到 6.5 万条/秒。这里有个完全没想到的问题Flink 作业的并发度设为 4但 Kafka 分区只有 12 个按理说是够的。可是当并行度超过 4 后Flink 作业抛异常“No resource available”。原因是我改了taskmanager.numberOfTaskSlots: 4Flink 只有 4 个 slot并行度不能超过这个数字。把并行度调整为 4 后恢复正常。4.4 压测结果和瓶颈分析用kafka-producer-perf-test模拟高吞吐写入测试bin/kafka-producer-perf-test.sh \ --topic raw_logs --num-records 1000000 \ --record-size 800 --throughput 50000 \ --producer-props bootstrap.servers192.168.10.11:9092Kafka 单节点接收 50MB/s 的写入完全无压力CPU 占用率只有 35%。真正的瓶颈在磁盘因为容器默认存储驱动是 overlay2写容器层会有额外开销数据落到/data/kafka-logs时如果路径挂载的是普通目录IO 会有明显抖动。我的解决办法是把 Kafka 的log.dirs直接指向宿主机 RAID0 的裸盘路径不经过 Docker 数据卷映射。毕竟流数据处理的核心就是磁盘 IO把这个绕开比什么都重要。5. 双机协作混合任务调度和监控告警的实用配置5.1 任务怎么分节点角色与标签两台机器各管一摊其实还不够。因为有些任务会串起来比如ANV32AA1WDK66上跑完 ETL 后需要把结果同步给R7KA8T2LFLCAC去供实时查询。这种“混合任务流”我的做法是用序列号作为“节点角色标签”来管理。每台机器上的/etc/hosts都做了解析192.168.10.10 node-a-anv32aa1wdk66 192.168.10.11 node-b-r7ka8t2lf lcac 192.168.10.11 node-b-r7ka8t2lflcac另外我维护了一份简单的任务清单表每次新任务进来先看它属于哪一类再分配节点任务ID任务类型分配节点序列号调度方式ETL_SECOND批处理ANV32AA1WDK66systemd timer 每日 06:00INFER_IMAGEAI 推理ANV32AA1WDK66systemd timer 每晚 22:00LOG_CLEAN流处理R7KA8T2LFLCACFlink 常驻作业SYNC_DATA文件搬运R7KA8T2LFLCACpython 脚本 每小时5.2 调度脚本systemd timer Ansible复杂调度全用 K8s 或 Airflow 是杀鸡用牛刀我用了一个非常朴素的方案每个定时任务写成一个 shell 脚本由 systemd timer 触发。以ANV32AA1WDK66上的 ETL 为例# /etc/systemd/system/etl-daily.service [Unit] DescriptionDaily ETL with DuckDB [Service] Userdataops ExecStart/opt/scripts/etl_daily.sh# /etc/systemd/system/etl-daily.timer [Unit] DescriptionSchedule ETL every day at 6am [Timer] OnCalendar*-*-* 06:00:00 Persistenttrue [Install] WantedBytimers.target这种方式的优点是每台机器上的任务状态通过systemctl list-timers一目了然出问题直接查 unit 的日志不用额外搭调度平台。混合任务之间偶尔有依赖关系我在脚本里加了一步很简单的检查先curl一下文件搬运节点上的健康检查接口返回 200 才继续执行 ETL。两台机器之间有 SSH 免密跨机执行命令也就一行的功夫ssh dataops192.168.10.11 /opt/scripts/check_ready.sh5.3 监控Prometheus node_exporter alertmanager轻量监控这套组合是我的首选三台进程就能覆盖全部指标。在每台机器上装node_exporter然后在ANV32AA1WDK66上跑一个单机版 Prometheus抓取两个节点的数据scrape_configs: - job_name: node static_configs: - targets: [192.168.10.10:9100, 192.168.10.11:9100]告警规则我配了三个先在用的CPU 持续 5 分钟超过 90%磁盘剩余空间低于 20%节点 ping 不通。实际上最有用的告警是磁盘空间。流处理日志非常吃硬盘我的目录/data/kafka-logs三天就涨了 800GB要是不盯着磁盘满了整个 Kafka 会直接崩溃。5.4 避免任务互相干扰的配置我一个很深刻的体验是即使两台机器分工明确但一台机器上批任务和推理任务同时跑可能会因为内存争抢而互相拖垮。解决办法很简单在跑重任务前用cgset或者 Docker 的--memory限制住单个容器的内存上限。我给批处理容器设了 24GB 内存上限给推理容器留了 20GB这样即使一个任务内存泄漏也不会导致另一台上的服务 OOM。这种“宁可错杀不可放过”的隔离策略用在数据任务上很管用。6. 这三天里踩过的坑比官方文档值钱6.1 序列号引起的驱动绑定问题别觉得奇怪真实遇到过。我最初在ANV32AA1WDK66上装 NVIDIA 驱动时驱动一直报错“no devices found”。查了一圈发现系统里确实有 GPU但设备 ID 没有被驱动识别。定位过程非常依赖序列号我在lspci -v里记录了 GPU 的 BDF 地址然后查lshw -class display -json确认这块 GPU 的物理插槽位置。最后发现是因为之前机器被拿去做过别的测试设备被切换到了PCIe Gen3模式驱动侧需要的gen4配置没对上重设 PCIe link 速度后解决。这就是为什么要始终带着“序列号”意识做硬件资产管理。不然同类问题换个板卡插槽你根本不知道它是哪台机器、哪块卡出了问题。6.2 透明大页导致 Flink GC 问题流处理节点跑了一段时间后Flink 的 GC 日志频繁出现 Full GC而且每次 Full GC 要停顿 4 秒以上。查了一圈问题出在 Linux 的透明大页Transparent Huge Pages。Kafka 和 Flink 在堆内存分配时使用的是mmap和直接内存而透明大页机制会把 4KB 的小页合并成 2MB 的大页看似提高了内存管理效率实际在频繁分配释放时会造成内存碎片并触发长时间停顿。解决方法是永久关闭透明大页echo never /sys/kernel/mm/transparent_hugepage/enabled同时在/etc/default/grub里加上transparent_hugepagenever让重启后也生效。这个坑在流处理节点上特别明显批处理节点倒是没怎么受影响。6.3 Docker 容器内打开文件数限制跑日志清洗容器时我发现 Flink 作业稳定运行一天后会突然报“Too many open files”。原因是容器默认的ulimit -n是 1024而 Flink 在消费 Kafka 时每个并行子任务会建立大量 TCP 连接和文件句柄这个默认值完全不够。修改容器启动参数即可docker run --ulimit nofile65536:65536 ...如果你用 docker-compose在服务定义里加ulimits: nofile: soft: 65536 hard: 65536这个坑非常隐蔽因为平时看不到报错只有流量峰值时才会暴露。建议所有跑数据处理任务的容器一律把nofile调高。6.4 Kafka 分区数设太小的教训一开始我没想太多把 Kafka 的num.partitions设成了 3结果日志消费经常出现积压明明 Flink 消费端有 4 个并行度却只有 3 个分区可用白白浪费一个并发。这里有个基本常识Kafka 的消费并发度上限就是分区数而且分区数只能在创建 topic 时确定之后只能通过kafka-topics.sh --alter增加但不能减少。我当时只能新建了一个 12 分区的新 topic再把老的积压数据重新导入。所以宁可多分也不要少分数据增长是必然的。6.5 断电重启后的服务自启问题这个不算技术坑但特别影响使用体验。前两天的折腾没做开机自启结果有一次机柜断电重启后Kafka、Flink、Prometheus 一台都没起来而我第二天到公司才发现数据任务停了整整 8 小时。补了一轮 systemd unit 文件把关键服务设置成Restartalways[Service] Restartalways RestartSec10这个改动看似不起眼但在无人值守的夜里太重要了。我现在每次部署一套新服务第一件事就是写 unit 文件而不是先nohup跑起来。6.6 我个人的一点体会这三天的折腾最大的收获不是学会用了多少新工具而是确认了一个结论对于这种服务器数量不多、任务种类繁杂的场景把序列号作为资产管理和任务调度的核心标识配合轻量容器化方案比一味追求上 K8s 或者盲目堆云服务更可靠也更省心。两台机器到现在已经稳定跑了一个多月每天的 ETL、日志流清洗、图片推理任务几乎没有中断过。流处理节点在高峰期也能保持 5 万条/秒的摄入速度批处理和推理任务各自吃掉了超过 60% 的算力但彼此都没有互相拖垮。如果你的团队也是这种规模数据任务也很杂手头正好有几台带序列号的空闲服务器完全可以照这个思路搭一套。后续我准备把 ClickHouse 加进去做实时指标查询再把任务调度从 systemd 升级成轻量的 Airflow。这个折腾的方向应该还能持续很久。
返回列表