ARTICLE DETAIL

资讯详情

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

数据处理和存储系统建设方案:从Excel到千万级吞吐的落地路径

数据处理和存储系统建设方案:从Excel到千万级吞吐的落地路径 简介这份《10数据处理和存储系统建设方案》面向智能化管控平台的项目规划人员、系统集成工程师与信息化方案撰写者围绕化工园区等场景解决数据处理能力、存储容量与传输带宽如何量化测算的问题。资源为1个doc文档压缩包约183KB内容以方案正文与测算表格为主便于直接参考或改写为项目立项材料。文档从用户构成入手将用户划分为化工企业、政府及园区职能部门、互联网公众三类按300个用户、3至5年需求展开分析数据计算部分以TPC-C基准推算数据库服务器峰值处理能力给出并发数、事务数与冗余系数等参数并据此估算CPU总核数与虚拟化服务器数量存储部分区分系统数据、企业报送数据与非结构化数据测算出三年约8.1TB容量配置传输部分覆盖平台用户、物联网感知与视频监控带宽需求并附软硬件选型配置表。已有223人学习适合需要完整测算逻辑与设备选型依据的读者参考。1. 数据处理和存储系统建设方案从一堆 Excel 到每天千万级吞吐的落地路径很多团队第一次认真考虑“数据处理和存储系统建设方案”不是因为技术驱动而是因为某个周一早上运营在群里问了一句“上周的日活到底是多少”结果三个人给出了三个数。Excel 散落在不同人手里日志文件堆在服务器上数据库里只有业务表没有明细层。这时候你才意识到数据处理和存储系统不是“要不要建”的问题而是“再不建就要被数据淹死”的问题。这个方案要解决的核心就三件事数据怎么进来、怎么存、怎么算。适合谁看适合正在从脚本Excel 向系统化过渡的中小团队也适合需要把高通量数据处理、流式数据处理和离线数仓打通的一线工程师。下面我按实际落地顺序把选型、步骤、参数和踩坑点拆开讲。2. 先定分层再选型数据处理和存储系统的四层架构怎么切2.1 为什么不能一上来就选 Kafka Hudi我见过不少团队一提到数据处理和存储系统建设方案第一反应就是“上实时湖仓”。结果三个月后Kafka 里堆了 200 个 topicHudi 表小文件泛滥查询慢到业务方直接放弃。问题不在技术选型而在分层没想清楚。常见做法是先把数据流切成四层采集层、存储层、计算层、服务层。采集层负责把日志、数据库变更、API 回调统一收进来存储层按“原始数据保真、明细数据可查、汇总数据加速”分三档计算层区分流式和批式服务层只暴露指标和标签不让业务直接查明细。这个分层的好处是你可以在每一层独立选型。采集层用 Filebeat Kafka 还是 Flume Pulsar取决于你的运维习惯存储层用 HDFS Hive 还是 S3 Iceberg取决于你的云环境计算层用 Flink 还是 Spark取决于你的延迟要求。不要试图用一个组件解决所有问题。提示分层不是画架构图用的是让你在出故障时能快速定位——是采集丢了还是存储没落盘还是计算逻辑错了。2.2 存储选型行存、列存、对象存储各管什么存储层最容易翻车的地方是把所有数据都塞进一个数据库。MySQL 扛不住每天千万级明细写入Elasticsearch 存全量明细成本高得离谱HBase 又不太适合做聚合分析。我的经验是按访问模式分三档存储类型典型组件适合数据保留周期查询特征行存MySQL / PostgreSQL维度表、配置表、任务元数据长期点查、小范围更新列存ClickHouse / Doris明细宽表、指标汇总3-12 个月聚合、分组、排序对象存储S3 / OSS / HDFS原始日志、备份、冷数据1-3 年批量扫描、归档参数上列存表的分区键一般选日期排序键选查询最频繁的维度。比如 ClickHouse 的PARTITION BY toYYYYMM(create_time) ORDER BY (tenant_id, create_time)这样按租户和时间范围查都能命中分区裁剪。2.3 计算层流式和批式怎么分工流式处理适合“事件发生到指标可见”在秒级到分钟级的场景比如实时风控、监控告警。批式处理适合“T1 出报表”的场景比如日活、留存、GMV 汇总。两者不是替代关系是互补关系。我一般会这样分工Flink 消费 Kafka 做实时清洗和窗口聚合结果写入 ClickHouse 或 Doris 供实时看板查询同时 Kafka 的数据通过 Connector 落到对象存储Spark 每天凌晨跑一遍全量批处理修正流式计算中因为乱序或迟到数据导致的不一致。这样既保证了实时性又保证了最终一致性。注意流式和批式的指标口径必须用同一套 SQL 或同一份配置否则业务方会发现“实时看板和日报对不上”然后你就得花一周时间对数。3. 采集和清洗把数据从源头搬进存储层的具体步骤3.1 数据库变更采集CDC 配置与字段映射数据库变更采集最常见的是 MySQL binlog 同步。用 Flink CDC 或 Debezium 都可以核心是配置好 server-id、binlog 格式和位点保存。# Flink CDC MySQL source 配置示例 source: type: mysql-cdc hostname: 10.0.0.12 port: 3306 username: cdc_reader password: ${CDC_PASSWORD} database-name: order_db table-name: order_db\.order_info server-id: 5401-5404 scan.startup.mode: initial scan.incremental.snapshot.enabled: true debezium.snapshot.mode: schema_only_recovery这段配置的关键参数server-id要避开 MySQL 主从集群已用的范围否则会冲突scan.startup.mode选initial表示先全量快照再增量scan.incremental.snapshot.enabled开启后全量阶段不锁表对线上影响小。debezium.snapshot.mode设为schema_only_recovery是为了在任务重启时从位点恢复而不是重新做全量。字段映射上我习惯在采集层就把字段名统一成小写下划线格式时间字段统一转成yyyy-MM-dd HH:mm:ss字符串或毫秒时间戳。这样下游计算层不用再为每个源做适配。3.2 日志清洗用 Python 做结构化解析和脏数据分流日志类数据往往半结构化常见做法是用 Python 脚本做第一轮解析。下面是一个处理 Nginx 日志的示例把原始行解析成 JSON并把解析失败的行单独写到错误目录。import re import json import os from datetime import datetime LOG_PATTERN re.compile( r(?Pip\S) - - \[(?Ptime[^\]])\] r(?Pmethod\S) (?Ppath\S) \S r(?Pstatus\d) (?Pbytes\d) ) def parse_line(line): match LOG_PATTERN.match(line) if not match: return None d match.groupdict() d[status] int(d[status]) d[bytes] int(d[bytes]) d[time] datetime.strptime(d[time], %d/%b/%Y:%H:%M:%S %z).isoformat() return d def process_file(src_path, ok_path, err_path): with open(src_path, r, encodingutf-8, errorsignore) as f: for line in f: parsed parse_line(line.strip()) if parsed: with open(ok_path, a, encodingutf-8) as ok: ok.write(json.dumps(parsed, ensure_asciiFalse) \n) else: with open(err_path, a, encodingutf-8) as err: err.write(line) if __name__ __main__: process_file(/data/logs/nginx.log, /data/clean/ok.jsonl, /data/clean/err.log)逻辑说明正则提取 IP、时间、方法、路径、状态码、字节数时间转成 ISO 格式方便后续入库解析失败的行不丢弃写到错误文件供人工排查。参数上errorsignore防止编码问题导致整个任务中断ensure_asciiFalse保证中文正常写入。提示清洗脚本不要直接写数据库先落成 JSONL 文件再用批量导入工具入库。这样脚本可以重跑不会因为数据库连接抖动丢数据。3.3 流式接入Kafka 分区和副本参数怎么设Kafka 作为采集层的缓冲分区数决定了并行度副本数决定了可靠性。我的经验是分区数 峰值吞吐 / 单分区处理能力一般单分区 5-10 MB/s副本数至少 2生产环境建议 3。# 创建 topic 示例 kafka-topics.sh --create \ --bootstrap-server kafka-broker:9092 \ --topic order_events \ --partitions 12 \ --replication-factor 3 \ --config retention.ms604800000 \ --config cleanup.policydeleteretention.ms设为 7 天给下游消费留足重放窗口cleanup.policydelete表示按时间删除不用 compact。如果 topic 是变更日志类可以改成compact保留每个 key 的最新值。分区数不是越多越好。分区太多会导致小文件问题Flink 或 Spark 消费时也会增加调度开销。一般从 6 或 12 起步观察消费延迟再调整。4. 存储层落地建表、分区、索引和生命周期管理4.1 列存表设计ClickHouse 建表模板与参数解释ClickHouse 是明细和汇总层常用的列存引擎。下面是一个订单明细表的建表语句包含分区、排序键和 TTL。CREATE TABLE dwd_order_detail ( order_id String, user_id UInt64, product_id UInt32, amount Decimal(18, 2), status String, create_time DateTime, dt Date DEFAULT toDate(create_time) ) ENGINE MergeTree() PARTITION BY toYYYYMM(dt) ORDER BY (dt, user_id, order_id) TTL dt INTERVAL 12 MONTH SETTINGS index_granularity 8192;参数说明PARTITION BY toYYYYMM(dt)按月分区避免单分区过大ORDER BY (dt, user_id, order_id)让按日期和用户查询能命中主键索引TTL dt INTERVAL 12 MONTH自动删除一年前数据省去手动清理index_granularity 8192是默认值数据量特别大时可以调到 16384 减少索引体积。注意ClickHouse 不适合频繁更新单行如果业务需要更新订单状态建议用 ReplacingMergeTree 或把更新写到另一张表做关联。4.2 对象存储分区规范按日期和来源分目录对象存储上的原始数据目录结构直接影响后续批处理效率。我一般用来源/日期/小时三级分区s3://data-lake/raw/order_events/dt2025-01-15/hour08/ s3://data-lake/raw/nginx_logs/dt2025-01-15/hour08/ s3://data-lake/clean/order_detail/dt2025-01-15/这样 Spark 或 Hive 读取时可以直接做分区裁剪不用全量扫描。文件格式优先选 Parquet 或 ORC压缩比高列裁剪效果好。单个文件大小控制在 128 MB 到 512 MB 之间太小会拖慢任务太大不利于并行。4.3 生命周期策略热、温、冷数据怎么迁移不是所有数据都值得一直放在列存里。我的策略是最近 3 个月的数据放 ClickHouse 或 Doris 供实时查询3 到 12 个月的数据放对象存储用 Hive 或 Spark 按需查超过 12 个月的数据转成归档格式只保留汇总指标。-- Doris 中设置冷热分层示例 ALTER TABLE dwd_order_detail SET ( storage_policy hot_to_cold, cold_boundary 2024-10-01 );这个配置表示 2024-10-01 之前的数据自动迁到冷存储。不同组件语法不同核心思路是让热数据占用的高性能存储尽量小把成本花在真正需要快速查询的数据上。5. 避坑与排查数据处理和存储系统建设中最容易翻车的 5 个点5.1 现象实时看板和日报数字对不上原因流式计算和批式计算用了不同的时间窗口和过滤条件。流式按事件时间算批式按入库时间算迟到数据在两边的归属不同。解决统一用事件时间做窗口流式任务设置 allowedLateness批式任务读取时按事件时间过滤。两边共用同一份指标定义 SQL不要各写各的。5.2 现象Kafka 消费延迟越来越高但 CPU 和内存都不高原因分区数不够或者消费者组里有个别消费者卡住导致 rebalance。也可能是下游写入太慢反压到 Kafka 消费。解决先看消费者 lag如果所有分区都积压加分区或加消费者如果个别分区积压检查该分区对应的消费者实例是否异常。下游写入慢的话在 Flink 里加批量写入和异步 IO。5.3 现象ClickHouse 查询越来越慢明明数据量没涨多少原因小文件太多或者分区键选得不好导致扫描范围过大。也可能是索引粒度太细主键索引膨胀。解决用OPTIMIZE TABLE ... FINAL合并小文件但不要频繁执行。检查查询是否命中分区裁剪EXPLAIN看扫描行数。分区键一般选日期不要选高基数字段。5.4 现象CDC 任务重启后数据重复或丢失原因位点保存不可靠或者全量阶段和增量阶段切换时没有做幂等。MySQL binlog 过期也会导致位点失效。解决开启 Flink Checkpoint把位点保存到可靠存储下游写入用主键去重比如 ClickHouse 的 ReplacingMergeTree 或 Doris 的 Unique 模型。binlog 保留时间至少 7 天给任务恢复留窗口。5.5 现象对象存储上文件数量爆炸NameNode 压力大原因流式任务每分钟甚至每秒写一个小文件没有做文件合并。解决在 Flink 里设置滚动策略按大小或时间滚动比如每 128 MB 或每 5 分钟写一个文件。Spark 批处理任务结束后做一次小文件合并。HDFS 上可以开 Archive对象存储上可以用生命周期规则转成低频访问。6. 进阶技巧用数据质量校验和血缘追踪给系统上保险系统跑起来之后最怕的不是性能问题而是数据悄悄错了没人发现。我一般会在关键链路加两道保险数据质量校验和血缘追踪。数据质量校验不用一上来就搞 Great Expectations 这种重框架先用 SQL 做核心指标的空值率、唯一性、值域检查。比如每天凌晨批处理结束后跑一遍校验 SQL-- 校验订单明细主键唯一性和金额非负 SELECT COUNT(*) AS total, COUNT(DISTINCT order_id) AS distinct_orders, SUM(CASE WHEN amount 0 THEN 1 ELSE 0 END) AS negative_amount FROM dwd_order_detail WHERE dt 2025-01-15;如果total和distinct_orders差距超过阈值或者negative_amount大于 0就触发告警。这个校验脚本可以放在调度系统里和批处理任务串起来。血缘追踪方面小团队不用急着上 Atlas 或 DataHub先在调度任务里记录输入表和输出表的映射关系存到一张元数据表里。出问题时能快速知道“这个指标是从哪张表算出来的”比翻代码快得多。我自己的习惯是每建一个新表就在元数据表里登记来源、更新频率、负责人和校验规则。这个习惯帮我省过很多次“这个字段谁改的”的扯皮时间。数据处理和存储系统建设方案不是一次性的架构设计而是持续迭代的工程习惯。希望帮到你。本文还有配套的精品资源点击获取
返回列表