☰
BSJ私有协议数据采集:Maven多模块与TCP/RabbitMQ链路实践
2026/10/3 2:54:56 网站建设 项目流程

简介:bsj协议数据采集.zip 是一份围绕BSJ数据采集协议的完整资源包,面向从事数据采集、后端通信开发的工程师及协议研究者,可帮助解决采集链路效率不高、协议实现不透明等问题。压缩包共84个文件,体积仅116KB,以Java为主(61个源码文件),配合11个XML配置、4个properties属性文件及IDEA工程模块文件,整体按Maven多模块结构组织,便于按采集模拟、发送、接收、HTTP服务、消息队列等环节分类查阅。资源实现了从模拟器到客户端、服务端再到RabbitMQ的完整采集链路,既有可直接运行的工程样例,也有协议格式、连接处理、数据解析与存储等核心逻辑,可支撑二次开发或协议移植。目前已有182人学习下载,适合用于协议分析、数据采集系统设计与性能优化。

1. bsj协议数据采集:从压缩包里拆出一条可复用的采集链路

做数据采集的同行应该都有同感:现场设备只认私有协议,文档缺失,抓包抓到凌晨还要靠猜来判断报文含义。我拆完这份bsj协议数据采集.zip之后,最大的体会是——它把设备侧模拟、TCP通道采集、HTTP接口接入、RabbitMQ消息投递和协议解析这几件事整合成了一个完整的Maven工程。说人话就是,你拿到的不只是某个协议的文本说明,而是一条可以直接参照搭建的采集链路,尤其适合给注塑机、传感器这类工业设备做数据接入,或者需要把私有协议收编成标准消息的场景。搞懂这套工程,比单独啃协议文档更省时间,因为你在看的是别人已经跑通的落地形态。

2. Maven多模块拆解:七个模块各管一段数据流水线

刚解压这个zip时,文件树并不复杂,但bsj-master内部的分层方式值得先讲清楚。这不是一个把所有代码堆在一个src里的单体项目,而是把采集这件事按数据流向拆成了多个Maven模块。每个模块名字都带着明确的职责,比如simulator负责模拟设备端,collector-tcp-client负责从TCP端口收数据,collector-rabbitmq-client负责对接消息队列。先看整体结构:

bsj-master/ ├── simulator # 模拟器:模拟 BSJ 设备端,产生报文 ├── collector-common # 公共模块:编解码、工具类、常量 ├── collector-tcp-client # TCP 采集客户端:连接设备端口 ├── collector-sender # 发送器:把采集结果发往接收端 ├── collector-receiver # 接收器:承接采集数据并处理 ├── collector-http-server # HTTP 服务:对外提供采集接口 └── collector-rabbitmq-client # MQ 客户端:投递/消费消息队列

这个结构说明一件事:BSJ协议数据采集不是单点程序,而是由数据源、采集端、传输层、接入层组成的流水线。simulator是模拟源头,它生成的报文经过collector-tcp-client进入系统,再通过collector-sender和collector-receiver完成内部流转;collector-http-server则给外部系统留了一个HTTP入口;collector-rabbitmq-client负责把数据异步投递到MQ。collector-common承载公共代码,避免多个模块重复写解析逻辑。

2.1 模块职责与它们在链路中的位置

把模块和实际数据流向对应起来,比单独记模块名更直观。整理成下表:

模块职责定位在链路中的位置典型动作
simulator模拟BSJ协议设备端链路最上游监听端口,定时生成设备报文
collector-tcp-client采集端连接simulator建立TCP连接,读取数据帧
collector-sender内部转发采集端之后把解析后的数据发往接收端
collector-receiver内部接收sender之后接收并处理采集结果
collector-http-server对外接入旁路接口提供HTTP接口,接收外部数据
collector-rabbitmq-client异步消息任意环节连接MQ,投递或消费消息
collector-common公共支撑被其他模块依赖报文解析、常量、工具类

