2026-06-26 20b96473f2520590a0dca6b775b81e3ea06a77a0
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
package cn.iocoder.yudao.module.iot.gateway.protocol.modbus.common.utils;
 
import cn.iocoder.yudao.module.iot.core.biz.dto.IotModbusPointRespDTO;
import cn.iocoder.yudao.module.iot.gateway.protocol.modbus.tcpclient.manager.IotModbusTcpClientConnectionManager;
import com.ghgande.j2mod.modbus.io.ModbusTCPTransaction;
import com.ghgande.j2mod.modbus.msg.*;
import com.ghgande.j2mod.modbus.procimg.InputRegister;
import com.ghgande.j2mod.modbus.procimg.Register;
import com.ghgande.j2mod.modbus.procimg.SimpleRegister;
import com.ghgande.j2mod.modbus.util.BitVector;
import io.vertx.core.Future;
import lombok.experimental.UtilityClass;
import lombok.extern.slf4j.Slf4j;
 
import static cn.iocoder.yudao.module.iot.gateway.protocol.modbus.common.utils.IotModbusCommonUtils.*;
 
/**
 * IoT Modbus TCP 客户端工具类
 * <p>
 * 封装基于 j2mod 的 Modbus TCP 读写操作:
 * 1. 根据功能码创建对应的 Modbus 读/写请求
 * 2. 通过 {@link IotModbusTcpClientConnectionManager.ModbusConnection} 执行事务
 * 3. 从响应中提取原始值
 *
 * @author 芋道源码
 */
@UtilityClass
@Slf4j
public class IotModbusTcpClientUtils {
 
    /**
     * 读取 Modbus 数据
     *
     * @param connection Modbus 连接
     * @param slaveId    从站地址
     * @param point      点位配置
     * @return 原始值(int 数组)
     */
    public static Future<int[]> read(IotModbusTcpClientConnectionManager.ModbusConnection connection,
                                     Integer slaveId,
                                     IotModbusPointRespDTO point) {
        return connection.executeBlocking(tcpConnection -> {
            try {
                // 1. 创建请求
                ModbusRequest request = createReadRequest(point.getFunctionCode(),
                        point.getRegisterAddress(), point.getRegisterCount());
                request.setUnitID(slaveId);
 
                // 2. 执行事务(请求)
                ModbusTCPTransaction transaction = new ModbusTCPTransaction(tcpConnection);
                transaction.setRequest(request);
                transaction.execute();
 
                // 3. 解析响应
                ModbusResponse response = transaction.getResponse();
                return extractValues(response, point.getFunctionCode());
            } catch (Exception e) {
                throw new RuntimeException(String.format("Modbus 读取失败 [slaveId=%d, identifier=%s, address=%d]",
                        slaveId, point.getIdentifier(), point.getRegisterAddress()), e);
            }
        });
    }
 
    /**
     * 写入 Modbus 数据
     *
     * @param connection Modbus 连接
     * @param slaveId    从站地址
     * @param point      点位配置
     * @param values     要写入的值
     * @return 是否成功
     */
    public static Future<Boolean> write(IotModbusTcpClientConnectionManager.ModbusConnection connection,
                                        Integer slaveId,
                                        IotModbusPointRespDTO point,
                                        int[] values) {
        return connection.executeBlocking(tcpConnection -> {
            try {
                // 1. 创建请求
                ModbusRequest request = createWriteRequest(point.getFunctionCode(),
                        point.getRegisterAddress(), point.getRegisterCount(), values);
                if (request == null) {
                    throw new RuntimeException("功能码 " + point.getFunctionCode() + " 不支持写操作");
                }
                request.setUnitID(slaveId);
 
                // 2. 执行事务(请求)
                ModbusTCPTransaction transaction = new ModbusTCPTransaction(tcpConnection);
                transaction.setRequest(request);
                transaction.execute();
                return true;
            } catch (Exception e) {
                throw new RuntimeException(String.format("Modbus 写入失败 [slaveId=%d, identifier=%s, address=%d]",
                        slaveId, point.getIdentifier(), point.getRegisterAddress()), e);
            }
        });
    }
 
