注解与 Handler 实现详解
注解与 Handler 实现详解
传送门
当前示例源码位于 quick-start/custom-annotation-server
概述
本例通过自定义注解 + 继承框架 HandlerMapping 的方式,将 X-IoT Demo 协议的消息分发逻辑与业务处理解耦。核心流程:
- 定义自己的
@DemoMessageHandler和@DemoMessageMapping注解 - 编写继承
AbstractSimpleXtreamRequestMappingHandlerMapping的DemoMessageHandlerMapping - 编写
MyDemoHandler,用自定义注解声明消息处理逻辑 - 在服务端入口注册 HandlerMapping 和自定义调度器
1. 自定义注解
@DemoMessageHandler
@Target({ElementType.TYPE})
@Retention(RetentionPolicy.RUNTIME)
@Documented
@XtreamRequestHandler
public @interface DemoMessageHandler {
}- 元注解
@XtreamRequestHandler让框架扫描器在类路径上自动发现被标记的类 - 如果你的应用只需要一个 Handler 类,这个注解并非必须(可以手动注册),但注解方式更符合声明式风格
@DemoMessageMapping
@Target({ElementType.METHOD})
@Retention(RetentionPolicy.RUNTIME)
@Documented
@XtreamRequestHandlerMapping
public @interface DemoMessageMapping {
/// @return 当前处理器方法能处理的消息类型(对应 X-IoT Demo 协议的 `msgType` 字段)
int[] msgType();
/// `scheduler()` 别名 —— 对应 `@XtreamRequestHandlerMapping.scheduler()`。
///
/// 指定当前方法运行在哪个 `Scheduler` 上。
///
/// - 默认空字符串:使用框架默认的调度器(非阻塞)
/// - 指定名称:从 `XtreamSchedulerRegistry` 中查找
@AliasFor(annotation = XtreamRequestHandlerMapping.class, attribute = "scheduler")
String scheduler() default "";
}msgType()指定该方法能处理的消息类型(对应协议头中的 msgType 字段)scheduler()通过@AliasFor代理到框架的调度器机制,可指定方法运行在哪个调度器上- 元注解
@XtreamRequestHandlerMapping让框架将方法注册为可调用的XtreamHandlerMethod
2. HandlerMapping
public class DemoMessageHandlerMapping extends AbstractSimpleXtreamRequestMappingHandlerMapping {
private static final Logger log = LoggerFactory.getLogger(DemoMessageHandlerMapping.class);
public DemoMessageHandlerMapping() {
this(new String[]{XtreamUtils.detectMainClassPackageName()},
cls -> BeanUtils.createNewInstance(cls, new Object[0]));
}
public DemoMessageHandlerMapping(String[] basePackages,
Function<Class<?>, Object> instanceFactory) {
super(
new DefaultXtreamSchedulerRegistry(
Schedulers.parallel(), Schedulers.boundedElastic(), Schedulers.boundedElastic()),
new DefaultXtreamBlockingHandlerMethodPredicate(),
basePackages, instanceFactory
);
}
public DemoMessageHandlerMapping(String[] basePackages,
Function<Class<?>, Object> instanceFactory,
XtreamSchedulerRegistry schedulerRegistry) {
super(schedulerRegistry, new DefaultXtreamBlockingHandlerMethodPredicate(),
basePackages, instanceFactory);
}
@Override
public Mono<Object> getHandler(XtreamExchange exchange) {
// msgType 在报文头的第 5 字节(偏移 4),无符号
final int msgType = exchange.request().payload().getByte(4) & 0xFF;
log.info("Dispatching request with msgType={}(0x{})",
msgType, FormatUtils.toHexString(msgType, 2));
// FIXME 这里可以缓存到 Map 里,不用每次都遍历(为演示方便,这里直接遍历)
for (final var handlerMethod : handlerMethods) {
final DemoMessageMapping mapping =
handlerMethod.getMethod().getAnnotation(DemoMessageMapping.class);
if (mapping != null) {
for (final int type : mapping.msgType()) {
if (type == msgType) {
return Mono.just(handlerMethod);
}
}
}
}
log.warn("No handler found for msgType={}(0x{})", msgType,
FormatUtils.toHexString(msgType, 2));
return Mono.empty();
}
}关键点:
handlerMethods是父类AbstractSimpleXtreamRequestMappingHandlerMapping的protected字段,在构造时自动扫描并填充getHandler()的职责:从请求报文提取msgType,遍历已注册的handlerMethods,匹配到对应的方法后返回- 返回
Mono.empty()表示无匹配,框架会继续尝试其他 HandlerMapping 或返回 404
扫描机制
构造函数通过 XtreamUtils.detectMainClassPackageName() 自动推断扫描包路径:
public DemoMessageHandlerMapping() {
this(new String[]{XtreamUtils.detectMainClassPackageName()},
cls -> BeanUtils.createNewInstance(cls, new Object[0]));
}detectMainClassPackageName() 会寻找含有 main() 方法的类,并将其包名作为扫描根路径。如果自动推断不准确,可以手动指定包路径并将 schedulerRegistry 传入:
new DemoMessageHandlerMapping(
new String[]{"io.github.hylexus.xtream.quickstart.custom.annotation"},
cls -> BeanUtils.createNewInstance(cls, new Object[0]),
schedulerRegistry
)3. Handler 方法
3.1 心跳处理 — 通用应答
直接引用后端源码,避免文档和实现分叉:
/// 处理心跳请求 (`msgType=0x10`)。
///
/// ## 测试报文 (hex)
///
/// ```java
/// // 心跳请求:magic(4) + msgType(1) + bodyLength(2) = 7 字节,无消息体
/// 12 34 56 78 10 00 00
/// ```
///
/// 回复 `GenericAckResponse` 告知客户端心跳成功。
@DemoMessageMapping(msgType = {0x10})
@XtreamResponseBody
public Mono<GenericAckResponse> handleHeartbeat() {
log.info("Received heartbeat");
return Mono.just(new GenericAckResponse(0x10, 0));
}@XtreamResponseBody 是框架内置的回写注解。handler 方法只需返回 POJO,XtreamResponseBodyHandlerResultHandler 会自动使用 EntityCodec 将其编码为二进制并写回响应,无需手动操作 ByteBuf。
3.2 时间查询 — @XtreamResponseBody 自动编码响应
对应的后端实体也直接引用源码:
/// 服务器时间响应 (`msgType=0x11` 的响应,响应类型为 `0x81`)。
///
/// ## 报文格式
///
/// ```
/// +-------------+--------+-----------------------------------------------------+
/// | 偏移 | 长度 | 说明 |
/// +-------------+--------+-----------------------------------------------------+
/// | 0 | 4 | magic: u32, 固定 0x12345678 |
/// | 4 | 1 | msgType: u8, 固定 0x81 (时间查询响应) |
/// | 5 | 2 | bodyLength: u16, 消息体长度(大端),固定为 6 |
/// | 7 | 6 | serverTime: BCD[6], yyMMddHHmmss |
/// +-------------+--------+-----------------------------------------------------+
/// ```
@Getter
@Setter
@ToString(callSuper = true)
@ReferencedByDocs("guide/server/samples/custom-demo-protocol/handler-demo.md")
public class ServerTimeResponse extends AbstractEntity {
// region 消息体 (6 bytes)
// byte[7-13) 服务器时间 BCD[6] yyMMddHHmmss
@Preset.JtStyle.BcdDateTime
private LocalDateTime serverTime;
// endregion
public ServerTimeResponse() {
this.msgType = 0x81;
}
public ServerTimeResponse(LocalDateTime serverTime) {
this();
this.serverTime = serverTime;
}
}响应体使用 @Preset.JtStyle.BcdDateTime 注解编码为 BCD[6] yyMMddHHmmss 格式,非 i64 时间戳。
为什么推荐 @XtreamResponseBody?
- 避免手动
ByteBuf操作(容易内存泄漏) - 框架自动管理 buffer 生命周期(编码后自动 release)
- 统一编解码逻辑,请求和响应走同一套注解体系
3.3 温湿度上报 — @XtreamRequestBody 注入实体
@XtreamRequestBody 是框架内置注解,自动将消息体解码为指定实体类。对应方法如下:
/// 处理温湿度上报 (`msgType=0x12`)。
///
/// ## 测试报文 (hex)
///
/// ```java
/// // 温度 23.5°C (235=0x00EB), 湿度 60.0%RH (120=0x78)
/// // 完整报文: magic(4) + msgType(1) + bodyLength(2) + temperature(2) + humidity(1)
/// 12 34 56 78 12 00 03 00 EB 78
/// ```
///
/// 演示 `@XtreamRequestBody` 注解的参数注入:将消息体自动解码为实体类。
/// 回复通用应答告知客户端上报成功。
@DemoMessageMapping(msgType = {0x12})
@XtreamResponseBody
public Mono<GenericAckResponse> handleTemperatureReport(@XtreamRequestBody TemperatureReport report) {
log.info("Received temperature report: {}°C, {}%RH",
report.temperatureInCelsius(), report.humidityInPercent());
return Mono.just(new GenericAckResponse(0x12, 0));
}3.4 多传感器上报 — @XtreamRequestBody 注入实体
对应方法如下:
/// 处理多传感器数据上报 (`msgType=0x13`)。
///
/// ## 测试报文 (hex)
///
/// ```java
/// // 温度 22.5°C (225=0x00E1), 湿度 55.0%RH (110=0x6E)
/// // 气压 1013.2hPa (10132=0x2794), 风速 3.5m/s (35=0x0023)
/// // 时间戳 1700000000000ms (0x0000018B3F3B5A00)
/// // 完整报文: magic(4) + msgType(1) + bodyLength(2) + body(15)
/// 12 34 56 78 13 00 0F 00 E1 6E 27 94 00 23 00 00 01 8B 3F 3B 5A 00
/// ```
///
/// 演示 `@XtreamRequestBody` 注解的参数注入:将消息体自动解码为实体类。
/// 回复通用应答告知客户端上报成功。
@DemoMessageMapping(msgType = {0x13})
@XtreamResponseBody
public Mono<GenericAckResponse> handleMultiSensorReport(@XtreamRequestBody MultiSensorData report) {
log.info("Received multi-sensor report: {}°C, {}%RH, {}hPa, {}m/s, ts={}",
report.temperatureInCelsius(), report.humidityInPercent(),
report.pressureInHpa(), report.windSpeedInMeterPerSecond(), report.getTimestamp());
return Mono.just(new GenericAckResponse(0x13, 0));3.5 设备注册 — 混合注入 + 自定义调度器
同时注入 XtreamExchange 和 @XtreamRequestBody 实体,框架按参数类型自动选择解析器。对应方法如下:
/// 处理设备注册 (`msgType=0x14`)。
///
/// ## 测试报文 (hex)
///
/// ```java
/// // imei="868105040876543", productKey="AB"
/// // 完整报文: magic(4) + msgType(1) + bodyLength(2) + imeiLen(1) + imei(15) + productKeyLen(1) + productKey(2)
/// 12 34 56 78 14 00 13 0F 38 36 38 31 30 35 30 34 30 38 37 36 35 34 33 02 41 42
/// ```
///
/// 演示同时注入 `XtreamExchange` 和 `@XtreamRequestBody` 实体。
/// 回复注册结果告知客户端是否注册成功。
///
/// 使用 `scheduler = "business"` 指定在自定义的业务调度器上运行。
@DemoMessageMapping(msgType = {0x14}, scheduler = "business")
@XtreamResponseBody
public Mono<RegisterAckResponse> handleDeviceRegister(
XtreamExchange exchange,
@XtreamRequestBody DeviceRegisterRequest request) {
log.info("Received device register: imei={}, productKey={}, remote={}",
request.getImei(), request.getProductKey(), exchange.request().remoteAddress());
return Mono.just(new RegisterAckResponse(0, "registered OK"));
}通过 scheduler = "business" 指定该方法运行在自定义的 business 调度器上(Schedulers.newBoundedElastic(4, 100, "business")),适用于设备注册等可能涉及 IO 或耗时操作的场景。
3.6 报警上报
对应方法如下:
/// 处理报警上报 (`msgType=0x15`)。
///
/// ## 测试报文 (hex)
///
/// ```java
/// // alarmType=1 (通用报警), desc="overheat" (UTF-8)
/// // 完整报文: magic(4) + msgType(1) + bodyLength(2) + alarmType(2) + descLen(1) + desc(8)
/// // bodyLength = 0x0B (11 bytes)
/// 12 34 56 78 15 00 0B 00 01 08 6F 76 65 72 68 65 61 74
/// ```
///
/// 演示 `@XtreamRequestBody` 注解的参数注入:将消息体自动解码为实体类。
/// 回复通用应答告知客户端上报成功。
@DemoMessageMapping(msgType = {0x15})
@XtreamResponseBody
public Mono<GenericAckResponse> handleAlarmReport(@XtreamRequestBody AlarmReport report) {
log.info("Received alarm report: type={}, desc={}", report.getAlarmType(), report.getDesc());
return Mono.just(new GenericAckResponse(0x15, 0));
}4. 服务端入口
直接引用后端服务端入口源码:
public static void main(String[] args) {
// 如果你不了解 ResourceLeakDetector 是做什么的,请务必注释掉下面这行代码
// io.netty.util.ResourceLeakDetector.setLevel(io.netty.util.ResourceLeakDetector.Level.PARANOID);
final XtreamSchedulerRegistry schedulerRegistry = new DefaultXtreamSchedulerRegistry(
Schedulers.parallel(),
Schedulers.boundedElastic(),
Schedulers.boundedElastic()
);
schedulerRegistry.registerScheduler("business", Schedulers.newBoundedElastic(4, 100, "business"));
XtreamServers.tcp()
.name("custom-annotation-server")
.bind("0.0.0.0", 9527)
.customize(server -> server.doOnConnection(conn -> log.info("New connection: {}", conn)))
.pipeline(pipeline -> pipeline.addFirst(
// bodyLength 在报文头偏移 5 处,长度 2 字节
new LengthFieldBasedFrameDecoder(
1024, // maxFrameLength
5, // lengthFieldOffset: magic(4) + msgType(1)
2, // lengthFieldLength: bodyLength 占 2 字节
0, // lengthAdjustment
0 // initialBytesToStrip
)
))
.dispatch(dispatcher -> dispatcher
.addHandlerMappings(new DemoMessageHandlerMapping(
new String[]{"io.github.hylexus.xtream.quickstart.custom.annotation"},
cls -> io.github.hylexus.xtream.codec.core.utils.BeanUtils.createNewInstance(cls, new Object[0]),
schedulerRegistry
))
.enableBuiltinHandlers(EntityCodec.DEFAULT)
.addFilter(new LoggingXtreamFilter())
.addExceptionHandler(new LoggingXtreamRequestExceptionHandler())
)
.build()
.start();关键配置:
| 配置项 | 说明 |
|---|---|
LengthFieldBasedFrameDecoder(1024, 5, 2, 0, 0) | 解决 TCP 粘包,从第 5 字节读取 2 字节 bodyLength |
new DemoMessageHandlerMapping(...) | 注册自定义 HandlerMapping(含自定义调度器) |
schedulerRegistry.registerScheduler("business") | 注册业务调度器,用于设备注册等耗时操作 |
enableBuiltinHandlers(EntityCodec.DEFAULT) | 同时启用内置参数解析器和内置返回值处理器 |
如果你需要更底层地控制 Reactor Netty 的 customizer 顺序,仍然可以直接使用 XtreamServerBuilder。
5. 测试
启动 XtreamCustomAnnotationServerApp 后,使用 nc 命令测试:
# 心跳 (msgType=0x10)
echo -ne '\x12\x34\x56\x78\x10\x00\x00' | nc localhost 9527
# 服务器时间查询 (msgType=0x11)
echo -ne '\x12\x34\x56\x78\x11\x00\x00' | nc localhost 9527
# 温湿度上报 (msgType=0x12) — 23.5°C, 60.0%RH
echo -ne '\x12\x34\x56\x78\x12\x00\x03\x00\xeb\x78' | nc localhost 9527
# 多传感器数据上报 (msgType=0x13)
echo -ne '\x12\x34\x56\x78\x13\x00\x0f\x00\xe1\x6e\x27\x94\x00\x23\x00\x00\x01\x8b\x3f\x3b\x5a\x00' | nc localhost 9527
# 设备注册 (msgType=0x14) — imei=868105040876543, productKey=AB
echo -ne '\x12\x34\x56\x78\x14\x00\x13\x0f\x38\x36\x38\x31\x30\x35\x30\x34\x30\x38\x37\x36\x35\x34\x33\x02\x41\x42' | nc localhost 9527
# 报警上报 (msgType=0x15) — alarmType=1, desc="overheat"
echo -ne '\x12\x34\x56\x78\x15\x00\x0b\x00\x01\x08\x6f\x76\x65\x72\x68\x65\x61\x74' | nc localhost 9527