这里有个容易被新手忽略的设计点:sender和receiver分开,意味着采集与处理可以部署在两台机器上。模拟器在一台机器跑,采集端在另一台机器连它,接收端放在数据中心——这种拓扑在工业数据采集里非常常见。如果你只在本地学习,全放一个进程里也完全能跑通,但理解了这个拆分逻辑,后续做多机部署时就不会懵。

2.2 先构建再研究:Maven多模块的编译要点

拿到压缩包后不要急着看单个模块的代码,先把整个工程构建一遍。多模块项目的父子依赖关系决定了编译顺序,不按顺序构建的话,collector-tcp-client会找不到collector-common的类。项目根目录有pom.xml,在bsj-master下执行标准构建命令:

mvn clean install -DskipTests

命令里包含两个关键动作:clean把每个模块的target目录清掉,避免上次编译的旧class干扰;install则把每个模块的jar装进本地Maven仓库,这样下游模块才能通过依赖坐标引用到它们。-DskipTests跳过单元测试,只编译不跑测试用例,节省时间。如果本地还没装过Maven,先执行mvn -version确认环境,缺依赖时加-U强制更新快照。

构建完成后,在每个模块的target目录下能看到对应jar包。此时再用IDE打开项目,Maven会自动识别模块依赖,collector-common会被其他模块引用。我习惯在构建时额外开一个终端扫一眼日志,确认有没有BUILD SUCCESS字样,而不是只看IDE有没有报错。

2.3 模块内部的主流程猜测与验证路径

光看结构还不够,想弄明白BSJ协议怎么跑起来,最好在源码里搜索关键入口。simulator模块里大概率会有一个main方法或者Spring Boot启动类,负责启动Socket服务;collector-tcp-client里则会有发起连接的客户端代码。顺着ServerSocket.accept()或Socket.connect()往上追,就能定位报文生成的逻辑。

常见做法是在每个模块的src/main/resources下找application.yml或application.properties,端口号、主题名、队列名这些可配参数都集中在那里。比如simulator监听哪个端口、tcp-client连接哪个地址、rabbitmq的交换机叫什么,都是先看配置文件再去看硬编码。这样做的原因是,实际部署时大概率不会用代码里写死的默认值,配置文件的优先级一定高于代码直觉。

3. 三路落地:模拟器、TCP客户端、RabbitMQ消费的完整链路

到了这一步,手里已经有一个能编译通过的工程,接下来要做的是让它真正转起来。BSJ协议数据采集这套资源里,simulator是设备端的模拟,collector-tcp-client是采集源头,collector-rabbitmq-client是消息出口。把这三个模块串起来,就形成了一条最小的完整链路。

3.1 先把模拟器拉起来:它扮演的是设备端

模拟器的作用是生成符合BSJ协议的数据,让采集端有东西可收。没有真实设备时,它是调试采集逻辑的关键。用Maven的exec插件可以直接从源码运行某个模块,不需要手动打jar包:

mvn -pl simulator -am exec:java -Dexec.mainClass=com.bsj.simulator.SimulatorServer

-pl simulator指定当前构建的模块,-am表示同时构建它依赖的兄弟模块,比如collector-common。exec.mainClass指定启动类的全限定名,实际类名要以仓库源码为准,解压后直接搜public static void main就能找到。如果模拟器是一个Spring Boot工程,也可以先mvn -pl simulator -am package,再java -jar simulator/target/simulator*.jar启动。

启动成功后,模拟器会监听某个端口,同时周期性打印发送日志。看到类似send packet -> 192.168.1.10:9001的日志,就说明设备端已经往外吐数据了。此时可以用telnet 127.0.0.1 端口号连上去看一眼原始报文,这比直接看代码更直观。

3.2 TCP采集端:主动连接还是被动接入

