ARTICLE DETAIL

资讯详情

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

BSJ私有协议数据采集:Maven多模块与TCP/RabbitMQ链路实践

BSJ私有协议数据采集:Maven多模块与TCP/RabbitMQ链路实践 简介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.mainClasscom.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.bsjDEBUG这样能看到每条报文的解析结果。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的完整走通路径把上面三路合并就是一套合法的本地复现流程启动simulator设备端口开始监听。启动collector-tcp-client连接并读取设备报文。启动collector-sender把采集结果转发到接收端。启动collector-receiver确认数据到达。启动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开发机默认可能是GBKTCP流里的字节被按GBK解码自然乱码。解决所有字符转换的地方强制写明字符集。读取用StandardCharsets.UTF_8写文件用Files.write(path, data, StandardCharsets.UTF_8)不要依赖new String(bytes)这种不指定字符集的方法。启动JVM时再加-Dfile.encodingUTF-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 21 \ 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里加一个消息归一化层把原始报文重组成标准格式// 消息归一化统一为下游期望的字段命名 MapString, 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这三项写在一张纸上再启动进程。这个动作帮我少熬了不少夜——因为绝大多数采集问题不是算法难题而是配置版面上的低级不一致。希望帮到你。本文还有配套的精品资源点击获取
返回列表