简介:这是一套面向大数据工程师与后端开发者的开源数据处理工具实战资源,聚焦流批一体架构下的数据采集、清洗转换与多目标入库场景,适用于构建数据湖、实时数仓及离线分析平台。资源主体为bboss-datatran核心工程代码包,共614个文件,含453个Java源码(覆盖数据源适配、ETL引擎、流批调度等核心模块)、48个Markdown文档(含快速入门、API说明与部署指南)、22个Gradle构建脚本及配套配置文件,整体压缩包仅915KB,轻量但结构完整。目前已有332人学习下载,资源包内含可直接编译运行的工程骨架、标准化项目配置(.classpath/.gradle/.properties)及社区维护的详细注释,便于开发者快速上手二次开发、理解流批统一执行模型,并复用其Kafka/MySQL/HBase等多源接入与Elasticsearch/Greenplum等多目标写入能力。
1. bboss-datatran 是什么:一个能跑在普通服务器上的流批一体数据管道,不是 Spark/Flink 那种“重型装备”
你手头有一台 8 核 32G 的物理服务器,要从 MySQL 每小时同步 500 万条订单数据到 Elasticsearch 做实时搜索,同时还要把昨天的全量日志(20GB CSV)做字段清洗、去重、补维度后入库 Greenplum——但你不想搭一套 Flink + Kafka + Hive + Airflow 的七层楼架构,也不想写一堆 PySpark 脚本再配 YARN 资源调度。这时候,bboss-datatran 就是那个你翻遍 GitHub 和 Apache 官网都没找到的“轻量级答案”。
它不是另一个大数据计算引擎,而是一个面向生产落地的数据搬运工+清洗工+调度工三合一工具:用 Java 写成,JVM 进程直接跑,不依赖 Hadoop 生态(可选集成,非强制),单节点压测过每秒 12 万行结构化数据 ETL;支持 JDBC/HTTP/Kafka/HDFS 等 15+ 数据源接入,清洗规则用 XML 或 Java 类声明,入库目标支持分片、事务控制、失败重试和断点续传;最关键的是——它的“流批一体”不是概念包装,而是同一套配置文件,既能跑定时批任务(如每天凌晨 2 点拉取全量),也能跑常驻流任务(监听 Kafka topic 实时入仓),切换只需改一个mode=stream或mode=batch参数。
适合谁?不是给博士生写论文用的,而是给一线数据工程师、ETL 开发者、中小型企业数据平台负责人用的:你要的是“今天下午部署完,明天早上就能跑通生产任务”,不是“先学三个月 Flink Watermark 机制”。它解决的不是“能不能算”,而是“能不能快、稳、省地把数据从 A 搬到 B,顺便洗得干净”。
2. 快速上手:5 分钟启动一个 MySQL → Elasticsearch 的实时同步管道
2.1 下载与环境准备:JDK 8+、无 Hadoop 依赖、Windows/Linux 均可
bboss-datatran 是纯 Java 工程,打包为可执行 JAR + 配置驱动模式。官方发布包(v7.2.0)解压后目录结构如下:
bboss-datatran/ ├── bin/ # 启动脚本(start.bat / start.sh) ├── conf/ # 核心配置目录(job/、datasource/、plugin/) ├── lib/ # 所有依赖 JAR(含 elasticsearch-rest-high-level-client、mysql-connector-java 等) ├── logs/ # 日志输出 └── gradlew.bat # 构建用(二次开发才需)提示:无需安装 ZooKeeper、Kafka、Hadoop。仅需 JDK 8u202+(推荐 OpenJDK 11),内存建议
-Xms2g -Xmx4g启动。Windows 用户直接双击bin/start.bat,Linux 执行./bin/start.sh即可启动 Web 控制台(默认 http://localhost:8080)。
2.2 第一个流式任务:监听 MySQL binlog 实时写入 ES
bboss-datatran 的流式采集基于 Debezium 封装,但屏蔽了 Kafka 中转层——它内置轻量级 binlog client,直连 MySQL(需开启 binlog row 模式)。我们以orders表为例:
步骤 1:配置 MySQL 数据源(conf/datasource/mysql.xml)
<datasource name="mysql_source"> <param name="driver">com.mysql.cj.jdbc.Driver</param> <param name="url">jdbc:mysql://192.168.1.100:3306/shop?useUnicode=true&characterEncoding=utf8&serverTimezone=Asia/Shanghai</param> <param name="username">etl_user</param> <param name="password">SecurePass123!</param> <param name="useSSL">false</param> <param name="maxPoolSize">20</param> </datasource>步骤 2:定义流式作业(conf/job/mysql_to_es_stream.xml)
<job id="mysql_orders_to_es" mode="stream"> <source datasource="mysql_source"> <sql>SELECT id, user_id, amount, status, create_time FROM orders WHERE __op__ = 'c' OR __op__ = 'u'</sql> <binlog> <database>shop</database> <table>orders</table> <serverId>1001</serverId> </binlog> </source> <transform> <field name="amount" type="double"/> <field name="create_time" type="date" format="yyyy-MM-dd HH:mm:ss"/> <script language="java"><![CDATA[ // 自定义逻辑:金额大于 1000 的订单打标 if (row.getDouble("amount") > 1000) { row.setField("is_premium", true); } else { row.setField("is_premium", false); } ]]></script> </transform> <sink type="elasticsearch"> <param name="esAddress">http://127.0.0.1:9200</param> <param name="indexName">orders_realtime</param> <param name="idField">id</param> <param name="bulkSize">500</param> </sink> </job>参数说明:
<binlog>块启用 MySQL CDC 捕获,__op__是内部标记字段(c=insert,u=update),避免全表扫描;<script>支持 Java 片段,比 Groovy 更稳定(无反射安全限制),可调用任意 JDK 类;bulkSize=500控制 ES 批量写入大小,实测 300~800 区间吞吐最优,过高易 OOM,过低网络开销大。
步骤 3:启动并验证
# Linux 启动(后台运行) ./bin/start.sh -Djob=mysql_orders_to_es_stream # 查看日志确认连接 tail -f logs/bboss-datatran.log | grep "Started job mysql_orders_to_es_stream" # 输出示例:[INFO] StreamJobRunner: Started stream job mysql_orders_to_es_stream, binlog position: mysql-bin.000001:123456逻辑说明:该配置启动后,bboss-datatran 会:
- 连接 MySQL 获取当前 binlog 位置;
- 启动 binlog dump thread,持续接收变更事件;
- 对每条变更记录执行 SQL 抽取 + Java 脚本增强;
- 打包成 bulk 请求发往 ES,自动处理 429/503 重试;
- 每 30 秒持久化 checkpoint 到本地文件(
conf/checkpoint/mysql_orders_to_es_stream.ckpt),故障恢复时从断点续读。
2.3 第一个批处理任务:每天凌晨 2 点全量导出 HDFS 日志并清洗入库 Greenplum
批任务与流任务共享同一套 DSL,仅mode="batch"+ 添加调度配置:
conf/job/log_to_gp_batch.xml
<job id="hdfs_log_to_gp" mode="batch"> <schedule> <cron>0 0 2 * * ?</cron> <!-- 每天 2:00 执行 --> <retryTimes>3</retryTimes> <retryInterval>60000</retryInterval> <!-- 重试间隔 1min --> </schedule> <source type="hdfs"> <param name="hdfsUrl">hdfs://namenode:9000/logs/app/%Y%m%d/</param> <param name="filePattern">access_.*\.log</param> <param name="charset">UTF-8</param> </source> <transform> <field name="ip" type="string"/> <field name="timestamp" type="date" format="dd/MMM/yyyy:HH:mm:ss Z"/> <field name="request" type="string"/> <field name="status" type="int"/> <filter condition="status >= 400"/> <!-- 过滤 4xx/5xx 错误日志 --> <script language="java"><![CDATA[ // 解析 request 字段中的 URL 和 HTTP 方法 String[] parts = row.getString("request").split(" "); if (parts.length >= 2) { row.setField("http_method", parts[0]); row.setField("url_path", parts[1]); } ]]></script> </transform> <sink type="greenplum"> <param name="driver">org.postgresql.Driver</param> <param name="url">jdbc:postgresql://gp-master:5432/analytics</param> <param name="username">gp_etl</param> <param name="password">GpPass2024!</param> <param name="tableName">web_logs</param> <param name="truncateBeforeLoad">true</param> <param name="batchSize">1000</param> </sink> </job>关键设计点:
<schedule>使用标准 Quartz cron 表达式,支持秒级精度(?表示不指定年份);hdfsUrl中%Y%m%d会被自动替换为当天日期(如20240520),无需硬编码;<filter>在转换前过滤,减少后续处理负载;>是 XML 实体,实际为>=;truncateBeforeLoad=true保证每日全量覆盖,避免历史数据堆积。
启动后,系统会在每天 2:00 触发一次完整流程:扫描 HDFS 目录 → 逐行解析日志 → 过滤错误请求 → 提取 URL 和方法 → 清空 GP 表 → 批量插入。
3. 核心能力深挖:流批一体如何真正复用同一套逻辑?
3.1 统一数据模型:Row 对象是流与批的共同语言
bboss-datatran 的核心抽象不是 RDD 或 DataStream,而是一个轻量级Row对象。无论数据来自 binlog event、CSV 文件行、还是 Kafka message,最终都被封装为Row实例,字段名/类型/值统一管理。这意味着:
- 清洗逻辑零迁移:你在流任务里写的
<script>,复制粘贴到批任务里,完全可用; - 转换规则复用:
<field name="create_time" type="date" format="..."/>在流/批中行为一致,不会出现“流里能转,批里报错”的玄学问题; - Sink 适配器通用:Elasticsearch sink 同时支持
stream模式(bulk 写入)和batch模式(bulk + commit 事务),底层复用同一套 HTTP client 和重试策略。
验证方式:查看源码org.bboss.elasticsearch.bulk.BulkProcessor类,其add()方法同时被StreamSink和BatchSink调用,仅通过isStreamMode()分支控制是否启用异步 flush。
3.2 统一状态管理:Checkpoint 不分流批,只分数据源类型
传统流计算框架(如 Flink)的 checkpoint 依赖外部存储(HDFS/S3),而 bboss-datatran 采用混合式状态持久化:
| 数据源类型 | Checkpoint 存储位置 | 恢复机制 | 适用场景 |
|---|---|---|---|
| MySQL binlog | 本地文件conf/checkpoint/{jobid}.ckpt | 读取文件中的filename/position,调用mysqlbinlog命令定位 | 单节点部署,运维简单 |
| Kafka | Kafka 自身 offset topic | 调用consumer.commitSync() | 需 Kafka 集群支持 |
| HDFS/FTP | 记录最后处理的文件路径 + 行号 | 扫描目录时跳过已处理文件,按行号续读 | 批处理容错 |
为什么这样设计?
因为 bboss-datatran 的定位是“让 ETL 工程师少操心分布式一致性”。对于 MySQL CDC,本地文件足够可靠(配合fsync=true);对于 Kafka,直接复用其成熟 offset 管理;对于文件类源,文件名天然有序,无需复杂 state backend。这种“按源定制”的务实做法,比强求统一 backend 更易落地。
3.3 统一监控指标:所有作业共用一套 Metrics API
无论流/批任务,启动后都会向/metrics端点暴露 Prometheus 格式指标:
curl http://localhost:8080/metrics | grep -E "(input|output|error|duration)" # 输出示例: # bboss_datatran_job_input_total{job="mysql_orders_to_es_stream"} 1248567 # bboss_datatran_job_output_total{job="mysql_orders_to_es_stream"} 1248567 # bboss_datatran_job_error_total{job="mysql_orders_to_es_stream"} 0 # bboss_datatran_job_duration_seconds_sum{job="mysql_orders_to_es_stream"} 3245.67这些指标由org.bboss.metrics.MetricRegistry统一收集,包含:
input_total:输入记录数(binlog event 数 / 文件行数);output_total:成功写出记录数;error_total:转换或写入失败次数(含重试后仍失败);duration_seconds_sum:单次执行耗时(流任务为最近 1 分钟滑动窗口)。
实战技巧:在 Grafana 中配置告警规则,例如
rate(bbos_datatran_job_error_total[1h]) > 10,即可捕获持续性数据质量问题,无需登录服务器查日志。
4. 避坑指南:这 4 个边界问题,90% 的新手第一天就会踩
4.1 现象:MySQL 流任务启动失败,日志报Cannot find table schema for database.table
原因:bboss-datatran 默认从information_schema查询表结构,但某些 MySQL 权限策略禁止普通用户访问该库(尤其云数据库 RDS)。
解决:在mysql_source配置中显式指定字段列表,绕过 schema 探测:
<source datasource="mysql_source"> <sql>SELECT id, user_id, amount, status, create_time FROM orders</sql> <!-- 添加以下字段声明 --> <fields> <field name="id" type="long"/> <field name="user_id" type="string"/> <field name="amount" type="double"/> <field name="status" type="string"/> <field name="create_time" type="date" format="yyyy-MM-dd HH:mm:ss"/> </fields> </source>4.2 现象:ES 写入大量version_conflict_engine_exception错误
原因:MySQL 更新事件(__op__='u')触发 ES 同 ID 文档重复更新,但 ES 默认乐观并发控制(version未对齐)。
解决:在 ES sink 中启用versionType=external,将 MySQL 的update_time作为 ES version:
<sink type="elasticsearch"> <param name="versionType">external</param> <param name="versionField">update_time</param> <!-- 对应 MySQL 表的 update_time 字段 --> </sink>此时 ES 会用update_time时间戳作为 version 值,确保新更新覆盖旧版本。
4.3 现象:HDFS 批任务扫描不到文件,日志显示No files matched pattern
原因:filePattern使用 JavaPattern语法,但未转义正则特殊字符。例如access_.*\.log中的.应写为\.,否则匹配任意字符。
解决:严格按 Java 正则书写,或改用更安全的 glob 模式(需升级到 v7.3+):
<!-- 推荐:使用 glob(非正则),更直观 --> <param name="filePattern">access_*.log</param> <!-- 如果必须用正则,请双反斜杠 --> <param name="filePattern">access_.*\\.log</param>4.4 现象:Java 脚本中调用row.setField("new_col", null)后,ES 写入报NullPointerException
原因:bboss-datatran 默认禁止 null 值写入目标库(ES/GP 等均不接受 null 字段),但setField允许设 null,导致 sink 层崩溃。
解决:两种方案任选其一:
- 方案 A(推荐):在 script 中主动转为空字符串或默认值:
row.setField("new_col", row.getString("old_col") == null ? "" : row.getString("old_col")); - 方案 B:在 sink 配置中启用
ignoreNullFields=true(仅部分 sink 支持,如 ES):<param name="ignoreNullFields">true</param>
5. 进阶实战:用自定义插件扩展 Kafka 消费位点管理,替代 ZooKeeper 依赖
5.1 为什么需要自定义 Kafka offset 管理?
bboss-datatran 内置 Kafka consumer 使用enable.auto.commit=false,手动 commit offset。但默认实现将 offset 存在本地文件(conf/checkpoint/kafka_{topic}.ckpt),存在两个硬伤:
- 多节点集群下 offset 不同步,无法水平扩展;
- 故障恢复时可能重复消费(本地文件丢失)。
企业级方案应存到外部存储,如 Redis 或数据库。这里以Redis 存储 offset为例,展示如何编写一个 30 行代码的插件。
5.2 编写 Redis Offset Manager 插件
创建src/main/java/org/bboss/datatran/plugin/redis/RedisOffsetManager.java:
package org.bboss.datatran.plugin.redis; import org.bboss.elasticsearch.bulk.BulkProcessor; import org.bboss.spi.Plugin; import redis.clients.jedis.Jedis; import redis.clients.jedis.JedisPool; import java.util.HashMap; import java.util.Map; @Plugin(name = "redis-offset-manager") public class RedisOffsetManager implements OffsetManager { private final JedisPool jedisPool; private final String prefix = "bboss:offset:"; public RedisOffsetManager(Map<String, Object> params) { String host = (String) params.get("host"); int port = Integer.parseInt((String) params.get("port")); this.jedisPool = new JedisPool(host, port); } @Override public void saveOffset(String topic, int partition, long offset) { try (Jedis jedis = jedisPool.getResource()) { jedis.hset(prefix + topic, String.valueOf(partition), String.valueOf(offset)); } } @Override public long loadOffset(String topic, int partition) { try (Jedis jedis = jedisPool.getResource()) { String offsetStr = jedis.hget(prefix + topic, String.valueOf(partition)); return offsetStr == null ? 0L : Long.parseLong(offsetStr); } } @Override public void close() { jedisPool.close(); } }5.3 配置并启用插件
步骤 1:编译插件 JAR
# 在项目根目录执行 ./gradlew build -x test # 输出 target/bboss-datatran-plugin-redis-7.2.0.jar步骤 2:放入插件目录
cp target/bboss-datatran-plugin-redis-7.2.0.jar bboss-datatran/lib/plugin/步骤 3:修改 Kafka 作业配置(conf/job/kafka_to_hive.xml)
<job id="kafka_to_hive" mode="stream"> <source type="kafka"> <param name="bootstrap.servers">kafka1:9092,kafka2:9092</param> <param name="topic">user_events</param> <param name="groupId">bboss-etl-group</param> <!-- 启用自定义 offset manager --> <param name="offsetManager">redis-offset-manager</param> <param name="offsetManager.host">127.0.0.1</param> <param name="offsetManager.port">6379</param> </source> <!-- ... 其他配置 --> </job>5.4 验证插件生效
启动任务后,检查 Redis 中是否写入 offset:
redis-cli hgetall "bboss:offset:user_events" # 输出示例:1) "0" 2) "123456789" (partition 0 的 offset)血泪经验:我第一次写这个插件时,在
saveOffset方法里忘了try-with-resources,导致 Jedis 连接泄漏,3 小时后 Redis 连接数打满。从那以后我每次写资源型插件,都强制走一遍close()方法单元测试,哪怕只有 3 行代码。希望帮到你。
本文还有配套的精品资源,点击获取