BSJ采集通常采用客户端主动连接设备端的模式,collector-tcp-client就是干这件事的。它启动后会根据配置连接simulator监听的地址,读入字节流,再交给collector-common里的解析逻辑。一个典型的连接核心代码长这样:

// 这是采集端发起连接的核心逻辑,实际类名以源码为准 Socket socket = new Socket(); socket.connect(new InetSocketAddress(host, port), 3000); socket.setSoTimeout(5000); InputStream in = socket.getInputStream(); BufferedReader reader = new BufferedReader(new InputStreamReader(in, StandardCharsets.UTF_8)); String line; while ((line = reader.readLine()) != null) { // 每读到一行完整报文,就交给解析器处理 collector.onMessage(line); // 这里可以做本地归档,也可以立即转发 } socket.close();

代码里两个超时参数值得说明:connect的3000毫秒是建连超时,设备不在线时不会让主线程卡死;setSoTimeout(5000)是读超时,防止连接半死时一直阻塞。用BufferedReader按行读取,是文本协议最简单的处理方式,前提是报文以换行符结尾。StandardCharsets.UTF_8必须显式指定,后面避坑章节会专门讲原因。

这个模块的配置项通常是bsj.host、bsj.port、bsj.charset这类键值,改端口时只需要动配置文件,不需要重新编译。我一般会在配置文件里把日志擦写到底,logging.level.com.bsj=DEBUG,这样能看到每条报文的解析结果。

3.3 sender与receiver:内部数据如何流转

采集端拿到报文后不会一直攒在内存里,而是交给collector-sender发出,由collector-receiver接收。这两个模块有点像生产者和消费者:sender把采集结果序列化后通过网络发送,receiver反序列化后做后续处理。如果链路中间没有这两个模块,采集端就得自己承担所有业务逻辑,一旦处理逻辑变重,采集性能会直线下降。

在本地验证时,可以先把这两端的日志级别调到INFO,观察sender的发送确认和receiver的接收确认。数据能成功流转,说明链路是通的;如果sender发出去但receiver没反应,优先检查两边的IP和端口是否互通,用nc -zv测一下端口连通性比看日志快得多。常见的落地方式是sender发JSON字符串、receiver收到后写入文件或数据库,这个阶段不必过度设计。

3.4 HTTP接入与RabbitMQ客户端:两个旁路方案

collector-http-server解决的是外部系统主动把数据推过来的场景。比如某些设备不支持TCP,但支持定期回调HTTP接口。这个模块启动后,外部系统向它的URL发送POST请求,body里带BSJ报文或者标准JSON,服务端解析后进入统一处理流程。用curl自测:

curl -X POST http://127.0.0.1:8081/collect \ -H "Content-Type: application/json" \ -d '{"device_id":"SN-2024-08B","type":"measure","temp":26.4,"press":1.82}'

-X POST指定请求方法,-H声明内容类型,-d携带报文主体。服务端返回200说明数据已接收,返回4xx则需要检查字段命名是否与接收类的属性对得上。HTTP接入的好处是外部系统不需要理解BSJ协议细节,只要按约定POST数据即可。

collector-rabbitmq-client则提供异步消息能力。采集端把数据投递到队列后,消费端可以按自己的节奏处理,天然具备削峰填谷的作用。配置项里最重要的是交换机名称、队列名称和路由键,三者不一致时消息会静默丢失,这在避坑章节会展开。投递代码通常是注入RabbitTemplate后调用convertAndSend(exchange, routingKey, message),值得留意的是消息体序列化格式,默认JDK序列化会有兼容性问题,JSON序列化更稳妥。

3.5 从设备到MQ的完整走通路径

把上面三路合并,就是一套合法的本地复现流程:

  1. 启动simulator,设备端口开始监听。
  2. 启动collector-tcp-client,连接并读取设备报文。
  3. 启动collector-sender,把采集结果转发到接收端。
  4. 启动collector-receiver,确认数据到达。
  5. 启动collector-rabbitmq-client,观察队列中的消息积压曲线。