    /**
     * 创建读取请求
     */
    @SuppressWarnings("EnhancedSwitchMigration")
    private static ModbusRequest createReadRequest(Integer functionCode, Integer address, Integer count) {
        switch (functionCode) {
            case FC_READ_COILS:
                return new ReadCoilsRequest(address, count);
            case FC_READ_DISCRETE_INPUTS:
                return new ReadInputDiscretesRequest(address, count);
            case FC_READ_HOLDING_REGISTERS:
                return new ReadMultipleRegistersRequest(address, count);
            case FC_READ_INPUT_REGISTERS:
                return new ReadInputRegistersRequest(address, count);
            default:
                throw new IllegalArgumentException("不支持的功能码: " + functionCode);
        }
    }
 
    /**
     * 创建写入请求
     */
    @SuppressWarnings("EnhancedSwitchMigration")
    private static ModbusRequest createWriteRequest(Integer functionCode, Integer address, Integer count, int[] values) {
        switch (functionCode) {
            case FC_READ_COILS: // 写线圈(使用功能码 5 或 15)
                if (count == 1) {
                    return new WriteCoilRequest(address, values[0] != 0);
                } else {
                    BitVector bv = new BitVector(count);
                    for (int i = 0; i < Math.min(values.length, count); i++) {
                        bv.setBit(i, values[i] != 0);
                    }
                    return new WriteMultipleCoilsRequest(address, bv);
                }
            case FC_READ_HOLDING_REGISTERS: // 写保持寄存器(使用功能码 6 或 16)
                if (count == 1) {
                    return new WriteSingleRegisterRequest(address, new SimpleRegister(values[0]));
                } else {
                    Register[] registers = new SimpleRegister[count];
                    for (int i = 0; i < count; i++) {
                        registers[i] = new SimpleRegister(i < values.length ? values[i] : 0);
                    }
                    return new WriteMultipleRegistersRequest(address, registers);
                }
            case FC_READ_DISCRETE_INPUTS: // 只读
            case FC_READ_INPUT_REGISTERS: // 只读
                return null;
            default:
                throw new IllegalArgumentException("不支持的功能码: " + functionCode);
        }
    }
 
    /**
     * 从响应中提取值
     */
    @SuppressWarnings("EnhancedSwitchMigration")
    private static int[] extractValues(ModbusResponse response, Integer functionCode) {
        switch (functionCode) {
            case FC_READ_COILS:
                ReadCoilsResponse coilsResponse = (ReadCoilsResponse) response;
                int bitCount = coilsResponse.getBitCount();
                int[] coilValues = new int[bitCount];
                for (int i = 0; i < bitCount; i++) {
                    coilValues[i] = coilsResponse.getCoilStatus(i) ? 1 : 0;
                }
                return coilValues;
            case FC_READ_DISCRETE_INPUTS:
                ReadInputDiscretesResponse discretesResponse = (ReadInputDiscretesResponse) response;
                int discreteCount = discretesResponse.getBitCount();
                int[] discreteValues = new int[discreteCount];
                for (int i = 0; i < discreteCount; i++) {
                    discreteValues[i] = discretesResponse.getDiscreteStatus(i) ? 1 : 0;
                }
                return discreteValues;
            case FC_READ_HOLDING_REGISTERS:
                ReadMultipleRegistersResponse holdingResponse = (ReadMultipleRegistersResponse) response;
                InputRegister[] holdingRegisters = holdingResponse.getRegisters();
                int[] holdingValues = new int[holdingRegisters.length];
                for (int i = 0; i < holdingRegisters.length; i++) {
                    holdingValues[i] = holdingRegisters[i].getValue();
                }
                return holdingValues;
            case FC_READ_INPUT_REGISTERS:
                ReadInputRegistersResponse inputResponse = (ReadInputRegistersResponse) response;
                InputRegister[] inputRegisters = inputResponse.getRegisters();
                int[] inputValues = new int[inputRegisters.length];
                for (int i = 0; i < inputRegisters.length; i++) {
                    inputValues[i] = inputRegisters[i].getValue();
                }
                return inputValues;
            default:
                throw new IllegalArgumentException("不支持的功能码: " + functionCode);
        }
    }
 
}