资讯动态

Java多模块数据采集工程:TCP粘包处理与RabbitMQ接入避坑实战

发布时间:2026/10/9 17:37:51 来源:尧图企业网站定制
简介bsj协议数据采集.zip 是一份面向Java开发者与数据采集工程师的完整资源包围绕BSJ协议提供从理论解析到工程落地的闭环。内含项目源码、采集工具与数据集覆盖数据格式、编码方式及错误处理机制可用于构建高效、稳定的采集系统并验证其准确性。压缩包共84个文件以61个Java源码为主配合11个XML配置以Maven的pom.xml为主、4个properties属性文件和6个iml模块描述另有readme.md说明文档整体仅116KB目录结构清晰按 collector-common、collector-tcp-client、collector-sender、collector-receiver、collector-http-server、collector-rabbitmq-client 等模块组织便于按需研读和二次开发。目前已有185人学习适合希望深入了解采集协议实现、或准备自建数据抓取通道的中级及以上Java开发者参考实践。通过阅读源码可掌握连接建立、收发数据、解析存储等关键环节借助模拟器与数据集可快速测试采集性能并指导真实业务部署。1. bsj协议数据采集拆开这份资源的第一眼结论bsj协议数据采集这个 zip 拆开之前我以为是某个单一协议的实现包拆完才发现它是一整套覆盖“采集—转发—接收—存储”的 Java 多模块工程跑在 Maven 的聚合结构上。它对应的场景很典型设备或上游服务按固定协议把数据推过来你要做的是收下来、校验、转发、落库而不是只写一个爬虫脚本。这个包里最大的亮点是同时提供了 TCP 长连接采集、HTTP 服务端接收、RabbitMQ 消息队列三种数据通道外加一个模拟器生成数据适合正在搭采集管道、又不想从网络层写起的开发者和运维。但前提是先把模块职责和消息流转搞清楚否则连先启动哪个模块、先配哪份参数都会懵。2. 先拆协议再拆工程七个模块的职责与数据流向2.1 BSJ 协议不是标准是一份双方约定的消息契约拆完这个包的第一个结论是别去搜什么 BSJ 协议标准文档它不是公开规范而是这套采集系统上下游之间约定好的通信契约。协议存在的意义是让数据在链路上传输时双方对“边界怎么切、字段怎么排、错误怎么发现”这三件事达成一致。只要遵守同一个约定采集端和发送端各自独立演进都没有问题。从代码层面看collector-common 就是协议的翻译层。它定义了公共消息体、字段常量、编解码工具模拟器造数据用它TCP 客户端解析也用它HTTP 服务端接收还用它。换句话说任何新模块想接入这套链路唯一要遵守的就是 collector-common 里那套消息定义其他模块对它来说是透明的。常见的帧结构是五段式魔数 版本号 长度 载荷 校验。魔数用来快速识别这是一个合法帧版本号用来区分协议演进长度字段告诉解析方这次要读多少字节载荷是真正的业务数据校验码兜底检测数据在传输过程中有没有被改坏。这套设计在采集场景里几乎是标准答案你在读源码时看到类似结构不用惊讶。帧字段作用典型长度魔数快速判断是否为合法消息帧2 字节版本号区分协议迭代版本1 字节长度载荷字节数指导粘包切分2 字节载荷业务数据本体不定长校验检测数据完整性2 字节这里最容易被低估的是长度字段。TCP 是流式协议数据没有天然的分隔线一次 read 可能读到多条完整的消息也可能只读到半条。长度字段就是用来做消息切分的锚点先读固定大小的头部拿到载荷长度之后再决定继续读多少字节。这个动作贯穿所有采集模块也是后面粘包/半包问题的根源我会在避坑章展开讲。另外要说清楚这套资源的定位是“协议数据采集”不是网页爬虫。它处理的是设备或系统按协议推送的结构化数据比如传感器读数、日志流、业务事件而不是从 HTML 页面里抽取内容。如果你的需求是网页抓取这套 TCP/MQ 架构可能用不上但如果场景是设备接入、数据管道、消息流转这套结构可以直接抄。2.2 七个模块的职责边界与依赖关系核心工程是 bsj-master用 Maven 聚合了多个子模块。我在下面把每个模块的角色和技术要点列出来模块角色技术要点bsj-master父工程统一依赖版本聚合子模块simulator数据源模拟器按协议生成消息帧控制频率与条数collector-common公共模块消息体定义、编解码、校验工具collector-tcp-clientTCP 采集入口长连接、粘包处理、断线重连collector-sender数据转发批量聚合推给下游collector-receiver数据接收端接收 sender 数据落库或再分发collector-http-serverHTTP 接收入口提供 POST 接口接收结构化数据collector-rabbitmq-client消息队列接入消费/投递 RabbitMQ 消息读这张表的时候重点看依赖方向collector-common 是最底层所有采集模块都依赖它simulator 和 collector-tcp-client 是同一对一个造数据一个收数据collector-sender 和 collector-receiver 是第二对负责把采集到的数据搬运到下一站HTTP 和 RabbitMQ 则是两种不同的出口。理解这个依赖方向比背模块名有用得多——改协议时只动 common加采集通道时只抄已有模块的结构不用动整条链路。为什么要拆这么多模块而不是写成一个单体程序因为采集链路每一跳的稳定性要求不一样。TCP 入口要处理断线重连转发端要做批量聚合接收端可能要落库混在一起的话某一环节翻车会拖垮整条链。拆开之后你可以单独重启 sender 而不影响 TCP 连接也可以在 receiver 端加消费逻辑而不动采集入口排障时边界清晰很多。阅读源码时建议按这个顺序先看 collector-common 的消息体定义与解码器这是全链路的地基再看 simulator 如何构造数据了解数据长什么样最后看 tcp-client 如何消费。这个顺序和调试链路的顺序一致能减少很多困惑。2.3 一条数据从产生到落库要过几关把各模块串起来数据流向是这样的模拟器按协议生成消息帧通过 TCP 长连接发给 collector-tcp-client。TCP 端读完字节流后做两件事用长度字段做粘包拆分再做字段校验。解析出来的结构化消息交给 collector-sender。sender 是聚合器攒够一批或者时间窗口到了就批量转发给 collector-receiver。receiver 拿到数据后可以选择直接落库也可以投递到 HTTP 接口或者 RabbitMQ交给后续业务消费。这条链路里有个容易被忽略的设计数据是“多跳”的每一跳都有可能失败。模拟器发出后 TCP 连接断了、sender 转发时对端拒绝、receiver 落库时数据库抖动都会让数据停在半路。所以链路里每一跳都必须有重试机制这不是能不能省的问题而是采集系统的基本功。你后面测试时会发现真正花时间的不是让链路通而是让链路在断断续续的情况下还能保证数据不丢不漏。如果你只想验证某一段逻辑不需要全链路启动。单独跑 simulator tcp-client sender 就能看解析结果只想验证 HTTP 出口直接拿 curl 往 http-server 发数据就行。下面两章就按“先拆 TCP 链路、再拆两个出口”的顺序展开。3. 把 TCP 采集链路跑起来模拟器、客户端与参数配置3.1 环境准备JDK、Maven 与依赖顺序先说结论JDK 8 和 Maven 3.6 是必须的RabbitMQ 只有在用到 collector-rabbitmq-client 时才需要只跑 TCP HTTP 链路可以先不装。确认环境的命令很简单java -version mvn -versionjava -version 输出 1.8 或 11 以上都行mvn -version 主要看 Apache Maven 那一行有没有正常打印。如果提示找不到命令先把 JDK 的 bin 目录配进 PATH再回来继续。工程用 Maven 聚合第一次构建要做一次整体 install把 collector-common 装进本地仓库否则其他模块单独跑会报依赖找不到。我一般会在根目录先跑cd bsj-master mvn clean install -DskipTests这条命令做了两件事clean 清掉上次编译产物install 把每个模块装到本地 ~/.m2 仓库。collector-common 不 install后面单跑任意模块时 IDE 或命令行都会解析不到它。跳过测试是建议而不是必须因为这份资源里的测试用例往往要连真实端口新手环境容易因为端口占用直接挂掉。等你熟悉了再打开测试跑也不迟。IDE 导入时直接把 bsj-master 的 pom.xml 作为 Maven Project 打开等右下角依赖索引转完即可。如果打开后某个模块标红大概率是本地仓库里没有对应依赖回到命令行执行一次 mvn clean install 基本能解决。多模块工程排错永远先看依赖有没有 install 到位再看代码本身。3.2 先起模拟器确认数据源是通的采集链路里第一环是 simulator。它做的工作很简单按协议格式化消息通过 TCP 推出去。运行前要确认三个参数监听端口、推送频率、数据条数。端口决定了 TCP 客户端连哪里频率和条数决定你能不能快速看到流量。启动命令cd simulator mvn exec:java -Dexec.mainClasscom.bsj.simulator.SimulatorMain \ -Dexec.args--port 9000 --interval 100 --count 1000参数含义--port 9000 是模拟器绑定的 TCP 端口--interval 100 表示每 100 毫秒发一条--count 1000 表示最多发 1000 条。如果工程里模拟器类名与这个不一样先看 src 目录下的实际包路径再替换 mainClass这个不影响整体流程。参数含义调试建议--port模拟器监听端口与 TCP 客户端保持一致--interval发送间隔毫秒先 200链路稳定再调小--count发送总条数首次验证 500 够用interval 是最容易踩坑的参数。它控制的是生产速率而生产速率决定了下游会不会积累 backlog。本地验证时先 200 毫秒一条等链路稳定了再往 20 甚至 10 毫秒调否则可能 TCP 粘包还没处理完下一批数据又到了。启动后让这个进程保持运行它得像水龙头一样持续供水后面才有数据可采。3.3 collector-tcp-client长连接、粘包切分与字段解析TCP 客户端是采集的核心入口。它做的事有三件建立长连接、从字节流里切出完整消息帧、把帧解析成结构化对象。连接参数集中在配置文件里核心三项是 host、port、bufferSize。核心循环逻辑如下// TCP 客户端采集主循环关键逻辑 Socket socket new Socket(config.getHost(), config.getPort()); InputStream in socket.getInputStream(); ByteBuffer buffer ByteBuffer.allocate(config.getBufferSize()); while (running) { int read in.read(buffer.array(), buffer.position(), buffer.remaining()); if (read 0) { // 对端关闭连接走重连逻辑不直接退出 reconnect(config); continue; } buffer.position(buffer.position() read); ListMessageFrame frames MessageCodec.decode(buffer); for (MessageFrame frame : frames) { // 解析出的消息交给 sender 批量转发 sender.offer(frame); } }这里要重点说明 buffer.remaining() 的作用——它保证每次 read 不会越界写坏 ByteBuffer。decode 返回的是 List 而不是单条消息因为 TCP 流里一次 read 可能包含多条帧也可能只读了半条帧这是粘包/半包问题的根源。MessageCodec 内部会先读完固定长度的头部再根据长度字段读取完整载荷解析不了的部分留在 buffer 里等下一次 read。这段逻辑是采集端最容易写错的地方有人会把读完的 buffer 直接 clear导致半包数据被丢弃还有人会一次性读固定大小把两条帧拆错位置解析出来的字段全乱。这两个错误都属于没有理解“流式”这个概念。host 和 port 要与模拟器对得上port 不一致是最常见的连不上原因。bufferSize 决定了单次 read 能吞多少字节本地调试 4096 够用线上建议 16384。太小会频繁触发多次读取浪费 CPU 和 IO太大浪费内存而且单次 read 返回的数据量也不会超过 TCP 接收窗口所以不必追求极致大。这几个参数没有“最优值”只有“适合你的流量”的值。调试思路是先低频率跑通再逐步加压观察 CPU 和内存变化找到你场景下的稳定区间。如果在实际项目里设备侧不方便维护长连接或者采集端要临时接第三方数据TCP 入口不是唯一选择。你也可以直接用 collector-http-server 的 POST 接口作为入口这种方式对调用方最友好不需要懂协议底层。但要说明白HTTP 每条消息都有请求头和连接建立开销吞吐量和实时性都不如 TCP 长连接。所以选型逻辑一般是这样设备数量少、数据量大、实时要求高走 TCP外部系统对接、数据量中等、希望接入简单走 HTTP。这个包把两条路都留好了这也是我推荐先完整跑一遍的原因——同样的数据你能直观看到两条路的差异。4. 数据出口怎么选HTTP 服务端与 RabbitMQ 的接入参数4.1 collector-http-server把采集结果暴露成 POST 接口当接收端希望以 HTTP 方式拿数据时collector-http-server 就派上用场了。它本质上是一个内嵌 HTTP 服务对外提供一个 POST 接口sender 或任意上游把数据推给它。好处是接入方不需要懂协议只需要会发 HTTP 请求对跨语言、跨团队协作特别友好。配置项集中在 properties 文件里核心四项如下# collector-http-server 配置示例 server.port8080 server.path/api/collect/receive server.max-threads200 server.batch-size100max-threads 控制并发处理线程数batch-size 是单次批量写入的条数。两者配合决定吞吐上限。batch-size 设太大单次请求体就大HTTP 超时风险上升设太小请求次数变多网络开销变大。我自己的习惯是在 100 到 500 之间调再根据单条消息体大小反推——单条 1KB 时 500 条就是 500KB常见网关都扛得住单条 10KB 时 500 条就到 5MB很多网关会直接拒绝。接口调试用 curl 就能做不需要把整个链路跑起来curl -X POST http://127.0.0.1:8080/api/collect/receive \ -H Content-Type: application/json \ -d [{id:msg-001,type:sensor,payload:2024-01-01T00:00:00Z|35.6}]调用方拿到 200 就表示服务端接收成功非 2xx 状态码要能识别并重试。用 HTTP 出口时有一个必须考虑的点接口幂等。采集系统重发是常态TCP 断线重连后会补发HTTP 接口必须对重复消息去重否则下游会拿到重复数据。常见方案是拿消息帧里的消息 ID 做唯一键落库前先查一次或者用数据库唯一索引兜底。4.2 collector-rabbitmq-client消息队列接入与 ack 策略走消息队列是为了解耦。采集端不直接面对下游业务而是把消息投到 RabbitMQ谁需要谁去消费。队列天生具备削峰能力突发流量到来时消息堆在队列里消费端按自己的节奏处理不会把下游打爆。启动命令mvn exec:java -Dexec.mainClasscom.bsj.collector.rabbitmq.RabbitMqClientMain \ -Dexec.args--host 127.0.0.1 --queue bsj.data.queue --prefetch 50host 是 RabbitMQ 地址queue 是消费的队列名prefetch 是消费端预取数量。prefetch 这个参数很容易被忽略但它直接影响吞吐设太小消费端一次取太少网络往返多设太大消息堆在消费端内存里消费端一旦挂掉这些消息就处于 unacked 状态恢复后要重新处理。50 到 100 是大部分场景的安全区间。消费端拿到消息后有一个动作必须做——手动 ack。RabbitMQ 默认自动 ack消息一投递给消费者就认为成功如果消费逻辑抛异常消息就丢了。改用手动 ack 后处理成功才 basicAck失败可以重回队列或进入死信// RabbitMQ 消费回调关键逻辑 channel.basicConsume(queue, false, (consumerTag, delivery) - { try { MessageFrame frame MessageCodec.decode(delivery.getBody()); repository.save(frame); channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } catch (Exception e) { // 处理失败重回队列等待下次消费或投递到死信队列 channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true); } }, consumerTag - {});这段代码里basicConsume 的第二个参数传 false 表示手动确认。try 块里业务处理成功后调用 basicAckcatch 块里调用 basicNack 并把第三个参数 requeue 设为 true 让消息重回队列。注意如果消费逻辑有幂等性要求重试之前要先做好去重否则一条处理失败的消息反复重试会反复触发副作用。4.3 sender 与 receiver批量聚合与重试策略sender 和 receiver 是链路中段的运输工具。sender 从采集入口拿消息聚合到一定数量或时间窗就批量发出receiver 负责接收并确认。核心价值是降低传输次数一条条发网络开销高且容易被打爆批量发效率高但需要处理“批内部分失败”的情况。批量触发条件通常有两个数量窗口和时间窗口哪个先到就先执行。数量窗口让吞吐可预期时间窗口保证低流量时数据不会积压太久。比如配置 500 条或 2 秒触发一次流量大时每 500 条发一批流量小时每 2 秒发一批两条路都能保证数据的及时性。receiver 处理部分失败时我看到最多的是“整批拒绝”和“逐条确认”两种策略。整批拒绝逻辑简单但一条坏消息会拖累整批重发逐条确认效率低但能保证坏消息不牵连好数据。如果对数据完整性要求高建议用逐条确认把失败消息单独路由到重试队列不要和正常消息混在一起。这里没有银弹只能根据下游的容错能力做取舍。5. 避坑与排查采集链路最常见的五个翻车点5.1 模拟器启动了TCP 客户端就是连不上现象模拟器打印了监听端口TCP 客户端启动后反复报 Connection refused。原因模拟器绑定了 127.0.0.1 而不是 0.0.0.0客户端连接的却是机器的局域网 IP 或 localhost 以外的地址数据包根本没到模拟器。这个情况在本地调试时非常容易发生——服务端默认绑定 localhost客户端却自作聪明地填了机器的局域网 IP。解决先把两端统一成 127.0.0.1 或 localhost 验证连通性确认通了之后再考虑要不要绑 0.0.0.0 对外服务。排查步骤是先分别在两端执行 netstat 看监听地址再做一次 telnet 127.0.0.1 9000 确认端口能通最后再去看代码。这个顺序能省下大量时间。5.2 收到的数据字段错位看起来像乱码现象解析出来的字段值和模拟器发出的对不上字符串字段出现移位或者乱码。原因帧长度字段和实际载荷长度不一致。常见于改了模拟器的消息体字段却没有同步更新 collector-common 里的长度计算逻辑导致解析时按旧长度截取后面的字段全部错位。解决确认 collector-common 的编解码版本和模拟器版本一致重点检查长度字段的单位是字节数还是字符数中文字符是否按 UTF-8 计算。我建议准备一组固定测试数据跑一次采集后做逐字段比对任何错位立刻能定位到是长度计算还是编码问题。这一步多做一次后面排障能少花一半时间。5.3 HTTP 服务端在大流量下丢数据现象低频测试一切正常加大推送频率后 HTTP 接口返回超时部分数据没有落库。原因server.max-threads 太小请求排队处理不过来的请求超过超时时间被客户端判失败或者 batch-size 太大单请求处理时间过长拖垮了整体吞吐。解决先调大 max-threads再降 batch-size两个参数交叉调整。我的经验是优先保证单请求在 1 秒内完成再靠并发把吞吐撑起来这个思路比盲目加内存靠谱。调整完一定要看两个指标接口平均响应时间和线程池活跃线程数而不是只看“好像没报错”。5.4 RabbitMQ 消息积压消费端日志却干干净净现象队列消息数持续上涨消费端日志没有报错但消息就是不减少。原因消费端开启了手动 ack但代码里没有调用 basicAck消息全部处于 unacked 状态。RabbitMQ 认为这些消息还没被确认所以不会投递给其他消费者也不会把它从队列里删掉。解决检查消费回调末尾是否有 channel.basicAck 调用。很多人在这个点上翻车因为系统不报错、日志很干净只能靠队列监控发现消息数只增不减。加了手动 ack 就必须在成功分支显式确认这是硬规矩。我一般会在确认前后各加一条 debug 日志确认后被消费的消息数能和入口数对得上。5.5 mvn package 报错找不到 collector-common 依赖现象单独构建 collector-tcp-client 模块时提示 Cannot resolve collector-common 或者找不到符号。原因collector-common 没有先 install 到本地仓库其他模块解析不到它。这是新导入多模块 Maven 工程最常见的问题和代码本身无关。解决回到工程根目录执行 mvn clean install -DskipTests确保公共模块先进入本地仓库再单独构建具体模块。批量构建时注意保持根目录的 reactor 顺序不要只构建单个模块。如果你用 IDE 打开了工程但右上角显示 Maven 未导入也要先重新导入再执行构建。这几个坑看似零散背后有一条共同的主线先确认谁能通、再看数据对不对、最后盯有没有丢。连不上是网络层的问题字段错位是协议层的问题丢数据是容量层的问题消息积压是确认机制的问题依赖缺失是构建顺序的问题——每类问题都有自己固定的排查入口顺着链路一层层看比在代码里乱翻高效得多。6. 进阶用数据集做完整性与一致性的交叉验证采集链路跑通只是第一步能不能放心上线靠的是验证。这套资源里的数据集正好用来做整套管线的对账。方法不复杂把数据集作为模拟器的输入源逐条发送采集端全部接收后比对两边的条数、关键字段值、校验和。数量对不上说明链路有丢字段对不上说明编解码有错。我习惯在 sender 和 receiver 各加一个计数埋点入口统计收到多少条出口统计转发成功多少条。两边数字相等只能说明传输没丢还要抽几条原始数据对比内容是否一致。字段级别的不一致往往比数量不一致更隐蔽——数量对得上但内容错位说明长度字段或编码逻辑有问题这在第 5 章讲过。一个值得养成的习惯是每次跑链路都把采集结果落一份 CSV然后用 diff 工具和数据集源文件做逐行对比。写一段 shell 脚本循环批量发送、批量比对整个过程可以自动化几分钟跑完一轮。第一次发现差异时优先怀疑长度字段和编码方式而不是网络。从那以后我每次跑采集链路都会强制走一遍“启动模拟器→采集落盘→diff 数据集”的闭环。链路改过配置、改过协议、换过出口都先对一遍账再考虑上线。希望帮到你。本文还有配套的精品资源点击获取

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

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

免费获取报价 →
↑