只要队列里出现消息,这就算正式跑通了。最容易忽略的是配置文件中主机地址写成了远端测试机地址,而模拟器跑在本地。我的排查习惯是先从simulator和tcp-client的视角分别执行netstat -an | grep 端口号,确认监听端存在且连接状态是ESTABLISHED,再往上查消息队列。

4. BSJ报文结构与解析器:从ByteBuf到可落库的Java对象

采集链路搭通之后,真正考验工程能力的是协议解析这一段。很多采集项目的翻车现场都发生在这一层:报文解析错一位,后续所有统计都跟着错。从simulator的源码和collector-common的公共代码来看,BSJ协议大概率是文本行式协议——每行一条记录,字段用固定前缀或分隔符组织。下面给出一种非常典型的报文模板:

#BSJ/1.0 device_id:SN-2024-08B seq:881240 type:measure ts:1728118400 payload:{"temp":26.4,"press":1.82,"status":1}

这种结构的好处是肉眼可读、调试方便,缺点是字段顺序一旦变化,解析逻辑必须跟着调整。逐行说明一下字段含义:

字段含义解析注意事项
#BSJ/1.0协议版本标识用于判断是否BSJ报文
device_id设备标识可能出现中划线等特殊字符
seq序列号用于去重和乱序检测
type消息类型心跳、测量、告警等
ts时间戳单位通常是秒或毫秒
payload业务数据内部可能是JSON或键值对

4.1 粘包与半包:文本协议解析的第一个坎

TCP是流式协议,数据到达应用层时可能一次读入多条报文,也可能一条报文被拆成多个包——这就是粘包和半包问题。用readLine()按行读取时,如果设备端发送时没有严格遵守换行符结束,粘包会导致解析器拿到拼接后的垃圾数据。我处理这类问题的固定思路是:先定义一个ByteBuf累积缓冲区,每收到新数据先追加进去,再按分隔符扫描出完整帧。

以下是基于Netty常见的解码思路,也适用于原生Socket:

// 按行拆帧,解决粘包和半包问题 ByteBuf buf = Unpooled.buffer(); while (buf.isReadable()) { int idx = buf.indexOf(buf.readerIndex(), buf.writerIndex(), (byte) '\n'); if (idx < 0) { // 还没读到完整一行,继续等下一次数据到达 break; } ByteBuf line = buf.readSlice(idx - buf.readerIndex() + 1); String raw = line.toString(StandardCharsets.UTF_8).trim(); if (raw.startsWith("#BSJ/1.0")) { // 进入一条新报文的状态 parseBsjFrame(raw); } } buf.discardReadBytes(); // 释放已读空间

这段代码里,indexOf负责扫描换行符位置,找不到则说明当前缓冲区内没有完整报文,必须等下一批数据到来。找到换行符后,用readSlice把这一整行切出来交给解析器。注意最后的discardReadBytes(),它把已读区域清掉,避免缓冲区无限增长。

4.2 从原始行到业务对象的解析流程

拿到完整文本行后,接下来就是纯粹的字符串处理。常见做法是按换行把报文头、字段区、payload区分开,再逐字段填充到POJO。为了不让解析代码散落各处,这类工程通常会在collector-common里统一封装:

// 简化版BSJ报文解析器,核心是字段拆分与类型转换 public class BsjMessageParser { public BsjMessage parse(String rawText) { BsjMessage msg = new BsjMessage(); String[] lines = rawText.split("\n"); for (String line : lines) { if (line.startsWith("#BSJ")) { msg.setVersion(line.substring(4).trim()); } else if (line.startsWith("device_id:")) { msg.setDeviceId(line.substring(10).trim()); } else if (line.startsWith("seq:")) { msg.setSeq(Long.parseLong(line.substring(4).trim())); } else if (line.startsWith("type:")) { msg.setType(line.substring(5).trim()); } else if (line.startsWith("ts:")) { msg.setTs(Long.parseLong(line.substring(3).trim())); } else if (line.startsWith("payload:")) { msg.setPayload(line.substring(8).trim()); } } return msg; } }

