资讯动态

自研Java大数据清洗工具DataJet:架构设计与亿级数据性能优化实战

发布时间:2026/9/13 3:12:36 来源:尧图企业网站定制
去年六月大促结束第二天业务方凌晨三点把我从床上叫醒。BI大屏上全国门店的销售实时曲线在零点之后直接断崖归零而数据库里当天的订单表一切正常。查了半天才发现数据管道里多了一步清洗脚本——Python脚本在处理某家门店新接入的设备上报数据时抛了异常整条链路直接停掉了。这种问题已经不是第一次出现每次都要有人去翻Shell日志、改脚本、重新跑数。那天凌晨我坐在工位上想明白了一件事在Java技术栈这么成熟的大数据领域我们不能再靠零散的临时脚本去维护数据质量必须自己设计并实现一个真正意义上的大数据清洗工具。这篇文章就把这个工具从设计到落地的完整过程拆开讲一讲包括整体架构、核心模块实现、亿级数据下的性能优化以及上线后踩过的几个大坑。如果你准备做Java大数据方向的项目或者正在为数据清洗选型发愁这里面的经验和代码思路可以直接参考。1. 凌晨三点的脏数据事故逼我们自研清洗工具1.1 事故现场Python脚本为什么撑不住了先说那次事故的导火索。门店上报的销售数据里某个新接入的智能POS机厂商返回的时间字段格式变成了2024/06/15 14:32:08而之前的脚本里写死了datetime.strptime(row[time], %Y-%m-%d %H:%M:%S)异常没有被捕获整个任务退出。后面的数据同步、指标计算全部跟着断掉。这其实暴露出了脚本模式处理数据清洗的几大痛点容错能力为零。一条脏数据就能让全链路崩溃没有跳过、没有隔离、没有告警。规则改不动。业务规则变了要改代码改完还要手动重跑上线流程非常随意。性能顶不住。单机Python脚本处理百万级数据还能忍到亿级就完全跑不动内存直接爆掉。没人敢维护。写了一年的清洗脚本散落在不同的服务器上有的还是clean_v7_final.py这种名字。而我们的数据量增长非常快日均订单和门店上报数据加起来接近1.2亿行大约35GB。靠脚本已经不是优化的问题而是从根本上就不适合这个量级。1.2 市面工具的短板以及DataJet的定位决定自研之前我们也不是没调研过现成的清洗工具。粗略对比下来各有各的问题工具优点在我们场景下的问题DataX数据同步能力强插件丰富主要做迁移和同步复杂清洗逻辑支持弱Kettle界面化配置功能全面太重大任务跑起来性能一般团队不熟Spark DataFrame分布式计算能力强还是要写代码业务团队没法直接配置规则自研脚本灵活可维护性差性能瓶颈明显市面上的工具要么太“重”要么太“底层”缺一个能让我们Java团队快速定制规则、又能扛住日亿级数据的清洗中间件。所以我们决定自己动手项目代号定为DataJet。DataJet的核心定位很明确一个基于Java的离线数据清洗引擎数据通过接入层变成统一记录流经过可编排的规则算子链完成过滤、转换、去重、标准化最终输出到目标存储。所有规则支持JSON配置化新增规则通过SPI接口扩展不需要改动主流程代码。2. DataJet整体架构一条数据流经过哪些关卡2.1 任务编排层Job-Step-Pipeline三层模型第一批版本设计的时候我们刻意先把“模型”定清楚再写代码。DataJet采用典型的Job-Step-Pipeline三层模型Job一个完整的清洗任务包含数据源配置、清洗规则集、目标端配置、调度策略。StepJob内部的具体执行单元比如“读取门店订单表”“执行去重规则”“写入目标库”。PipelineStep按依赖关系组成的执行链支持串行和并行。用这个模型的原因很简单——实际业务里的清洗任务从来不是一个线性的“读取-清洗-写入”而是有多张表关联、有前置依赖、有失败重试的复杂流程。比如我们有一个任务需要先清洗门店维表再清洗订单事实表维表清洗完成后才能做订单表里门店编码的映射。如果只做一个单纯的“循环处理”后面扩展起来就是灾难。Pipeline的依赖关系用DAG描述执行引擎跑拓扑排序只有当前置Step全部成功后后置Step才会开始。每个Step支持独立配置重试次数超过阈值就标记任务失败并发送告警。2.2 接入层不管数据从哪来都变成Record流数据清洗面对的数据源五花八门。我们的实际环境里有MySQL业务库、有CSV文件、有HDFS上的离线导出后面还接了Kafka消息。如果每个来源都写一套处理逻辑规则引擎就得针对每个数据源重复实现一遍这绝对不行。所以接入层做了一个统一抽象任何数据源经过Reader插件后都转换为Record流进入清洗管线。Record在设计上参考了数据库行的概念包含一个Schema和一组字段值。考虑到清洗场景下字段的读取频率极高我们没有用繁琐的Map而是直接用数组存值、用Schema保存字段名和类型的对应关系。public class Schema { private final String[] fieldNames; private final FieldType[] fieldTypes; } public class Record { private final Schema schema; private final Object[] values; }这样设计的好处是规则算子在访问字段时能通过索引直接取到值不用做Map的哈希计算内存上也更紧凑。1.2亿行数据经过十几个算子每行少做几次哈希运算整体性能差距是非常可观的。2.3 规则引擎责任链上的每个算子各干一件事规则引擎是DataJet的核心。我们把清洗动作抽象为两类操作校验类算子判断一条记录是否满足条件不满足就进入脏数据通道。转换类算子对Record的字段做变换比如格式标准化、字段映射、数据补全。每个算子实现同一个Transform接口多个算子串成责任链一条Record按顺序流过所有算子。这样做最大的价值在于单一职责非常清晰新增一种清洗规则只需要写一个新的算子类不用动别的代码。举个实际例子一条门店销售数据进入清洗链路的完整流程是校验订单号格式不匹配则进入脏数据通道将时间字段统一为yyyy-MM-dd HH:mm:ss格式并转为东八区时间手机号字段做11位校验和脱敏处理金额字段由double转为BigDecimal负数和超出合理范围的值拦截按门店编码、订单号、业务日期做去重补充渠道字段缺失的填unknown。这六步对应六个算子配置在同一个Job里按顺序执行。每个算子只关注自己负责的那一件事后续要调整脱敏规则只需要改手机号算子不影响其他任何环节。3. 核心模块实现这些代码决定了工具的灵活性3.1 算子SPI新增一种清洗规则不用改引擎DataJet的算子扩展采用SPI机制。主工程里只定义接口和基础抽象类具体的算子是独立模块通过ServiceLoader加载。这样不同业务线可以按需提供自己的算子包主流程完全不需要感知。public interface Transform { String getName(); void init(MapString, String config); TransformResult process(Record record); }TransformResult不仅是成功或失败还带着失败原因和告警级别这个设计在后面做数据质量报告时帮了大忙。每条被拦截的记录都会带上算子的name和具体的拦截原因最终统一汇集成一张脏数据统计表按规则聚类展示。运营人员看到的不再是“数据有问题”这样一句空话而是“订单号格式校验拦截了1234条记录占比0.01%”这种可量化的指标。3.2 规则配置化用JSON描述一条清洗链路工具能不能被广泛使用关键看配置化程度。DataJet的每个Job都是JSON描述文件存储在一张配置表里。整个清洗链路长什么样直接在JSON里一目了然。下面是一个简化版的Job配置示例{ jobName: store_sale_clean_job, reader: { type: jdbc, jdbcUrl: jdbc:mysql://..., sql: select id, order_no, store_code, sale_time, amount, channel from t_store_sale where biz_date ${bizDate} }, transforms: [ {name: orderNoValidator, pattern: ^[A-Z]{2}\\d{12}$}, {name: timeNormalizer, field: sale_time, targetFormat: yyyy-MM-dd HH:mm:ss, timezone: Asia/Shanghai}, {name: phoneMasker, field: mobile, maskType: middle4}, {name: moneyConverter, field: amount, scale: 2, min: 0, max: 99999999}, {name: deduplicator, keys: [store_code, order_no, biz_date]}, {name: channelFiller, field: channel, defaultValue: unknown} ], writer: { type: jdbc, targetTable: ods_store_sale_clean, batchSize: 500 } }我第一次给团队看这个配置的时候很多人第一反应是“这不就是ETL工具里的配置项吗”。对道理是相通的但关键是这个JSON直接由执行引擎解析并驱动算子链运行不需要生成代码也不需要额外的解释层。规则的增删改就是改一条数据发布后下一次调度自动生效比改脚本再上线快得多。3.3 脏数据分流与质量报告清洗不只是删数据很多初做清洗工具的人容易忽略一件事——清洗不应该只是把脏数据扔掉而是要“让脏数据可见、可分析、可追溯”。如果一条数据被拦截后直接丢弃业务方事后根本查不到原因清洗工具就成了黑盒。所以DataJet在清洗管线中专门设计了脏数据分流机制。正常数据继续流向Writer脏数据则写入独立的脏数据表。每条脏数据记录都包含完整原始数据、算子名称、拦截时间、机器节点、错误原因等信息。脏数据表同时承担了一个重要职能——数据质量报告的数据源。每天跑完调度后自动生成本日报表统计各数据源的总体数据量、清洗通过率、各类规则拦截量、占比趋势。业务团队可以从报表里直观看到数据质量的波动而不是等下游报错才来排查。4. 性能从8小时到40分钟亿级数据清洗的优化链路4.1 第一轮瓶颈单线程逐行处理第一个可用版本上线后跑全量数据耗时8个半小时这个结果肯定是不能接受的。第一轮性能分析发现最大的问题在接入层——Jdbc Reader就是一条SQL查出来然后在一个循环里逐行处理没有任何并发和批处理的概念。当时我们用JProfiler看了一眼CPU和堆内存发现大量的时间耗在Object创建上。每一行数据都被包成一个Record每个字段都被转成String再转回目标类型光这个序列化动作就吃掉了大量CPU。第一轮优化很简单也很粗暴引入并行读取按主键范围把数据切分成多个分片每个分片一个线程独立读取和处理。long minId queryMinId(jobId); long maxId queryMaxId(jobId); int shardCount 16; long step (maxId - minId) / shardCount 1; for (int i 0; i shardCount; i) { long start minId i * step; long end Math.min(start step, maxId 1); executor.execute(() - processShard(job, start, end)); }这里要注意一个隐蔽的问题分片数量不能拍脑袋定。我们测试下来单机的情况并行度太高反而会导致数据库连接和线程切换的开销超过收益。机器是16核32G内存HikariCP连接池配了16个连接分片数设在16到24之间效果最稳定。并行度设置成CPU核心数的一倍到一点五倍是相对稳妥的起点。4.2 第二轮瓶颈一次一万行和一次一行差在哪第一轮优化后耗时降到3小时左右但离目标还有距离。继续排查发现瓶颈转移到了数据库读取方式上。MySQL JDBC默认的fetchSize是10也就是说每调一次next()JDBC驱动就需要往返数据库取一次数据。在逐行处理的代码里这意味着一行数据一次网络往返1.2亿行就是1.2亿次网络请求性能可想而知。解决方法是显式设置fetchSize配合MySQL服务端的游标读取模式。PreparedStatement ps connection.prepareStatement(sql, ResultSet.TYPE_FORWARD_ONLY, ResultSet.CONCUR_READ_ONLY); ps.setFetchSize(Integer.MIN_VALUE); // MySQL长文本/流式读取时必须用这个值设置之后next()时会从服务端批量拉取数据到客户端内存网络往返次数大幅下降。同时我们把单行处理改成积攒到一定量再批量统一处理减少中间对象的创建频率。这一轮优化后耗时降到了1小时10分钟左右。4.3 批处理写入与事务边界计算端提速之后写入端又成了瓶颈。最开始写入目标库是单条INSERT每条都走一次事务提交数据库WAL日志刷盘成了最大的开销来源。优化方案无外乎两个方向批处理插入和合适的事务边界。JDBC提供了addBatch和executeBatchMySQL从Connector/J 8.0.13开始需要在JDBC URL上额外加一个参数rewriteBatchedStatementstrue。不加这个参数驱动仍然会一条一条发送SQL批处理的效果约等于零这个细节很多教程不会提。我们最终把批大小设为500每500条记录打包成一个批次提交事务。理论上批越大性能越好但实践下来超过500之后单批次执行的失败重试成本随之增大。一条数据有问题整批500条都要回滚重试成本太高了这是需要权衡的地方。4.4 序列化选型Java原生Serializable的教训在优化过程中我们还踩了一个Java开发很容易踩的坑——使用Java原生Serializable做对象序列化。DataJet内部本来有一步将清洗后的Record缓存到本地磁盘防止OOM场景。最初图省事直接让Record实现Serializable结果在8G堆内存下频繁触发Full GC数据写盘慢得离谱。Java原生序列化的性能太差了不仅仅是慢而且序列化后的字节体积大并且Java序列化机制在反序列化时会执行大量反射调用。我们的Record结构非常规整字段类型固定完全没有必要用这种重机制。后来我们直接用自研的二进制编码方案Schema信息只写一次后续每条记录只按字段顺序写定长或变长值。序列化和反序列化都变成了简单的字节操作速度提升了不止一个量级。这个优化也提醒了我在大数据量场景下任何“通用机制”都可能成为瓶颈越贴近数据结构的代码越高效。优化链路走完之后最直观的效果是同样的1.2亿行数据从最初的8个半小时压到了40分钟以内而且是在单机多线程条件下完成的。5. 上线后踩过的三个大坑5.1 时区错八小时时间字段必须在接入层统一归一化第一个坑发生在测试环境非常有意思。门店上报的数据里有sale_time源库存的时间是业务系统生成的应用服务器时区是Asia/Shanghai但清洗工具所在的服务器如果设置成UTCJDBC读取时间字段后再转时间戳写回目标库时就会出现整整八小时的偏移。刚开始排查时我们反复检查清洗规则觉得时间标准化逻辑写得没问题。后来才发现问题根本不在清洗算子而在接入层读取数据时的时区解释。解决方式是在接入层统一约定所有时间字段在进入Record之前必须转换为UTC的long类型时间戳存储所有清洗算子只处理时间戳不感知本地时区写回目标库时再统一转成业务约定的东八区。这样一来不管上游业务系统是什么时区清洗工具内部始终保持一致不会再出现偏移。5.2 热点key打散去重场景下的两阶段方案第二个坑是去重算子上线后出现的。我们在做门店订单去重时按store_code作为去重key。结果某一线城市的门店数据量占了全量数据的30%单线程处理这个任务被拖死其他门店的数据处理完了卡在这一个分片上干等。这就是典型的数据倾斜。一开始我们的实现是在每个分片内做去重以为只要分片足够多就均匀了。实际上按门店维度去重热点门店的所有数据都落在同一个分片里怎么切都没用。后来改成两阶段去重第一阶段给去重key加一个随机后缀均匀打散到多个分片每个分片内做初步去重 第二阶段把第一阶段的结果按原始key重新分组再做一次精确去重。这个方案的代价是中间结果量会膨胀膨胀系数取决于打散倍数我们用的是8倍。但它能让去重算子的吞吐量不再被热点key卡死。如果你遇到类似的场景这个思路可以直接抄。5.3 浮点金额的科学计算double必须退场第三个坑估计不少搞过数据的人都有阴影金额字段用double存储在计算时出现精度误差。我们清洗规则里有一项“金额合理性校验”最开始是用double比较客户实付金额是否超过订单总额结果偶尔出现误判。排查时发现根源在于double的二进制浮点表示对十进制小数天生不精确比如0.1 0.2在Java里等于0.30000000000000004。当两个金额做比较时微小的误差就被放大了。统一改用BigDecimal之后精度问题解决了但要注意一个细节BigDecimal的构造函数要传String不能直接传double。new BigDecimal(0.1)得到的仍然是不精确的二进制值只有new BigDecimal(0.1)才是精确的十进制值。这属于BigDecimal最经典的坑之一网上随便一搜都是但实际项目里踩中的人依然不少。6. 后续演进从离线定时清洗到准实时清洗6.1 对接Kafka把批处理引擎改成微批DataJet第一版是纯离线调度每天凌晨跑昨天的数据T1出结果。但业务方很快提出新需求部分核心指标需要做到准实时延迟不超过5分钟。我们做了一个折中方案不重写引擎而是把Kafka当作新的Reader和Writer接入层。原来JDBC Reader按主键范围分片Kafka Reader则按Topic分区消费。清洗规则完全复用只是数据源从“批量查询”变成“实时消息流”写入端也改成按批次累积写入目标库达到批大小或时间窗口就提交一次。这样从离线到准实时的演进改动集中在接入层规则引擎和执行引擎都保持了稳定。6.2 规则管理后台让业务方自己配规则DataJet虽然实现了JSON配置化但配置都在代码仓库里管理每次调整规则对业务方来说还是不够友好。后面我们做了一个配套的管理后台前端用的Spring Boot加Vue 3提供Job创建、规则配置、调度管理、脏数据报表查看等功能。这里想分享一个经验规则配置页面尽量做成“所见即所得”不要只放一堆下拉框让用户猜。我们最初的后台就是简单表单用户配置完一条规则压根不知道执行效果是什么。后来加了一个“试运行”功能——用户导入一小批样本数据后台立即用当前配置跑一遍清洗流程把通过和拦截的记录展示出来。业务方配置规则的意愿和准确率明显上升了。6.3 个人体会这类工具项目怎么做才有价值如果读者是Java方向的学生或者准备面试以“大数据清洗工具”作为项目题目面试官大概率会追问几个问题为什么不用现成工具去重怎么做数据倾斜怎么解决性能瓶颈在哪这些问题如果只看过几篇博客很容易被问穿。但如果真的跟进过一遍数据从接入、清洗、分流到写入的全过程哪怕是自己模拟的数据也能讲出细节来。我的体会是清洗工具真正难的地方不在于“写代码”而在于数据质量问题的抽象、规则如何组合编排、以及大规模数据处理时的性能权衡。DataJet从出现到稳定运行前后迭代了三个多月踩坑无数但最终带来的收益是实打实的之前每周都要人工处理的脏数据事故现在基本清零业务不再打电话来问为什么报表数据对不上数据质量的指标从拍脑袋变成了每天有报表可看。如果你也要做类似的东西建议一开始就盯着数据质量度量和规则可配置这两个方向发力代码量反而不是第一优先级。数据清洗工具存在的意义不是“处理了多少条脏数据”而是“让数据质量问题暴露在明确的位置并让修复的成本降到最低”。这个大方向想明白了后面的设计决策基本不会跑偏。

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

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

免费获取报价