注意substring的参数要和键名长度严格匹配,比如device_id:这个前缀长度是10,截取后要调用trim()消除首尾空格。seq和ts是数字字符串,转换时如果设备端填了非数字字符,NumberFormatException会直接抛出来。这里显示了所有线上的陷阱,我通常在外层加异常捕获,解析失败的报文记录到独立日志,而不是让整个采集线程挂掉。

4.3 二进制变体的处理思路

有些私有协议的实现会走二进制定长帧,不在文本行式报文这个范围内。BSJ如果存在二进制变体,解析思路也很接近:先按固定帧头找到帧起始位置,再按帧头里的长度字段截取一帧,最后按偏移量解析各字段。与之相比没有哪个更好的问题,只有哪个和真实设备端对齐的问题。判断方法是直接看设备端的发送代码或者抓包数据的前几个字节:能看到#BSJ就是文本协议,看见0xAA 0x55这类魔数就要走二进制解析。弄清楚这一点,你才知道该用BufferedReader还是该用ByteBuf.readInt()。

5. 踩坑记录:端口、编码、队列积压和多模块构建的五个典型翻车现场

链路跑通是一回事,跑得稳是另一回事。这套工程在本地复现时最容易翻车的地方,我整理成了五条,每一条都是拆项目时实际见过的现场。

5.1 模拟器起来了,TCP客户端却一直“连接被拒”

现象:simulator日志显示监听成功,但collector-tcp-client一直报ConnectException: Connection refused。

原因:端口不一致。simulator读取的是application.yml里默认端口,而tcp-client连接的是代码中硬编码或另一份配置里的端口。两个数字对不上,连接必然失败。

解决:先确认两边各自读的端口号。在simulator机器上执行netstat -anp | grep 端口号,确认监听地址是0.0.0.0还是127.0.0.1。如果监听的是127.0.0.1,外部机器连不上;如果监听的是0.0.0.0但客户端仍连不上,检查防火墙。统一做法是把端口配置从源码挪到配置中心或环境变量,避免代码里出现魔法数字。

5.2 RabbitMQ队列有积压,但消费者就是不消费

现象:rabbitmq管理后台看到bsj.queue里有几千条消息,消费者连接数却是0。

原因:消费者没有真正启动,或者消费者监听的是另一个队列。工程里@RabbitListener指定的队列名少写了一个字符是最常见的场景,消息投递到了bsj.queue,消费者却盯着bsj-queue。

解决:先用管理API确认队列名,在命令行执行rabbitmqctl list_queues name messages consumers,然后对照消费者代码里的@RabbitListener(queues = "队列名")。同时检查交换机绑定关系:rabbitmqctl list_bindings,看路由键是否匹配。路由键不匹配时,消息会进队列但消费者收不到,处理方式是把交换机、队列和路由键集中定义成常量类,三个地方引同一份值。

5.3 修改代码后重新构建,行为还是老样子

现象:改了collector-common里的解析逻辑,重启tcp-client后,解析结果一点没变。

原因:多模块项目里,collector-common改动后没有重新install到本地仓库,tcp-client构建时引用的还是旧jar。这是Maven多模块最经典的一个坑。

解决:每次修改公共模块后,在根目录执行mvn clean install -DskipTests,确保本地仓库里是最新版本。验证方式是用jar tf查看jar包里的class文件时间戳,或者直接解压看改动类是否在包里。从那以后我每次改完公共代码都会看一眼target/xxx.jar的生成时间,确认构建真的成功了再启动。

5.4 报文中中文乱码,字段解析全部错位

现象:payload里出现???或者温度,JSON解析直接崩。

原因:设备端发送的是UTF-8编码,而解析端用了系统默认字符集。Windows开发机默认可能是GBK,TCP流里的字节被按GBK解码,自然乱码。

解决:所有字符转换的地方强制写明字符集。读取用StandardCharsets.UTF_8,写文件用Files.write(path, data, StandardCharsets.UTF_8),不要依赖new String(bytes)这种不指定字符集的方法。启动JVM时再加-Dfile.encoding=UTF-8兜底,虽然不能完全替代代码层面的明确指定,但能降低运行环境差异导致的解析风险。

5.5 连接没有断,但采集数据突然停了

现象:tcp-client进程还活着,日志不再输出,接收端也收不到新数据。

原因:设备端如果长时间没数据,连接处于半开状态。上游设备悄然断开后,下游无法感知。代码里没做心跳检测的时候,TCP连接可能一直占着直到系统超时。

解决:在tcp-client里增加心跳读超时逻辑,客户端定期发送心跳报文,超过阈值没收到回应就主动重连。常见做法是每隔30秒检查一次连接状态,连接不可用时关闭旧Socket并重新connect。这个逻辑在避坑层面属于必备项,否则机房一个瞬时网络抖动,采集就静默中断。

6. 本地验证与扩展:一条命令确认链路健康,再把数据送进其他系统

本地调试到这步,最希望的就是有一条命令能把整条链路的状态一次看清楚,而不是开四个终端来回切。我习惯把验证逻辑写成一个检查脚本,涵盖端口连通性、HTTP探活和队列积压三个维度。

#!/bin/bash # 本地验证链路健康状态 check_tcp() { nc -zv 127.0.0.1 "$1" >/dev/null 2>&1 \ && echo "[OK] 端口 $1 可连接" \ || echo "[FAIL] 端口 $1 连接失败" } check_tcp 9001 # simulator监听端口 check_tcp 9002 # tcp-client连接端口 curl -sf http://127.0.0.1:8081/actuator/health >/dev/null \ && echo "[OK] http-server 心跳接口正常" \ || echo "[FAIL] http-server 异常" rabbitmqctl list_queues messages consumers \ | grep -E "bsj|^" \ && echo "[OK] 队列连接已确认"

脚本里三个检查点各有讲究:nc -zv只测端口通不通,不建立长连接;curl -sf中的-s静默模式避免输出响应体,-f让HTTP错误码直接使命令失败;rabbitmqctl list_queues messages consumers同时输出队列积压量和消费者数,光有积压没有消费者,说明链路后半段没通。端口号根据实际配置替换,这套骨架放到别的采集项目里也能直接用。

链路健康之后,再往前一步就是扩展。BSJ报文里的payload通常是JSON,但设备ID、时间戳、序列号这些字段在不同系统里命名不一致。为了让下游系统统一消费,我会在rabbitmq-client里加一个消息归一化层,把原始报文重组成标准格式:

// 消息归一化,统一为下游期望的字段命名 Map<String, Object> standard = new LinkedHashMap<>(); standard.put("deviceId", msg.getDeviceId()); standard.put("eventTime", Instant.ofEpochSecond(msg.getTs()).toString()); standard.put("data", parsePayload(msg.getPayload()));

这段代码做的核心事情是:让下游不再关心BSJ协议细节,只管收标准JSON。LinkedHashMap保持字段顺序,Instant.ofEpochSecond(...)把秒级时间戳统一成ISO8601字符串,方便其他系统直接按字符串时间处理。归一化之后,数据既可以继续推给其他MQ,也可以写进时序数据库做告警分析。

聊到这儿,最想跟你分享的习惯是:从那以后我每次新接一套设备采集,都会先把端口、字符集、队列路由key这三项写在一张纸上再启动进程。这个动作帮我少熬了不少夜——因为绝大多数采集问题不是算法难题,而是配置版面上的低级不一致。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询