From 45b4f79c16cec8129e6867f8cba1270214e1ead6 Mon Sep 17 00:00:00 2001 From: yanghanbin Date: Mon, 26 Jan 2026 18:08:01 +0800 Subject: [PATCH] init --- .gitignore | 44 +++ README.md | 196 +++++++++++ UDP_MULTICAST_RECEIVER.md | 178 ++++++++++ UDP_MULTICAST_SENDER.md | 216 ++++++++++++ pom.xml | 87 +++++ .../dc/tranlog/TranLogServiceApplication.java | 58 +++ .../afe/dc/tranlog/config/AsyncConfig.java | 53 +++ .../dc/tranlog/config/MulticastConfig.java | 51 +++ .../afe/dc/tranlog/config/ScheduleConfig.java | 18 + .../internal/DataQueryController.java | 111 ++++++ .../internal/DataReceiverController.java | 56 +++ .../internal/TradeDayController.java | 62 ++++ .../tranlog/domain/dto/AccumulatedData.java | 31 ++ .../tranlog/domain/dto/BuySellStrength.java | 39 +++ .../afe/dc/tranlog/domain/dto/COPItem.java | 94 +++++ .../tranlog/domain/dto/TradeDayResponse.java | 44 +++ .../tranlog/domain/model/entity/BsRecord.java | 39 +++ .../domain/model/entity/MtcRecord.java | 70 ++++ .../domain/model/entity/TranRecord.java | 87 +++++ .../domain/model/entity/TranTable.java | 148 ++++++++ .../service/AccumulatedLogicService.java | 111 ++++++ .../service/BuySellStrengthService.java | 166 +++++++++ .../com/afe/dc/tranlog/service/COPParser.java | 206 +++++++++++ .../dc/tranlog/service/MTCLogicService.java | 244 +++++++++++++ .../service/MulticastReceiverService.java | 272 ++++++++++++++ .../tranlog/service/TranDatabaseService.java | 193 ++++++++++ .../dc/tranlog/service/TranLogService.java | 262 ++++++++++++++ .../dc/tranlog/service/TranTableService.java | 237 +++++++++++++ .../service/multicast/COPMessageBuilder.java | 253 +++++++++++++ .../service/multicast/MulticastSender.java | 331 ++++++++++++++++++ .../multicast/MulticastSenderExample.java | 79 +++++ .../afe/dc/tranlog/task/ScheduledTask.java | 96 +++++ src/main/resources/application.yml | 58 +++ src/main/resources/banner.txt | 24 ++ src/main/resources/bootstrap.yml | 102 ++++++ src/main/resources/logback.xml | 74 ++++ 36 files changed, 4390 insertions(+) create mode 100644 .gitignore create mode 100644 README.md create mode 100644 UDP_MULTICAST_RECEIVER.md create mode 100644 UDP_MULTICAST_SENDER.md create mode 100644 pom.xml create mode 100644 src/main/java/com/afe/dc/tranlog/TranLogServiceApplication.java create mode 100644 src/main/java/com/afe/dc/tranlog/config/AsyncConfig.java create mode 100644 src/main/java/com/afe/dc/tranlog/config/MulticastConfig.java create mode 100644 src/main/java/com/afe/dc/tranlog/config/ScheduleConfig.java create mode 100644 src/main/java/com/afe/dc/tranlog/controller/internal/DataQueryController.java create mode 100644 src/main/java/com/afe/dc/tranlog/controller/internal/DataReceiverController.java create mode 100644 src/main/java/com/afe/dc/tranlog/controller/internal/TradeDayController.java create mode 100644 src/main/java/com/afe/dc/tranlog/domain/dto/AccumulatedData.java create mode 100644 src/main/java/com/afe/dc/tranlog/domain/dto/BuySellStrength.java create mode 100644 src/main/java/com/afe/dc/tranlog/domain/dto/COPItem.java create mode 100644 src/main/java/com/afe/dc/tranlog/domain/dto/TradeDayResponse.java create mode 100644 src/main/java/com/afe/dc/tranlog/domain/model/entity/BsRecord.java create mode 100644 src/main/java/com/afe/dc/tranlog/domain/model/entity/MtcRecord.java create mode 100644 src/main/java/com/afe/dc/tranlog/domain/model/entity/TranRecord.java create mode 100644 src/main/java/com/afe/dc/tranlog/domain/model/entity/TranTable.java create mode 100644 src/main/java/com/afe/dc/tranlog/service/AccumulatedLogicService.java create mode 100644 src/main/java/com/afe/dc/tranlog/service/BuySellStrengthService.java create mode 100644 src/main/java/com/afe/dc/tranlog/service/COPParser.java create mode 100644 src/main/java/com/afe/dc/tranlog/service/MTCLogicService.java create mode 100644 src/main/java/com/afe/dc/tranlog/service/MulticastReceiverService.java create mode 100644 src/main/java/com/afe/dc/tranlog/service/TranDatabaseService.java create mode 100644 src/main/java/com/afe/dc/tranlog/service/TranLogService.java create mode 100644 src/main/java/com/afe/dc/tranlog/service/TranTableService.java create mode 100644 src/main/java/com/afe/dc/tranlog/service/multicast/COPMessageBuilder.java create mode 100644 src/main/java/com/afe/dc/tranlog/service/multicast/MulticastSender.java create mode 100644 src/main/java/com/afe/dc/tranlog/service/multicast/MulticastSenderExample.java create mode 100644 src/main/java/com/afe/dc/tranlog/task/ScheduledTask.java create mode 100644 src/main/resources/application.yml create mode 100644 src/main/resources/banner.txt create mode 100644 src/main/resources/bootstrap.yml create mode 100644 src/main/resources/logback.xml diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..2c0e049 --- /dev/null +++ b/.gitignore @@ -0,0 +1,44 @@ +###################################################################### +# Build Tools + +.gradle +/build/ +!gradle/wrapper/gradle-wrapper.jar + +target/ +!.mvn/wrapper/maven-wrapper.jar + +###################################################################### +# IDE + +### STS ### +.apt_generated +.classpath +.factorypath +.project +.settings +.springBeans + +### IntelliJ IDEA ### +.idea +*.iws +*.iml +*.ipr + +### NetBeans ### +nbproject/private/ +build/* +nbbuild/ +dist/ +nbdist/ +.nb-gradle/ + +###################################################################### +# Others +*.log +*.xml.versionsBackup +*.swp + +!*/build/*.java +!*/build/*.html +!*/build/*.xml diff --git a/README.md b/README.md new file mode 100644 index 0000000..7433f74 --- /dev/null +++ b/README.md @@ -0,0 +1,196 @@ +# DC TranLog Server + +交易日志服务器 - 处理交易数据、MTC(分钟交易数据)、价格成交量等业务逻辑 + +## 项目概述 + +**TranLogServer** 是一个**实时交易日志服务器**,主要用于: +- 接收和处理实时交易数据(股票、期货等金融产品) +- 计算和生成分钟级交易数据(MTC - Minute Transaction Count) +- 提供价格-成交量分析、买卖强度计算等增值服务 +- 支持数据重建、备份、同步等功能 + +本项目是从C++版本转换而来的Java 17 + Spring Cloud Alibaba版本。 + +## 技术栈 + +- **Java 17** +- **Spring Boot** +- **Spring Cloud Alibaba** +- **Lombok** +- **MySQL** (可选,用于数据持久化) + +## 项目结构 + +``` +dc-tranlog-server/ +├── src/main/java/com/afe/dc/tranlog/ +│ ├── TranLogServerApplication.java # 主应用类 +│ ├── controller/ # 控制器层 +│ │ ├── DataReceiverController.java # 数据接收控制器 +│ │ ├── DataQueryController.java # 数据查询控制器 +│ │ └── TradeDayController.java # 交易日管理控制器 +│ ├── service/ # 服务层 +│ │ ├── TranLogService.java # 主服务(对应CTranLogServer) +│ │ ├── TranDatabaseService.java # 数据库服务(对应CTranDatabase) +│ │ ├── TranTableService.java # 交易表服务(对应CTranTable) +│ │ ├── MTCLogicService.java # MTC逻辑服务(对应CMTCLogic) +│ │ ├── AccumulatedLogicService.java # 累计数据服务(对应CAccumulatedLogic) +│ │ └── BuySellStrengthService.java # 买卖强度服务(对应CBuySellStrengthLogic) +│ ├── entity/ # 实体类 +│ │ ├── TranRecord.java # 交易记录 +│ │ ├── TranTable.java # 交易表 +│ │ ├── MtcRecord.java # MTC记录 +│ │ └── BsRecord.java # 买卖强度记录 +│ ├── dto/ # 数据传输对象 +│ │ ├── COPItem.java # COP数据项 +│ │ ├── TradeDayResponse.java # 交易日响应 +│ │ ├── AccumulatedData.java # 累计数据 +│ │ └── BuySellStrength.java # 买卖强度 +│ ├── config/ # 配置类 +│ │ ├── AsyncConfig.java # 异步配置 +│ │ └── ScheduleConfig.java # 定时任务配置 +│ └── task/ # 定时任务 +│ └── ScheduledTask.java # 定时任务类 +└── src/main/resources/ + └── application.yml # 配置文件 +``` + +## 核心功能 + +### 1. 数据接收和处理 +- 接收实时交易数据(通过REST API或消息队列) +- 处理不同类型的消息(UPDATE、DROP、REBUILD等) +- 管理交易日历和交易时间 + +### 2. MTC(分钟交易数据)计算 +- 将实时交易数据聚合成分钟级数据 +- 包含:开盘价、最高价、最低价、收盘价、成交量 + +### 3. 业务逻辑计算 +- **累计数据**:累计成交量和成交额 +- **买卖强度**:根据交易价格和成交量计算买卖强度(高/中/低) +- **价格-成交量分析**:价格-成交量分布 + +### 4. 定时任务 +- MTC发送任务(每50ms) +- 定期重建任务(每30s) +- 价格-成交量发送任务(每50ms) +- 开盘维护任务(每天9:00) +- 收盘维护任务(每天16:00) + +## 使用说明 + +### 启动应用 + +```bash +mvn spring-boot:run +``` + +### 配置说明 + +主要配置在 `application.yml` 中: + +```yaml +tranlog: + trade: + start-hour: 9 + start-minute: 30 + period: 240 # 交易时长(分钟) + mtc: + update-period: 50 # MTC更新周期(毫秒) + rebuild: + enabled: true + schedule: "0 0 2 * * ?" # 每天凌晨2点 +``` + +### API接口 + +#### 1. 接收交易数据 +``` +POST /api/tranlog/process +Content-Type: application/json + +{ + "itemNo": 123456, + "msgType": 1, + "fields": { + "1001": {...} // FID_TRAN_LOG + } +} +``` + +#### 2. 查询MTC记录 +``` +GET /api/data/mtc/{itemNo}?dateIndex=0&startIndex=0&endIndex=100 +``` + +#### 3. 查询累计数据 +``` +GET /api/data/accumulated/{itemNo} +``` + +#### 4. 查询买卖强度 +``` +GET /api/data/buysell/{itemNo}?is15Tick=false +``` + +## 架构说明 + +### C++ 到 Java 的映射关系 + +| C++ 类 | Java 对应 | 说明 | +|--------|----------|------| +| `CTranLogServer` | `TranLogService` + `DataReceiverController` | 主服务器类 | +| `CTranDatabase` | `TranDatabaseService` | 数据访问层 | +| `CTranTable` | `TranTableService` + `TranTable` | 交易表 | +| `CMTCLogic` | `MTCLogicService` | MTC计算逻辑 | +| `CAccumulatedLogic` | `AccumulatedLogicService` | 累计数据逻辑 | +| `CBuySellStrengthLogic` | `BuySellStrengthService` | 买卖强度逻辑 | + +### 数据流 + +``` +外部数据源 + ↓ +DataReceiverController (REST API) + ↓ +TranLogService::process() + ↓ +TranDatabaseService::updateTran() + ↓ +TranTableService::updateTran() + ↓ +业务逻辑更新 (MTCLogicService, AccumulatedLogicService等) +``` + +## 注意事项 + +1. **线程安全**:使用 `ReentrantReadWriteLock` 保证线程安全 +2. **内存管理**:数据主要存储在内存中(ConcurrentHashMap),类似Redis +3. **性能优化**:使用异步处理和定时任务提高性能 +4. **扩展性**:支持水平扩展,可以部署多个实例 + +## 开发说明 + +### 代码转换说明 + +本项目按照 `source/项目快速理解指南-面向Java程序员.md` 中的转换逻辑进行转换: + +1. **C++类 → Java类**:使用Spring的`@Service`、`@RestController`等注解 +2. **线程 → 异步**:使用`@Async`和`ThreadPoolTaskExecutor` +3. **定时器 → 定时任务**:使用`@Scheduled` +4. **多播通信 → REST API/消息队列**:使用Spring MVC或RabbitMQ/Kafka +5. **读写锁**:使用`ReentrantReadWriteLock` + +### 待完善功能 + +1. **消息队列集成**:集成RabbitMQ或Kafka接收实时数据 +2. **数据持久化**:完善MySQL持久化功能 +3. **Redis缓存**:集成Redis缓存热点数据 +4. **监控和日志**:集成Spring Boot Actuator和日志系统 +5. **单元测试**:添加单元测试和集成测试 + +## 许可证 + +Copyright (c) 2025 diff --git a/UDP_MULTICAST_RECEIVER.md b/UDP_MULTICAST_RECEIVER.md new file mode 100644 index 0000000..0764ef6 --- /dev/null +++ b/UDP_MULTICAST_RECEIVER.md @@ -0,0 +1,178 @@ +# UDP多播接收功能说明 + +## 概述 + +本功能实现了UDP多播接收逻辑,对应C++代码`TLogServer.cpp`中的接收广播机制。服务通过UDP多播接收COP(Common Object Protocol)协议数据,解析后交由`TranLogService`处理。 + +## 功能对应关系 + +### C++代码对应关系 + +| C++代码 | Java实现 | 说明 | +|---------|---------|------| +| `m_RecvCtrl.Start(m_iRecvPort, m_IP.c_str())` | `MulticastReceiverService.start()` | 启动UDP多播接收 | +| `m_RecvCtrl.AddGroup(m_GroupIPArray[GIPIndex].c_str(), m_IP.c_str())` | `MulticastReceiverService.start()` 中的 `socket.joinGroup()` | 加入多播组 | +| `m_RecvCtrl.RegisterCallback(CallbackFunc, this)` | `MulticastReceiverService.receiveLoop()` | 注册回调函数 | +| `CTranLogServer::CallbackFunc(COP_ITEM& item)` | `MulticastReceiverService.processReceivedData()` | 回调处理函数 | +| `CTranLogServer::Process(item)` | `TranLogService.process(item)` | 处理数据 | + +## 配置说明 + +在 `application.yml` 中配置UDP多播接收参数: + +```yaml +multicast: + receiver: + # 是否启用UDP多播接收 + enabled: true + # 接收端口(对应C++的RecvPort) + recv-port: 5000 + # 本地网络接口IP地址(对应C++的IPAddress) + ip-address: 0.0.0.0 + # 多播组IP地址列表(对应C++的GroupIP0, GroupIP1, ...) + group-ips: + - 225.6.7.8 + # 可以配置多个多播组 + # - 225.6.7.9 + # 接收缓冲区大小(字节) + buffer-size: 65536 + # 接收超时时间(毫秒) + timeout: 1000 +``` + +### 环境变量配置 + +也可以通过环境变量配置: + +- `MULTICAST_RECEIVER_ENABLED`: 是否启用(默认:true) +- `MULTICAST_RECEIVER_PORT`: 接收端口(默认:5000) +- `MULTICAST_RECEIVER_IP`: 本地IP地址(默认:0.0.0.0) +- `MULTICAST_GROUP_IP_0`: 第一个多播组IP +- `MULTICAST_BUFFER_SIZE`: 缓冲区大小(默认:65536) +- `MULTICAST_TIMEOUT`: 超时时间(默认:1000) + +## 核心组件 + +### 1. MulticastConfig + +配置类,读取UDP多播接收相关配置。 + +**位置**: `com.afe.dc.tranlog.config.MulticastConfig` + +### 2. MulticastReceiverService + +UDP多播接收服务,负责: +- 创建并绑定UDP多播Socket +- 加入配置的多播组 +- 接收UDP数据包 +- 调用COP解析器解析数据 +- 将解析后的数据传递给TranLogService处理 + +**位置**: `com.afe.dc.tranlog.service.MulticastReceiverService` + +**生命周期**: +- `@PostConstruct`: 应用启动时自动启动接收服务 +- `@PreDestroy`: 应用关闭时自动停止接收服务 + +### 3. COPParser + +COP协议解析器,负责将UDP接收到的字节数据解析为`COPItem`对象。 + +**位置**: `com.afe.dc.tranlog.service.COPParser` + +**功能**: +- 解析COP消息头(消息类型、ItemNo等) +- 解析FID字段数据 +- 支持多种数据类型(CHAR、SHORT、INT、LONG、FLOAT、DOUBLE、STRING等) + +**注意**: COP协议的具体格式可能需要根据实际协议规范进行调整。 + +## 数据流程 + +``` +外部数据源(Data Provider) + │ + │ (UDP 多播发送 COP 协议数据) + ▼ +MulticastReceiverService (接收端) + │ + │ (接收UDP数据包) + ▼ +COPParser.parse() (解析COP协议) + │ + │ (转换为COPItem对象) + ▼ +TranLogService.process() (处理数据) + │ + ▼ +TranDatabaseService (更新数据库) +``` + +## 启动和停止 + +### 自动启动 + +服务会在Spring Boot应用启动时自动启动(通过`@PostConstruct`注解)。 + +### 手动控制 + +可以通过配置`multicast.receiver.enabled=false`来禁用UDP多播接收功能。 + +### 停止 + +服务会在Spring Boot应用关闭时自动停止(通过`@PreDestroy`注解),包括: +- 离开所有多播组 +- 关闭Socket +- 停止接收线程 + +## 日志 + +服务会输出以下关键日志: + +- 启动成功:`[MulticastReceiverService] Started UDP multicast receiver on port: {port}, groups: {groups}` +- 加入多播组:`[MulticastReceiverService] Joined multicast group: {ip} on interface: {interface}` +- 接收数据:`[MulticastReceiverService] Processing COP item: itemNo={itemNo}, msgType={msgType}` +- 解析失败:`[COPParser] Failed to parse COP data` +- 处理失败:`[MulticastReceiverService] Failed to process COP item` + +## 注意事项 + +1. **COP协议格式**: 当前实现的COP解析器是基于通用协议的假设。如果实际的COP协议格式不同,需要调整`COPParser`的解析逻辑。 + +2. **网络权限**: 在某些操作系统上,加入多播组可能需要特殊权限。 + +3. **防火墙**: 确保防火墙允许UDP数据包通过配置的端口。 + +4. **网络接口**: 如果配置了特定的IP地址,确保该IP地址对应的网络接口存在且可用。 + +5. **性能**: 接收缓冲区大小和超时时间可以根据实际网络环境调整。 + +## 故障排查 + +### 无法接收数据 + +1. 检查配置是否正确(端口、多播组IP) +2. 检查网络接口是否正确 +3. 检查防火墙设置 +4. 查看日志中的错误信息 + +### 解析失败 + +1. 检查COP协议格式是否与实现一致 +2. 查看日志中的详细错误信息 +3. 可能需要调整`COPParser`的解析逻辑 + +### 处理失败 + +1. 检查`TranLogService`的日志 +2. 检查数据库连接 +3. 检查数据格式是否正确 + +## 扩展 + +如果需要支持更多的COP协议特性,可以: + +1. 扩展`COPParser`以支持更多的数据类型 +2. 添加更多的FID字段解析逻辑 +3. 添加数据验证和错误处理 +4. 添加性能监控和统计 diff --git a/UDP_MULTICAST_SENDER.md b/UDP_MULTICAST_SENDER.md new file mode 100644 index 0000000..d8ecfe3 --- /dev/null +++ b/UDP_MULTICAST_SENDER.md @@ -0,0 +1,216 @@ +# UDP多播发送功能说明 + +## 概述 + +本功能实现了UDP多播发送逻辑,对应C++代码`TLogServer.cpp`中的发送广播机制。作为独立功能,不依赖Spring框架,可以直接使用。 + +## 功能对应关系 + +### C++代码对应关系 + +| C++代码 | Java实现 | 说明 | +|---------|---------|------| +| `m_SendCtrl.Start(m_iSendPort, m_IP.c_str())` | `MulticastSender.initialize()` | 启动UDP多播发送 | +| `m_SendCtrl.AddChannel(m_SendGrpIP.c_str(), m_iSendPort)` | `MulticastSender`构造函数配置 | 配置发送通道 | +| `m_SendCtrl.Send(item)` | `MulticastSender.send(item)` | 发送COP_ITEM数据 | +| `item->GetMessage()` | `COPMessageBuilder.buildMessage(item)` | 构建COP协议消息 | + +## 核心组件 + +### 1. MulticastSender + +UDP多播发送器,负责: +- 创建并配置UDP多播Socket +- 发送COP协议数据 +- 管理发送状态和监听器 + +**位置**: `com.afe.dc.tranlog.service.multicast.MulticastSender` + +**主要方法**: +- `initialize()`: 初始化发送器 +- `send(COPItem item)`: 发送COPItem数据 +- `sendRaw(byte[] data)`: 发送原始字节数据 +- `close()`: 关闭发送器 +- `addListener(SendListener listener)`: 添加发送监听器 + +### 2. COPMessageBuilder + +COP协议消息构建器,负责将`COPItem`对象转换为字节数组(COP协议格式)。 + +**位置**: `com.afe.dc.tranlog.service.multicast.COPMessageBuilder` + +**主要方法**: +- `buildMessage(COPItem item)`: 构建COP协议消息 + +## 使用示例 + +### 基本使用 + +```java +// 1. 创建发送配置 +MulticastSender.Config config = new MulticastSender.Config( + 5001, // 发送端口 + "225.6.7.9" // 多播组IP +); +config.setLocalIp("0.0.0.0") // 本地IP(可选) + .setTtl(1) // TTL值(可选,默认1) + .setLoopbackDisabled(true); // 禁用回环(可选,默认true) + +// 2. 创建发送器 +MulticastSender sender = new MulticastSender(config); + +// 3. 初始化 +if (!sender.initialize()) { + System.err.println("初始化失败"); + return; +} + +try { + // 4. 创建COPItem数据 + COPItem item = COPItem.builder() + .itemNo(12345L) + .msgType(COPItem.MsgType.MGT_UPDATE) + .build(); + + // 添加FID字段 + item.addField(COPItem.FID.FID_TRAN_LOG, "交易数据"); + item.addField(COPItem.FID.FID_ASK, 100.5f); + item.addField(COPItem.FID.FID_BID, 100.3f); + + // 5. 发送数据 + boolean success = sender.send(item); + +} finally { + // 6. 关闭发送器 + sender.close(); +} +``` + +### 使用发送监听器 + +```java +// 添加发送监听器 +sender.addListener(new MulticastSender.SendListener() { + @Override + public void onSendSuccess(COPItem item, int bytesSent) { + System.out.println("发送成功 - itemNo: " + item.getItemNo() + + ", 大小: " + bytesSent + " 字节"); + } + + @Override + public void onSendError(COPItem item, Exception error) { + System.err.println("发送失败 - itemNo: " + item.getItemNo() + + ", 错误: " + error.getMessage()); + } +}); +``` + +### 发送原始字节数据 + +```java +byte[] rawData = new byte[]{0x01, 0x02, 0x03, 0x04}; +sender.sendRaw(rawData); +``` + +## 配置说明 + +### MulticastSender.Config + +| 参数 | 类型 | 必填 | 默认值 | 说明 | +|------|------|------|--------|------| +| sendPort | int | 是 | - | 发送端口 | +| groupIp | String | 是 | - | 多播组IP地址 | +| localIp | String | 否 | "0.0.0.0" | 本地网络接口IP地址 | +| ttl | int | 否 | 1 | TTL值(Time To Live) | +| loopbackDisabled | boolean | 否 | true | 是否禁用回环 | + +### 配置方法 + +```java +MulticastSender.Config config = new MulticastSender.Config(5001, "225.6.7.9") + .setLocalIp("192.168.1.100") // 链式调用设置本地IP + .setTtl(2) // 设置TTL + .setLoopbackDisabled(false); // 允许回环 +``` + +## COP协议格式 + +当前实现的COP协议格式: + +``` +消息头(16字节): + - 消息类型(1字节) + - ItemNo(4字节,小端序) + - 其他头部信息(11字节,当前用0填充) + +FID字段(每个字段): + - FID编号(2字节) + - 数据类型(1字节) + - 数据长度(2字节) + - 数据内容(变长) +``` + +### 支持的数据类型 + +| 数据类型 | 标识 | Java类型 | 大小 | +|---------|------|----------|------| +| CHAR | 0x01 | Byte, Character | 1字节 | +| SHORT | 0x02 | Short | 2字节 | +| INT | 0x04 | Integer | 4字节 | +| LONG | 0x08 | Long | 8字节 | +| FLOAT | 0x10 | Float | 4字节 | +| DOUBLE | 0x20 | Double | 8字节 | +| STRING | 0x40 | String | 变长 | +| BYTE_ARRAY | 0x80 | byte[] | 变长 | + +## 注意事项 + +1. **独立功能**: 本功能不依赖Spring框架,可以在任何Java应用中使用。 + +2. **线程安全**: `MulticastSender`不是线程安全的,如果需要在多线程环境中使用,需要外部同步。 + +3. **资源管理**: 使用完毕后务必调用`close()`方法释放资源。 + +4. **网络权限**: 在某些操作系统上,发送UDP多播数据可能需要特殊权限。 + +5. **防火墙**: 确保防火墙允许UDP数据包通过配置的端口。 + +6. **TTL值**: TTL值决定了数据包可以经过的路由器数量,通常设置为1(本地网络)或更大的值(跨网络)。 + +7. **回环**: 如果`loopbackDisabled`为`true`,发送的数据不会回环到本地接收端,适合生产环境。 + +## 故障排查 + +### 初始化失败 + +1. 检查多播组IP地址是否有效(必须是224.0.0.0到239.255.255.255之间) +2. 检查端口是否被占用 +3. 检查网络接口是否正确 +4. 查看日志中的详细错误信息 + +### 发送失败 + +1. 检查网络连接 +2. 检查防火墙设置 +3. 检查多播组地址和端口是否正确 +4. 查看日志中的详细错误信息 + +### 数据格式问题 + +1. 检查COP协议格式是否与接收端一致 +2. 检查数据类型是否正确 +3. 查看`COPMessageBuilder`的构建逻辑 + +## 扩展 + +如果需要扩展功能,可以: + +1. **自定义协议格式**: 修改`COPMessageBuilder`以支持不同的协议格式 +2. **批量发送**: 添加批量发送方法 +3. **异步发送**: 添加异步发送支持 +4. **发送统计**: 添加发送统计和监控功能 +5. **重试机制**: 添加发送失败重试机制 + +## 完整示例 + +参考 `MulticastSenderExample.java` 文件中的完整示例代码。 diff --git a/pom.xml b/pom.xml new file mode 100644 index 0000000..1567ab7 --- /dev/null +++ b/pom.xml @@ -0,0 +1,87 @@ + + + 4.0.0 + + + + com.afe.dc + dc-parent + 1.0.0-SNAPSHOT + ../dc-parent/pom.xml + + + dc-tranlog-server + DC TranLog Server + 交易日志服务器 - 处理交易数据、MTC(分钟交易数据)、价格成交量等业务逻辑 + + + 4.1.0 + + + + + + + com.afe.dc + dc-common-settting + ${project.version} + + + + + com.afe.dc + dc-common-core + ${project.version} + + + + + + com.afe.dc + dc-data-feign + ${project.version} + + + + + + ${project.artifactId} + + + org.apache.maven.plugins + maven-compiler-plugin + + 17 + 17 + + -parameters + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + org.projectlombok + lombok + + + + + + + repackage + + + + + + + + diff --git a/src/main/java/com/afe/dc/tranlog/TranLogServiceApplication.java b/src/main/java/com/afe/dc/tranlog/TranLogServiceApplication.java new file mode 100644 index 0000000..1d66ded --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/TranLogServiceApplication.java @@ -0,0 +1,58 @@ +package com.afe.dc.tranlog; + +import com.afe.dc.tranlog.service.TranDatabaseService; +import com.baomidou.mybatisplus.annotation.DbType; +import com.baomidou.mybatisplus.extension.plugins.MybatisPlusInterceptor; +import com.baomidou.mybatisplus.extension.plugins.inner.PaginationInnerInterceptor; +import lombok.extern.slf4j.Slf4j; +import org.mybatis.spring.annotation.MapperScan; +import org.springframework.boot.CommandLineRunner; +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.context.annotation.Bean; + +/** + * TranLogServer 主应用类 + * 对应C++的main函数和CTranLogServer::Initialize + */ +@Slf4j +@SpringBootApplication +@MapperScan("com.afe.dc.tranlog.domain.mapper") +public class TranLogServiceApplication { + + public static void main(String[] args) { + log.info("[TranLogServerApplication] Starting TranLogServer..."); + SpringApplication.run(TranLogServiceApplication.class, args); + } + + /** + * 应用启动后初始化 + * 对应C++的Initialize方法 + */ + @Bean + public CommandLineRunner init(TranDatabaseService tranDatabaseService) { + return args -> { + log.info("[TranLogServerApplication] Initializing TranLogServer..."); + + // 初始化数据库 + if (tranDatabaseService.initialize()) { + log.info("[TranLogServerApplication] TranDatabase initialized successfully"); + } else { + log.error("[TranLogServerApplication] Failed to initialize TranDatabase"); + System.exit(1); + } + + log.info("[TranLogServerApplication] TranLogServer started successfully"); + }; + } + + /** + * MyBatis-Plus 分页插件配置 + */ + @Bean + public MybatisPlusInterceptor mybatisPlusInterceptor() { + MybatisPlusInterceptor interceptor = new MybatisPlusInterceptor(); + interceptor.addInnerInterceptor(new PaginationInnerInterceptor(DbType.MYSQL)); + return interceptor; + } +} diff --git a/src/main/java/com/afe/dc/tranlog/config/AsyncConfig.java b/src/main/java/com/afe/dc/tranlog/config/AsyncConfig.java new file mode 100644 index 0000000..43b20ea --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/config/AsyncConfig.java @@ -0,0 +1,53 @@ +package com.afe.dc.tranlog.config; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.annotation.EnableAsync; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; + +import java.util.concurrent.Executor; +import java.util.concurrent.ThreadPoolExecutor; + +/** + * 异步配置类 + * 对应C++的线程池配置 + */ +@Slf4j +@Configuration +@EnableAsync +public class AsyncConfig { + + /** + * 异步任务执行器 + * 用于数据重建、备份等异步操作 + */ + @Bean(name = "tranlogTaskExecutor") + public Executor tranlogTaskExecutor() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(10); + executor.setMaxPoolSize(20); + executor.setQueueCapacity(100); + executor.setThreadNamePrefix("tranlog-"); + executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); + executor.initialize(); + log.info("[AsyncConfig] Initialized tranlogTaskExecutor"); + return executor; + } + + /** + * 重建任务执行器 + */ + @Bean(name = "rebuildTaskExecutor") + public Executor rebuildTaskExecutor() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(5); + executor.setMaxPoolSize(10); + executor.setQueueCapacity(50); + executor.setThreadNamePrefix("rebuild-"); + executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); + executor.initialize(); + log.info("[AsyncConfig] Initialized rebuildTaskExecutor"); + return executor; + } +} diff --git a/src/main/java/com/afe/dc/tranlog/config/MulticastConfig.java b/src/main/java/com/afe/dc/tranlog/config/MulticastConfig.java new file mode 100644 index 0000000..e589105 --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/config/MulticastConfig.java @@ -0,0 +1,51 @@ +package com.afe.dc.tranlog.config; + +import lombok.Data; +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.context.annotation.Configuration; + +import java.util.ArrayList; +import java.util.List; + +/** + * UDP多播接收配置 + * 对应C++的LoadConfig中的配置项 + */ +@Data +@Configuration +@ConfigurationProperties(prefix = "multicast.receiver") +public class MulticastConfig { + + /** + * 接收端口 + * 对应C++的RecvPort + */ + private Integer recvPort; + + /** + * 本地网络接口IP地址 + * 对应C++的IPAddress + */ + private String ipAddress = "0.0.0.0"; + + /** + * 多播组IP地址列表 + * 对应C++的GroupIP0, GroupIP1, ... + */ + private List groupIps = new ArrayList<>(); + + /** + * 是否启用UDP多播接收 + */ + private Boolean enabled = true; + + /** + * 接收缓冲区大小(字节) + */ + private Integer bufferSize = 65536; + + /** + * 接收超时时间(毫秒) + */ + private Integer timeout = 1000; +} diff --git a/src/main/java/com/afe/dc/tranlog/config/ScheduleConfig.java b/src/main/java/com/afe/dc/tranlog/config/ScheduleConfig.java new file mode 100644 index 0000000..1dbdbb3 --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/config/ScheduleConfig.java @@ -0,0 +1,18 @@ +package com.afe.dc.tranlog.config; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.annotation.EnableScheduling; + +/** + * 定时任务配置类 + * 对应C++的定时器配置 + */ +@Slf4j +@Configuration +@EnableScheduling +public class ScheduleConfig { + + // 定时任务配置已启用 + // 具体的定时任务在ScheduledTask类中定义 +} diff --git a/src/main/java/com/afe/dc/tranlog/controller/internal/DataQueryController.java b/src/main/java/com/afe/dc/tranlog/controller/internal/DataQueryController.java new file mode 100644 index 0000000..1e780d9 --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/controller/internal/DataQueryController.java @@ -0,0 +1,111 @@ +package com.afe.dc.tranlog.controller.internal; + +import com.afe.dc.tranlog.domain.dto.AccumulatedData; +import com.afe.dc.tranlog.domain.dto.BuySellStrength; +import com.afe.dc.tranlog.domain.model.entity.MtcRecord; +import com.afe.dc.tranlog.domain.model.entity.TranTable; +import com.afe.dc.tranlog.service.TranDatabaseService; +import com.afe.dc.tranlog.service.TranTableService; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.http.ResponseEntity; +import org.springframework.web.bind.annotation.*; + +import java.util.ArrayList; +import java.util.List; + +/** + * 数据查询控制器 + * 提供交易数据查询接口 + */ +@Slf4j +@RestController +@RequestMapping("/api/data") +@RequiredArgsConstructor +public class DataQueryController { + + private final TranDatabaseService tranDatabaseService; + private final TranTableService tranTableService; + + /** + * 获取MTC记录 + */ + @GetMapping("/mtc/{itemNo}") + public ResponseEntity> getMtcRecords( + @PathVariable Long itemNo, + @RequestParam(required = false, defaultValue = "0") Integer dateIndex, + @RequestParam(required = false) Integer startIndex, + @RequestParam(required = false) Integer endIndex) { + + try { + List records = new ArrayList<>(); + + // 获取MTC记录数量 + int size = tranTableService.getMtcRecord(itemNo, 0, dateIndex) != null ? + tranTableService.getMtcRecord(itemNo, 0, dateIndex).getSeq().intValue() : 0; + + int start = startIndex != null ? startIndex : 0; + int end = endIndex != null ? endIndex : size; + + for (int i = start; i < end && i < size; i++) { + MtcRecord record = tranTableService.getMtcRecord(itemNo, i, dateIndex); + if (record != null) { + records.add(record); + } + } + + return ResponseEntity.ok(records); + } catch (Exception e) { + log.error("[DataQueryController] Error getting MTC records", e); + return ResponseEntity.internalServerError().build(); + } + } + + /** + * 获取累计数据 + */ + @GetMapping("/accumulated/{itemNo}") + public ResponseEntity getAccumulatedData(@PathVariable Long itemNo) { + try { + AccumulatedData data = tranTableService.getAccumulatedData(itemNo); + return ResponseEntity.ok(data); + } catch (Exception e) { + log.error("[DataQueryController] Error getting accumulated data", e); + return ResponseEntity.internalServerError().build(); + } + } + + /** + * 获取买卖强度 + */ + @GetMapping("/buysell/{itemNo}") + public ResponseEntity getBuySellStrength( + @PathVariable Long itemNo, + @RequestParam(required = false, defaultValue = "false") Boolean is15Tick) { + + try { + BuySellStrength strength = tranTableService.getBuySellStrength(itemNo, is15Tick); + return ResponseEntity.ok(strength); + } catch (Exception e) { + log.error("[DataQueryController] Error getting buy sell strength", e); + return ResponseEntity.internalServerError().build(); + } + } + + /** + * 获取交易表信息 + */ + @GetMapping("/table/{itemNo}") + public ResponseEntity getTranTable(@PathVariable Long itemNo) { + try { + TranTable table = tranDatabaseService.getTranTable(itemNo); + if (table == null) { + return ResponseEntity.notFound().build(); + } + return ResponseEntity.ok(table); + } catch (Exception e) { + log.error("[DataQueryController] Error getting tran table", e); + return ResponseEntity.internalServerError().build(); + } + } +} diff --git a/src/main/java/com/afe/dc/tranlog/controller/internal/DataReceiverController.java b/src/main/java/com/afe/dc/tranlog/controller/internal/DataReceiverController.java new file mode 100644 index 0000000..75b506e --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/controller/internal/DataReceiverController.java @@ -0,0 +1,56 @@ +package com.afe.dc.tranlog.controller.internal; + +import com.afe.dc.tranlog.domain.dto.COPItem; +import com.afe.dc.tranlog.service.TranLogService; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.http.ResponseEntity; +import org.springframework.web.bind.annotation.*; + +/** + * 数据接收控制器 + * 对应C++的CallbackFunc回调函数 + * 接收实时交易数据 + */ +@Slf4j +@RestController +@RequestMapping("/api/tranlog") +@RequiredArgsConstructor +public class DataReceiverController { + + private final TranLogService tranLogService; + + /** + * 接收交易数据 + * 对应C++的Process方法 + */ + @PostMapping("/process") + public ResponseEntity processData(@RequestBody COPItem item) { + try { + boolean result = tranLogService.process(item); + return ResponseEntity.ok(result); + } catch (Exception e) { + log.error("[DataReceiverController] Error processing data", e); + return ResponseEntity.internalServerError().body(false); + } + } + + /** + * 批量接收交易数据 + */ + @PostMapping("/process/batch") + public ResponseEntity processBatchData(@RequestBody java.util.List items) { + try { + int successCount = 0; + for (COPItem item : items) { + if (tranLogService.process(item)) { + successCount++; + } + } + return ResponseEntity.ok(successCount); + } catch (Exception e) { + log.error("[DataReceiverController] Error processing batch data", e); + return ResponseEntity.internalServerError().body(0); + } + } +} diff --git a/src/main/java/com/afe/dc/tranlog/controller/internal/TradeDayController.java b/src/main/java/com/afe/dc/tranlog/controller/internal/TradeDayController.java new file mode 100644 index 0000000..7823b1b --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/controller/internal/TradeDayController.java @@ -0,0 +1,62 @@ +package com.afe.dc.tranlog.controller.internal; + +import com.afe.dc.tranlog.domain.dto.TradeDayResponse; +import com.afe.dc.tranlog.service.TranLogService; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.http.ResponseEntity; +import org.springframework.web.bind.annotation.*; + +import java.time.LocalDate; + +/** + * 交易日管理控制器 + */ +@Slf4j +@RestController +@RequestMapping("/api/trade-day") +@RequiredArgsConstructor +public class TradeDayController { + + private final TranLogService tranLogService; + + /** + * 获取交易日 + */ + @GetMapping("/{dayCode}") + public ResponseEntity getTradeDay(@PathVariable Integer dayCode) { + try { + LocalDate tradeDate = tranLogService.getTradeDay(dayCode); + LocalDate nextTradeDate = tranLogService.getNextTradeDay(); + + TradeDayResponse response = TradeDayResponse.builder() + .dayCode(dayCode) + .tradeDate(tradeDate) + .nextTradeDate(nextTradeDate) + .valid(tradeDate != null) + .build(); + + return ResponseEntity.ok(response); + } catch (Exception e) { + log.error("[TradeDayController] Error getting trade day", e); + return ResponseEntity.internalServerError().build(); + } + } + + /** + * 验证交易时间 + */ + @GetMapping("/validate-time") + public ResponseEntity validateTradeTime( + @RequestParam Byte hour, + @RequestParam Byte min) { + + try { + boolean valid = tranLogService.isValidTradeTime(hour, min); + return ResponseEntity.ok(valid); + } catch (Exception e) { + log.error("[TradeDayController] Error validating trade time", e); + return ResponseEntity.internalServerError().body(false); + } + } +} diff --git a/src/main/java/com/afe/dc/tranlog/domain/dto/AccumulatedData.java b/src/main/java/com/afe/dc/tranlog/domain/dto/AccumulatedData.java new file mode 100644 index 0000000..116639c --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/domain/dto/AccumulatedData.java @@ -0,0 +1,31 @@ +package com.afe.dc.tranlog.domain.dto; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.io.Serializable; + +/** + * 累计数据DTO + * 对应C++的CAccumulatedLogic计算结果 + */ +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class AccumulatedData implements Serializable { + + private static final long serialVersionUID = 1L; + + /** + * 累计成交量 + */ + private Double volume; + + /** + * 累计成交额 + */ + private Double turnover; +} diff --git a/src/main/java/com/afe/dc/tranlog/domain/dto/BuySellStrength.java b/src/main/java/com/afe/dc/tranlog/domain/dto/BuySellStrength.java new file mode 100644 index 0000000..7d1fadb --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/domain/dto/BuySellStrength.java @@ -0,0 +1,39 @@ +package com.afe.dc.tranlog.domain.dto; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.io.Serial; +import java.io.Serializable; + +/** + * 买卖强度DTO + * 对应C++的CBuySellStrengthLogic计算结果 + * @author ben.yang + */ +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class BuySellStrength implements Serializable { + + @Serial + private static final long serialVersionUID = 1L; + + /** + * 高强度 + */ + private Float high; + + /** + * 中强度 + */ + private Float mid; + + /** + * 低强度 + */ + private Float low; +} diff --git a/src/main/java/com/afe/dc/tranlog/domain/dto/COPItem.java b/src/main/java/com/afe/dc/tranlog/domain/dto/COPItem.java new file mode 100644 index 0000000..3a67b9c --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/domain/dto/COPItem.java @@ -0,0 +1,94 @@ +package com.afe.dc.tranlog.domain.dto; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.io.Serial; +import java.io.Serializable; +import java.util.HashMap; +import java.util.Map; + +/** + * COP数据项DTO + * 对应C++的COP_ITEM + * 用于数据传输和消息传递 + */ +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class COPItem implements Serializable { + + @Serial + private static final long serialVersionUID = 1L; + + /** + * 交易品种编号 + */ + private Long itemNo; + + /** + * 消息类型 + * MGT_UPDATE = 更新消息 + * MGT_FORCEUPDATE = 强制更新 + * MGT_DROP = 删除消息 + * MGT_INTRADAYREBUILD = 日内重建 + * MGT_CLOSINGRUN = 收盘运行 + */ + private Integer msgType; + + /** + * FID字段数据映射 + * Key: FID编号 + * Value: 字段值 + */ + @Builder.Default + private Map fields = new HashMap<>(); + + /** + * 添加FID字段 + */ + public void addField(Integer fid, Object value) { + if (fields == null) { + fields = new HashMap<>(); + } + fields.put(fid, value); + } + + /** + * 获取FID字段值 + */ + public Object getField(Integer fid) { + return fields != null ? fields.get(fid) : null; + } + + /** + * 消息类型常量 + */ + public static class MsgType { + public static final int MGT_UPDATE = 1; + public static final int MGT_FORCEUPDATE = 2; + public static final int MGT_DROP = 3; + public static final int MGT_INTRADAYREBUILD = 4; + public static final int MGT_CLOSINGRUN = 5; + } + + /** + * FID常量定义 + */ + public static class FID { + public static final int FID_TRAN_LOG = 1001; + public static final int FID_RICNAME = 1002; + public static final int FID_ASK = 1003; + public static final int FID_BID = 1004; + public static final int FID_MTCV_01 = 2001; + public static final int FID_MTCV_06 = 2006; + public static final int FID_PRC_VOL = 3001; + public static final int FID_ACC_VOL1 = 3002; + public static final int FID_ACC_TURN1 = 3003; + public static final int FID_BS_STRENGTH = 3004; + public static final int FID_BS_STRN_15TICK = 3005; + } +} diff --git a/src/main/java/com/afe/dc/tranlog/domain/dto/TradeDayResponse.java b/src/main/java/com/afe/dc/tranlog/domain/dto/TradeDayResponse.java new file mode 100644 index 0000000..4ca7310 --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/domain/dto/TradeDayResponse.java @@ -0,0 +1,44 @@ +package com.afe.dc.tranlog.domain.dto; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.io.Serial; +import java.io.Serializable; +import java.time.LocalDate; + +/** + * 交易日响应DTO + * @author ben.yang + */ +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class TradeDayResponse implements Serializable { + + @Serial + private static final long serialVersionUID = 1L; + + /** + * 交易日代码(索引) + */ + private Integer dayCode; + + /** + * 交易日日期 + */ + private LocalDate tradeDate; + + /** + * 下一个交易日 + */ + private LocalDate nextTradeDate; + + /** + * 是否有效 + */ + private Boolean valid; +} diff --git a/src/main/java/com/afe/dc/tranlog/domain/model/entity/BsRecord.java b/src/main/java/com/afe/dc/tranlog/domain/model/entity/BsRecord.java new file mode 100644 index 0000000..8c00b42 --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/domain/model/entity/BsRecord.java @@ -0,0 +1,39 @@ +package com.afe.dc.tranlog.domain.model.entity; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.io.Serializable; + +/** + * 买卖强度记录实体类 + * 对应C++的BsRecord结构体 + */ +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class BsRecord implements Serializable { + + private static final long serialVersionUID = 1L; + + /** + * 买入累计成交量 + */ + private Double buySumVolume; + + /** + * 卖出累计成交量 + */ + private Double sellSumVolume; + + /** + * 清空记录 + */ + public void empty() { + this.buySumVolume = 0.0; + this.sellSumVolume = 0.0; + } +} diff --git a/src/main/java/com/afe/dc/tranlog/domain/model/entity/MtcRecord.java b/src/main/java/com/afe/dc/tranlog/domain/model/entity/MtcRecord.java new file mode 100644 index 0000000..129f5bc --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/domain/model/entity/MtcRecord.java @@ -0,0 +1,70 @@ +package com.afe.dc.tranlog.domain.model.entity; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.io.Serializable; + +/** + * MTC记录实体类 + * 对应C++的MtcRecord结构体 + * MTC = Minute Transaction Count (分钟交易数据) + */ +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class MtcRecord implements Serializable { + + private static final long serialVersionUID = 1L; + + /** + * 序列号 + */ + private Long seq; + + /** + * 开盘价 + */ + private Float open; + + /** + * 最高价 + */ + private Float high; + + /** + * 最低价 + */ + private Float low; + + /** + * 收盘价 + */ + private Float exit; + + /** + * 成交量 + */ + private Double volume; + + /** + * 记录时间 (HHMM格式,如915表示9:15) + */ + private Short recordTime; + + /** + * 清空记录 + */ + public void empty() { + this.seq = 0L; + this.open = 0.0f; + this.high = 0.0f; + this.low = 0.0f; + this.exit = 0.0f; + this.volume = 0.0; + this.recordTime = null; + } +} diff --git a/src/main/java/com/afe/dc/tranlog/domain/model/entity/TranRecord.java b/src/main/java/com/afe/dc/tranlog/domain/model/entity/TranRecord.java new file mode 100644 index 0000000..2643205 --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/domain/model/entity/TranRecord.java @@ -0,0 +1,87 @@ +package com.afe.dc.tranlog.domain.model.entity; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.io.Serializable; + +/** + * 交易记录实体类 + * 对应C++的CTranRecord + */ +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class TranRecord implements Serializable { + + private static final long serialVersionUID = 1L; + + /** + * 序列号 + */ + private Long seq; + + /** + * 小时 + */ + private Byte hour; + + /** + * 分钟 + */ + private Byte min; + + /** + * 秒 + */ + private Byte sec; + + /** + * 成交量 + */ + private Double volume; + + /** + * 价格 + */ + private Float price; + + /** + * 标志位 (FLAG_ASK=0x1, FLAG_BID=0x2, FLAG_NOMINAL=0x4) + */ + private Byte flag; + + /** + * 新FID标识 (用于处理fid 2700) + */ + private Integer newFid; + + /** + * 是否已初始化 + */ + private Boolean initialized; + + /** + * 判断是否为ASK标志 + */ + public boolean isAsk() { + return flag != null && (flag & 0x1) != 0; + } + + /** + * 判断是否为BID标志 + */ + public boolean isBid() { + return flag != null && (flag & 0x2) != 0; + } + + /** + * 判断是否为NOMINAL标志 + */ + public boolean isNominal() { + return flag != null && (flag & 0x4) != 0; + } +} diff --git a/src/main/java/com/afe/dc/tranlog/domain/model/entity/TranTable.java b/src/main/java/com/afe/dc/tranlog/domain/model/entity/TranTable.java new file mode 100644 index 0000000..5a14a0d --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/domain/model/entity/TranTable.java @@ -0,0 +1,148 @@ +package com.afe.dc.tranlog.domain.model.entity; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.io.Serializable; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.locks.ReentrantReadWriteLock; + +/** + * 交易表实体类 + * 对应C++的CTranTable类 + * 存储单个交易品种的所有交易记录和业务逻辑数据 + */ +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class TranTable implements Serializable { + + private static final long serialVersionUID = 1L; + + /** + * 交易品种编号 + */ + private Long itemNo; + + /** + * RIC名称 + */ + private String ricName; + + /** + * 卖价 + */ + private Float ask; + + /** + * 买价 + */ + private Float bid; + + /** + * 上一个卖价 + */ + private Float lastAsk; + + /** + * 上一个买价 + */ + private Float lastBid; + + /** + * 交易记录列表 + */ + @Builder.Default + private List records = new ArrayList<>(); + + /** + * 是否发送增值数据 + */ + @Builder.Default + private Boolean sendValueAdded = true; + + /** + * 读写锁 - 保护基础数据(价格、名称等) + */ + private transient ReentrantReadWriteLock rtLock = new ReentrantReadWriteLock(); + + /** + * 读写锁 - 保护业务逻辑数据 + */ + private transient ReentrantReadWriteLock lock = new ReentrantReadWriteLock(); + + /** + * 添加交易记录 + */ + public void addRecord(TranRecord record) { + lock.writeLock().lock(); + try { + if (records == null) { + records = new ArrayList<>(); + } + records.add(record); + } finally { + lock.writeLock().unlock(); + } + } + + /** + * 更新价格 + */ + public boolean updatePrice(Float ask, Float bid, int flag) { + rtLock.writeLock().lock(); + try { + if ((flag & 0x1) != 0) { + this.lastAsk = this.ask; + this.ask = ask; + } + if ((flag & 0x2) != 0) { + this.lastBid = this.bid; + this.bid = bid; + } + return true; + } finally { + rtLock.writeLock().unlock(); + } + } + + /** + * 获取卖价(线程安全) + */ + public Float getAsk() { + rtLock.readLock().lock(); + try { + return ask; + } finally { + rtLock.readLock().unlock(); + } + } + + /** + * 获取买价(线程安全) + */ + public Float getBid() { + rtLock.readLock().lock(); + try { + return bid; + } finally { + rtLock.readLock().unlock(); + } + } + + /** + * 获取RIC名称(线程安全) + */ + public String getRicName() { + rtLock.readLock().lock(); + try { + return ricName; + } finally { + rtLock.readLock().unlock(); + } + } +} diff --git a/src/main/java/com/afe/dc/tranlog/service/AccumulatedLogicService.java b/src/main/java/com/afe/dc/tranlog/service/AccumulatedLogicService.java new file mode 100644 index 0000000..86ad34d --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/service/AccumulatedLogicService.java @@ -0,0 +1,111 @@ +package com.afe.dc.tranlog.service; + +import com.afe.dc.tranlog.domain.dto.AccumulatedData; +import com.afe.dc.tranlog.domain.model.entity.TranRecord; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; + +import java.util.concurrent.locks.ReentrantReadWriteLock; + +/** + * 累计数据逻辑服务 + * 对应C++的CAccumulatedLogic类 + * 负责计算累计成交量和成交额 + */ +@Slf4j +@Service +public class AccumulatedLogicService { + + /** + * 累计成交量 + */ + private Double volume = 0.0; + + /** + * 累计成交额 + */ + private Double turnover = 0.0; + + /** + * 非自动撮合交易数量 + */ + private Long nonAutomatch = 0L; + + /** + * 读写锁 + */ + private final ReentrantReadWriteLock lock = new ReentrantReadWriteLock(); + + /** + * 更新累计数据 + * 对应C++的internalUpdate方法 + */ + public void update(TranRecord record) { + lock.writeLock().lock(); + try { + Double recordVolume = record.getVolume() != null ? record.getVolume() : 0.0; + Float recordPrice = record.getPrice() != null ? record.getPrice() : 0.0f; + + // 更新累计成交量 + volume += recordVolume; + + // 更新累计成交额 + turnover += recordVolume * recordPrice; + + // 判断是否为自动撮合交易 + if (!isAutomatched(record.getFlag())) { + nonAutomatch++; + } + + } finally { + lock.writeLock().unlock(); + } + } + + /** + * 获取累计数据 + */ + public AccumulatedData getData() { + lock.readLock().lock(); + try { + return AccumulatedData.builder() + .volume(volume) + .turnover(turnover) + .build(); + } finally { + lock.readLock().unlock(); + } + } + + /** + * 清空成交量 + */ + public void dropVolume() { + lock.writeLock().lock(); + try { + volume = 0.0; + } finally { + lock.writeLock().unlock(); + } + } + + /** + * 清空成交额 + */ + public void dropTurnover() { + lock.writeLock().lock(); + try { + turnover = 0.0; + } finally { + lock.writeLock().unlock(); + } + } + + /** + * 判断是否为自动撮合交易 + */ + private boolean isAutomatched(Byte flag) { + // 根据业务逻辑判断 + return flag != null && (flag == ' ' || flag == 'A' || flag == 'B' || flag == 'Y' || flag == 'U'); + } +} diff --git a/src/main/java/com/afe/dc/tranlog/service/BuySellStrengthService.java b/src/main/java/com/afe/dc/tranlog/service/BuySellStrengthService.java new file mode 100644 index 0000000..548c5bf --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/service/BuySellStrengthService.java @@ -0,0 +1,166 @@ +package com.afe.dc.tranlog.service; + +import com.afe.dc.tranlog.domain.dto.BuySellStrength; +import com.afe.dc.tranlog.domain.model.entity.TranRecord; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; + +import java.util.concurrent.locks.ReentrantReadWriteLock; + +/** + * 买卖强度逻辑服务 + * 对应C++的CBuySellStrengthLogic类 + * 负责计算买卖强度(高/中/低) + */ +@Slf4j +@Service +public class BuySellStrengthService { + + /** + * 是否已更新 + */ + private Boolean updated = false; + + /** + * 买入成交量 + */ + private Double buyVolume = 0.0; + + /** + * 卖出成交量 + */ + private Double sellVolume = 0.0; + + /** + * 中性成交量 + */ + private Double neutralVolume = 0.0; + + /** + * 读写锁 + */ + private final ReentrantReadWriteLock lock = new ReentrantReadWriteLock(); + + /** + * 是否启用 + */ + public static boolean enabled = true; + + /** + * 更新买卖强度数据 + * 对应C++的internalUpdate方法 + */ + public void update(TranRecord record) { + if (!enabled) { + return; + } + + lock.writeLock().lock(); + try { + if (!isInterestedFlag(record)) { + return; + } + + Double volume = record.getVolume() != null ? record.getVolume() : 0.0; + Byte flag = record.getFlag(); + + // 根据flag判断买卖方向 + if (record.isAsk()) { + // 卖盘 + sellVolume += volume; + } else if (record.isBid()) { + // 买盘 + buyVolume += volume; + } else { + // 中性 + neutralVolume += volume; + } + + updated = true; + + } finally { + lock.writeLock().unlock(); + } + } + + /** + * 获取买卖强度数据 + */ + public BuySellStrength getData() { + lock.readLock().lock(); + try { + // 计算买卖强度 + // 这里简化处理,实际应该根据价格-成交量分布来计算 + Float high = calculateHigh(); + Float mid = calculateMid(); + Float low = calculateLow(); + + return BuySellStrength.builder() + .high(high) + .mid(mid) + .low(low) + .build(); + } finally { + lock.readLock().unlock(); + } + } + + /** + * 清空数据 + */ + public void drop() { + lock.writeLock().lock(); + try { + buyVolume = 0.0; + sellVolume = 0.0; + neutralVolume = 0.0; + updated = false; + } finally { + lock.writeLock().unlock(); + } + } + + /** + * 判断是否为感兴趣的标志 + */ + private boolean isInterestedFlag(TranRecord record) { + Byte flag = record.getFlag(); + return flag != null && (flag == ' ' || flag == 'A' || flag == 'B' || flag == 'Y' || flag == 'U'); + } + + /** + * 计算高强度 + */ + private Float calculateHigh() { + // 简化计算,实际应该根据业务逻辑 + double total = buyVolume + sellVolume + neutralVolume; + if (total == 0) { + return 0.0f; + } + return (float) (buyVolume / total * 100); + } + + /** + * 计算中强度 + */ + private Float calculateMid() { + // 简化计算 + double total = buyVolume + sellVolume + neutralVolume; + if (total == 0) { + return 0.0f; + } + return (float) (neutralVolume / total * 100); + } + + /** + * 计算低强度 + */ + private Float calculateLow() { + // 简化计算 + double total = buyVolume + sellVolume + neutralVolume; + if (total == 0) { + return 0.0f; + } + return (float) (sellVolume / total * 100); + } +} diff --git a/src/main/java/com/afe/dc/tranlog/service/COPParser.java b/src/main/java/com/afe/dc/tranlog/service/COPParser.java new file mode 100644 index 0000000..51f85f2 --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/service/COPParser.java @@ -0,0 +1,206 @@ +package com.afe.dc.tranlog.service; + +import com.afe.dc.tranlog.domain.dto.COPItem; +import lombok.extern.slf4j.Slf4j; + +import java.nio.ByteBuffer; +import java.nio.ByteOrder; +import java.nio.charset.StandardCharsets; +import java.util.HashMap; +import java.util.Map; + +/** + * COP协议解析器 + * 解析UDP多播接收到的COP协议数据 + * 对应C++的MsgObject和COP_ITEM解析逻辑 + */ +@Slf4j +public class COPParser { + + /** + * COP消息头大小(字节) + * 通常包含:消息类型(1字节) + ItemNo(4字节) + 其他头部信息 + */ + private static final int COP_HEADER_SIZE = 16; + + /** + * 解析COP协议数据 + * + * @param data UDP接收到的字节数组 + * @return 解析后的COPItem对象,如果解析失败返回null + */ + public static COPItem parse(byte[] data) { + if (data == null || data.length < COP_HEADER_SIZE) { + log.warn("[COPParser] Invalid data length: {}", data != null ? data.length : 0); + return null; + } + + try { + ByteBuffer buffer = ByteBuffer.wrap(data); + buffer.order(ByteOrder.LITTLE_ENDIAN); // COP协议通常使用小端序 + + // 读取消息类型(通常在头部) + int msgType = buffer.get() & 0xFF; + + // 读取ItemNo(交易品种编号,4字节) + long itemNo = buffer.getInt() & 0xFFFFFFFFL; + + // 解析FID字段 + // 注意:实际的COP协议格式可能更复杂,这里提供基础框架 + // 需要根据实际的COP协议规范来完善解析逻辑 + Map fields = parseFidFields(buffer); + + // 创建COPItem对象 + COPItem item = new COPItem(); + item.setItemNo(itemNo); + item.setMsgType(msgType); + item.setFields(fields != null && !fields.isEmpty() ? fields : new HashMap<>()); + + return item; + + } catch (Exception e) { + log.error("[COPParser] Error parsing COP data", e); + return null; + } + } + + /** + * 解析FID字段 + * + * @param buffer 字节缓冲区 + * @return FID字段Map,如果解析失败返回空Map + */ + private static Map parseFidFields(ByteBuffer buffer) { + Map fields = new HashMap<>(); + + try { + // 跳过头部剩余部分 + if (buffer.remaining() < 4) { + return fields; + } + + // 读取FID数量(如果协议中有此字段) + int fidCount = 0; + if (buffer.remaining() >= 2) { + fidCount = buffer.getShort() & 0xFFFF; + } + + // 解析每个FID字段 + // 注意:这里需要根据实际的COP协议格式来实现 + // 通常格式为:FID编号(2字节) + 数据类型(1字节) + 数据长度(2字节) + 数据内容 + while (buffer.remaining() >= 5) { + int fid = buffer.getShort() & 0xFFFF; + byte dataType = buffer.get(); + int dataLength = buffer.getShort() & 0xFFFF; + + if (buffer.remaining() < dataLength) { + log.warn("[COPParser] Insufficient data for FID {}: need {}, remaining {}", + fid, dataLength, buffer.remaining()); + break; + } + + Object value = parseFidValue(buffer, dataType, dataLength); + if (value != null) { + fields.put(fid, value); + } + } + + } catch (Exception e) { + log.error("[COPParser] Error parsing FID fields", e); + } + + return fields; + } + + /** + * 根据数据类型解析FID值 + * + * @param buffer 字节缓冲区 + * @param dataType 数据类型 + * @param length 数据长度 + * @return 解析后的值 + */ + private static Object parseFidValue(ByteBuffer buffer, byte dataType, int length) { + try { + switch (dataType) { + case 0x01: // CHAR + if (length >= 1) { + return (char) buffer.get(); + } + break; + case 0x02: // SHORT + if (length >= 2) { + return buffer.getShort(); + } + break; + case 0x04: // INT + if (length >= 4) { + return buffer.getInt(); + } + break; + case 0x08: // LONG + if (length >= 8) { + return buffer.getLong(); + } + break; + case 0x10: // FLOAT + if (length >= 4) { + return buffer.getFloat(); + } + break; + case 0x20: // DOUBLE + if (length >= 8) { + return buffer.getDouble(); + } + break; + case 0x40: // STRING + if (length > 0) { + byte[] strBytes = new byte[length]; + buffer.get(strBytes); + return new String(strBytes, StandardCharsets.UTF_8).trim(); + } + break; + default: + // 未知类型,读取原始字节 + if (length > 0) { + byte[] rawBytes = new byte[length]; + buffer.get(rawBytes); + return rawBytes; + } + break; + } + } catch (Exception e) { + log.error("[COPParser] Error parsing FID value, type={}, length={}", dataType, length, e); + } + return null; + } + + /** + * 简化版解析:直接从字节数组解析基础信息 + * 如果完整解析失败,可以使用此方法作为fallback + */ + public static COPItem parseSimple(byte[] data) { + if (data == null || data.length < 5) { + return null; + } + + try { + ByteBuffer buffer = ByteBuffer.wrap(data); + buffer.order(ByteOrder.LITTLE_ENDIAN); + + int msgType = buffer.get() & 0xFF; + long itemNo = buffer.getInt() & 0xFFFFFFFFL; + + COPItem item = new COPItem(); + item.setItemNo(itemNo); + item.setMsgType(msgType); + item.setFields(new HashMap<>()); + + return item; + + } catch (Exception e) { + log.error("[COPParser] Error in simple parse", e); + return null; + } + } +} diff --git a/src/main/java/com/afe/dc/tranlog/service/MTCLogicService.java b/src/main/java/com/afe/dc/tranlog/service/MTCLogicService.java new file mode 100644 index 0000000..12104af --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/service/MTCLogicService.java @@ -0,0 +1,244 @@ +package com.afe.dc.tranlog.service; + +import com.afe.dc.tranlog.domain.model.entity.MtcRecord; +import com.afe.dc.tranlog.domain.model.entity.TranRecord; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.locks.ReentrantReadWriteLock; + +/** + * MTC逻辑服务 + * 对应C++的CMTCLogic类 + * 负责计算和生成分钟级交易数据(MTC - Minute Transaction Count) + */ +@Slf4j +@Service +public class MTCLogicService { + + /** + * MTC记录映射 + * Key: DateIndex (日期索引,0=当前交易日,1=下一个交易日) + * Value: MTC记录列表 + */ + private final Map> mtcMap = new ConcurrentHashMap<>(); + + /** + * 参考价格(开盘价) + */ + private Float referencePrice = 0.0f; + + /** + * 参考价格是否已更新 + */ + private Boolean referencePriceUpdated = false; + + /** + * 当前Tick时间索引 + */ + private Integer currentTickTimeIndex = -1; + + /** + * 交易开始时间 (HHMM格式) + */ + private Short startTime; + + /** + * 交易结束时间 (HHMM格式) + */ + private Short endTime; + + /** + * MTC开始时间 (HHMM格式) + */ + private Short mtcStartTime; + + /** + * MTC结束时间 (HHMM格式) + */ + private Short mtcEndTime; + + /** + * 交易品种编号 + */ + private Long itemNo; + + /** + * 最后序列号 + */ + private Long lastSeqNum = 0L; + + /** + * 读写锁 + */ + private final ReentrantReadWriteLock lock = new ReentrantReadWriteLock(); + + /** + * 初始化MTC逻辑 + */ + public void initialize(Long itemNo, Short startTime, Short endTime, + Short mtcStartTime, Short mtcEndTime, int totalMinutes) { + this.itemNo = itemNo; + this.startTime = startTime; + this.endTime = endTime; + this.mtcStartTime = mtcStartTime; + this.mtcEndTime = mtcEndTime; + + // 初始化两个日期的MTC记录列表 + List currentDay = new ArrayList<>(totalMinutes); + for (int i = 0; i < totalMinutes; i++) { + currentDay.add(MtcRecord.builder() + .seq(0L) + .open(0.0f) + .high(0.0f) + .low(0.0f) + .exit(0.0f) + .volume(0.0) + .build()); + } + mtcMap.put(0, currentDay); + + List nextDay = new ArrayList<>(totalMinutes); + for (int i = 0; i < totalMinutes; i++) { + nextDay.add(MtcRecord.builder() + .seq(0L) + .open(0.0f) + .high(0.0f) + .low(0.0f) + .exit(0.0f) + .volume(0.0) + .build()); + } + mtcMap.put(1, nextDay); + + log.info("[MTCLogicService] Initialized for itemNo={}, totalMinutes={}", itemNo, totalMinutes); + } + + /** + * 更新MTC记录 + * 对应C++的internalUpdate方法 + */ + public void update(TranRecord record, int dateIndex) { + lock.writeLock().lock(); + try { + // 检查序列号连续性 + if (lastSeqNum > 0 && lastSeqNum + 1 != record.getSeq()) { + log.warn("Seq Incorrect: itemNo={}, prev={}, this={}", + itemNo, lastSeqNum, record.getSeq()); + } + lastSeqNum = record.getSeq(); + + // 判断是否为自动撮合交易(根据flag判断) + if (!isAutomatched(record.getFlag())) { + return; + } + + // 计算时间索引 + int recordTickTime = record.getHour() * 100 + record.getMin(); + int index = time2Index((short) recordTickTime, startTime); + + List mtcList = mtcMap.get(dateIndex); + if (index < 0 || index >= mtcList.size()) { + log.warn("Index out of range: index={}, size={}", index, mtcList.size()); + return; + } + + MtcRecord mtc = mtcList.get(index); + Float curPrice = record.getPrice(); + Double curVolume = record.getVolume(); + + // 更新参考价格(如果是第一条记录) + if (dateIndex == 0 && referencePrice == 0.0f) { + updateReferencePrice(curPrice, false); + } + + // 更新MTC记录 + if (mtc.getOpen() == 0.0f) { + mtc.setOpen(curPrice); + } + mtc.setHigh(Math.max(mtc.getHigh() != null ? mtc.getHigh() : curPrice, curPrice)); + mtc.setLow(Math.min(mtc.getLow() != null && mtc.getLow() > 0 ? mtc.getLow() : curPrice, curPrice)); + mtc.setExit(curPrice); + mtc.setVolume((mtc.getVolume() != null ? mtc.getVolume() : 0.0) + curVolume); + mtc.setSeq(record.getSeq()); + + currentTickTimeIndex = index; + + } finally { + lock.writeLock().unlock(); + } + } + + /** + * 更新参考价格 + */ + public void updateReferencePrice(Float refPrice, Boolean prvClUpdated) { + lock.writeLock().lock(); + try { + if (referencePrice == 0.0f || prvClUpdated) { + referencePrice = refPrice; + referencePriceUpdated = true; + } + } finally { + lock.writeLock().unlock(); + } + } + + /** + * 获取MTC记录 + */ + public MtcRecord getMtcRecord(int index, int dateIndex) { + lock.readLock().lock(); + try { + List mtcList = mtcMap.get(dateIndex); + if (mtcList == null || index < 0 || index >= mtcList.size()) { + return null; + } + return mtcList.get(index); + } finally { + lock.readLock().unlock(); + } + } + + /** + * 获取MTC记录数量 + */ + public int getSize(int dateIndex) { + lock.readLock().lock(); + try { + List mtcList = mtcMap.get(dateIndex); + return mtcList != null ? mtcList.size() : 0; + } finally { + lock.readLock().unlock(); + } + } + + /** + * 时间转索引 + */ + private int time2Index(short time, short startTime) { + int startMinutes = (startTime / 100) * 60 + (startTime % 100); + int timeMinutes = (time / 100) * 60 + (time % 100); + return timeMinutes - startMinutes; + } + + /** + * 判断是否为自动撮合交易 + */ + private boolean isAutomatched(Long itemNo, Byte flag) { + // 根据业务逻辑判断,这里简化处理 + // 实际应该根据flag的值来判断 + return flag != null && (flag == ' ' || flag == 'A' || flag == 'B' || flag == 'Y' || flag == 'U'); + } + + /** + * 判断是否为自动撮合交易(重载方法) + */ + private boolean isAutomatched(Byte flag) { + return isAutomatched(itemNo, flag); + } +} diff --git a/src/main/java/com/afe/dc/tranlog/service/MulticastReceiverService.java b/src/main/java/com/afe/dc/tranlog/service/MulticastReceiverService.java new file mode 100644 index 0000000..9fefc39 --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/service/MulticastReceiverService.java @@ -0,0 +1,272 @@ +package com.afe.dc.tranlog.service; + +import com.afe.dc.tranlog.config.MulticastConfig; +import com.afe.dc.tranlog.domain.dto.COPItem; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Service; + +import jakarta.annotation.PostConstruct; +import jakarta.annotation.PreDestroy; +import java.io.IOException; +import java.net.*; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.atomic.AtomicBoolean; + +/** + * UDP多播接收服务 + * 对应C++的CDataCtrl::m_RecvCtrl + * 负责接收UDP多播数据并解析COP协议 + * @author ben.yang + */ +@Slf4j +@Service +@RequiredArgsConstructor +@ConditionalOnProperty(prefix = "multicast.receiver", name = "enabled", havingValue = "true", matchIfMissing = true) +public class MulticastReceiverService { + + private final MulticastConfig config; + private final TranLogService tranLogService; + + private MulticastSocket socket; + private final AtomicBoolean running = new AtomicBoolean(false); + private ExecutorService executorService; + private Thread receiverThread; + + /** + * 初始化并启动UDP多播接收 + * 对应C++的m_RecvCtrl.Start()和AddGroup() + */ + @PostConstruct + public void start() { + if (!config.getEnabled()) { + log.info("[MulticastReceiverService] UDP multicast receiver is disabled"); + return; + } + + if (config.getRecvPort() == null || config.getRecvPort() <= 0) { + log.error("[MulticastReceiverService] Invalid recvPort: {}", config.getRecvPort()); + return; + } + + if (config.getGroupIps() == null || config.getGroupIps().isEmpty()) { + log.error("[MulticastReceiverService] No multicast group IPs configured"); + return; + } + + try { + // 创建多播Socket + socket = new MulticastSocket(config.getRecvPort()); + socket.setSoTimeout(config.getTimeout()); + socket.setReceiveBufferSize(config.getBufferSize()); + + // 绑定本地网络接口 + if (config.getIpAddress() != null && !"0.0.0.0".equals(config.getIpAddress())) { + try { + InetAddress localInterface = InetAddress.getByName(config.getIpAddress()); + socket.setInterface(localInterface); + log.info("[MulticastReceiverService] Set local interface to: {}", config.getIpAddress()); + } catch (UnknownHostException e) { + log.warn("[MulticastReceiverService] Failed to set local interface: {}", config.getIpAddress(), e); + } + } + + // 加入多播组 + NetworkInterface networkInterface = getNetworkInterface(); + for (String groupIp : config.getGroupIps()) { + try { + InetAddress group = InetAddress.getByName(groupIp); + if (networkInterface != null) { + // Java 7+ 使用新的API + InetSocketAddress groupAddress = new InetSocketAddress(group, config.getRecvPort()); + socket.joinGroup(groupAddress, networkInterface); + log.info("[MulticastReceiverService] Joined multicast group: {} on interface: {}", + groupIp, networkInterface.getName()); + } else { + // 使用默认接口 + socket.joinGroup(group); + log.info("[MulticastReceiverService] Joined multicast group: {} on default interface", groupIp); + } + } catch (IOException e) { + log.error("[MulticastReceiverService] Failed to join multicast group: {}", groupIp, e); + } + } + + // 启动接收线程 + running.set(true); + executorService = Executors.newSingleThreadExecutor(r -> { + Thread t = new Thread(r, "MulticastReceiver"); + t.setDaemon(false); + return t; + }); + + receiverThread = new Thread(this::receiveLoop, "MulticastReceiver-Thread"); + receiverThread.setDaemon(false); + receiverThread.start(); + + log.info("[MulticastReceiverService] Started UDP multicast receiver on port: {}, groups: {}", + config.getRecvPort(), config.getGroupIps()); + + } catch (IOException e) { + log.error("[MulticastReceiverService] Failed to start UDP multicast receiver", e); + running.set(false); + } + } + + /** + * 停止UDP多播接收 + * 对应C++的m_RecvCtrl.Stop()和DelGroup() + */ + @PreDestroy + public void stop() { + log.info("[MulticastReceiverService] Stopping UDP multicast receiver..."); + + running.set(false); + + // 离开多播组 + if (socket != null && !socket.isClosed()) { + NetworkInterface networkInterface = getNetworkInterface(); + for (String groupIp : config.getGroupIps()) { + try { + InetAddress group = InetAddress.getByName(groupIp); + if (networkInterface != null) { + // Java 7+ 使用新的API + InetSocketAddress groupAddress = new InetSocketAddress(group, config.getRecvPort()); + socket.leaveGroup(groupAddress, networkInterface); + } else { + // 使用默认接口 + socket.leaveGroup(group); + } + log.info("[MulticastReceiverService] Left multicast group: {}", groupIp); + } catch (IOException e) { + log.warn("[MulticastReceiverService] Failed to leave multicast group: {}", groupIp, e); + } + } + + socket.close(); + } + + // 等待接收线程结束 + if (receiverThread != null && receiverThread.isAlive()) { + try { + receiverThread.join(5000); + } catch (InterruptedException e) { + log.warn("[MulticastReceiverService] Interrupted while waiting for receiver thread", e); + Thread.currentThread().interrupt(); + } + } + + // 关闭线程池 + if (executorService != null) { + executorService.shutdown(); + } + + log.info("[MulticastReceiverService] UDP multicast receiver stopped"); + } + + /** + * 接收循环 + * 对应C++的回调函数CallbackFunc + */ + private void receiveLoop() { + byte[] buffer = new byte[config.getBufferSize()]; + DatagramPacket packet = new DatagramPacket(buffer, buffer.length); + + log.info("[MulticastReceiverService] Receiver loop started"); + + while (running.get() && socket != null && !socket.isClosed()) { + try { + // 接收UDP数据包 + socket.receive(packet); + + // 处理接收到的数据 + processReceivedData(packet.getData(), packet.getLength(), packet.getAddress()); + + // 重置数据包长度 + packet.setLength(buffer.length); + + } catch (SocketTimeoutException e) { + // 超时是正常的,继续接收 + continue; + } catch (IOException e) { + if (running.get()) { + log.error("[MulticastReceiverService] Error receiving UDP packet", e); + } + } + } + + log.info("[MulticastReceiverService] Receiver loop stopped"); + } + + /** + * 处理接收到的数据 + * 对应C++的Process方法 + * + * @param data 接收到的数据 + * @param length 数据长度 + * @param sourceAddress 数据源地址 + */ + private void processReceivedData(byte[] data, int length, InetAddress sourceAddress) { + if (data == null || length <= 0) { + return; + } + + try { + // 截取实际数据长度 + byte[] actualData = new byte[length]; + System.arraycopy(data, 0, actualData, 0, length); + + // 解析COP协议数据 + COPItem item = COPParser.parse(actualData); + + if (item == null) { + // 如果完整解析失败,尝试简单解析 + item = COPParser.parseSimple(actualData); + } + + if (item != null) { + // 调用处理逻辑 + boolean success = tranLogService.process(item); + + if (!success) { + log.warn("[MulticastReceiverService] Failed to process COP item: itemNo={}, msgType={}", + item.getItemNo(), item.getMsgType()); + } + } else { + log.warn("[MulticastReceiverService] Failed to parse COP data from: {}, length: {}", + sourceAddress, length); + } + + } catch (Exception e) { + log.error("[MulticastReceiverService] Error processing received data from: {}", sourceAddress, e); + } + } + + /** + * 获取网络接口 + * 用于加入多播组 + */ + private NetworkInterface getNetworkInterface() { + if (config.getIpAddress() == null || config.getIpAddress().equals("0.0.0.0")) { + return null; // 使用默认接口 + } + + try { + InetAddress localAddress = InetAddress.getByName(config.getIpAddress()); + return NetworkInterface.getByInetAddress(localAddress); + } catch (Exception e) { + log.warn("[MulticastReceiverService] Failed to get network interface for: {}", + config.getIpAddress(), e); + return null; + } + } + + /** + * 检查服务是否运行 + */ + public boolean isRunning() { + return running.get() && socket != null && !socket.isClosed(); + } +} diff --git a/src/main/java/com/afe/dc/tranlog/service/TranDatabaseService.java b/src/main/java/com/afe/dc/tranlog/service/TranDatabaseService.java new file mode 100644 index 0000000..fa4193d --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/service/TranDatabaseService.java @@ -0,0 +1,193 @@ +package com.afe.dc.tranlog.service; + +import com.afe.dc.tranlog.domain.model.entity.TranRecord; +import com.afe.dc.tranlog.domain.model.entity.TranTable; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; + +import java.util.List; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.locks.ReentrantReadWriteLock; +import java.util.stream.Collectors; + +/** + * 交易数据库服务 + * 对应C++的CTranDatabase类 + * 管理所有交易表的集合 + */ +@Slf4j +@Service +@RequiredArgsConstructor +public class TranDatabaseService { + + private final TranTableService tranTableService; + + /** + * 交易表映射 + * Key: itemNo (交易品种编号) + * Value: TranTable (交易表) + */ + private final ConcurrentHashMap tableMap = new ConcurrentHashMap<>(); + + /** + * 读写锁 + */ + private final ReentrantReadWriteLock lock = new ReentrantReadWriteLock(); + + /** + * 初始化数据库 + */ + public boolean initialize() { + lock.writeLock().lock(); + try { + log.info("[TranDatabaseService] Initialized"); + return true; + } finally { + lock.writeLock().unlock(); + } + } + + /** + * 更新交易记录 + * 对应C++的UpdateTran方法 + */ + public boolean updateTran(Long itemNo, TranRecord tranRec) { + TranTable table = forceGetTranTable(itemNo); + if (table != null) { + return tranTableService.updateTran(itemNo, tranRec); + } + return false; + } + + /** + * 删除交易记录 + */ + public boolean dropTran(Long itemNo, TranRecord tranRec) { + // 简化处理,实际应该从表中删除指定记录 + return true; + } + + /** + * 更新价格 + */ + public boolean updatePrice(Long itemNo, Float ask, Float bid, int flag) { + return tranTableService.updatePrice(itemNo, ask, bid, flag); + } + + /** + * 更新参考价格 + */ + public boolean updateRefPrice(Long itemNo, Float open) { + return tranTableService.updateRefPrice(itemNo, open); + } + + /** + * 更新RIC名称 + */ + public boolean updateRicName(Long itemNo, String ricName) { + return tranTableService.updateRicName(itemNo, ricName); + } + + /** + * 删除交易品种 + */ + public boolean dropItem(Long itemNo) { + lock.writeLock().lock(); + try { + tableMap.remove(itemNo); + tranTableService.dropItem(itemNo); + log.info("[TranDatabaseService] Dropped itemNo={}", itemNo); + return true; + } finally { + lock.writeLock().unlock(); + } + } + + /** + * 强制获取交易表(如果不存在则创建) + */ + public TranTable forceGetTranTable(Long itemNo) { + lock.readLock().lock(); + try { + TranTable table = tableMap.get(itemNo); + if (table == null) { + lock.readLock().unlock(); + lock.writeLock().lock(); + try { + table = tableMap.get(itemNo); + if (table == null) { + // 创建新表 + tranTableService.initialize(itemNo); + table = tranTableService.getTranTable(itemNo); + if (table != null) { + tableMap.put(itemNo, table); + } + } + } finally { + lock.writeLock().unlock(); + lock.readLock().lock(); + } + } + return table; + } finally { + lock.readLock().unlock(); + } + } + + /** + * 获取交易表 + */ + public TranTable getTranTable(Long itemNo) { + lock.readLock().lock(); + try { + return tableMap.get(itemNo); + } finally { + lock.readLock().unlock(); + } + } + + /** + * 获取第一个交易品种编号 + */ + public Long getFirstItem() { + lock.readLock().lock(); + try { + return tableMap.keySet().stream().findFirst().orElse(null); + } finally { + lock.readLock().unlock(); + } + } + + /** + * 获取下一个交易品种编号 + */ + public Long getNextItem(Long lastItem) { + lock.readLock().lock(); + try { + List sortedItems = tableMap.keySet().stream() + .sorted() + .collect(Collectors.toList()); + + int index = sortedItems.indexOf(lastItem); + if (index >= 0 && index < sortedItems.size() - 1) { + return sortedItems.get(index + 1); + } + return null; + } finally { + lock.readLock().unlock(); + } + } + + /** + * 获取数据库大小 + */ + public long getSize() { + lock.readLock().lock(); + try { + return tableMap.size(); + } finally { + lock.readLock().unlock(); + } + } +} diff --git a/src/main/java/com/afe/dc/tranlog/service/TranLogService.java b/src/main/java/com/afe/dc/tranlog/service/TranLogService.java new file mode 100644 index 0000000..1a899cb --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/service/TranLogService.java @@ -0,0 +1,262 @@ +package com.afe.dc.tranlog.service; + +import com.afe.dc.tranlog.domain.dto.COPItem; +import com.afe.dc.tranlog.domain.model.entity.TranRecord; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; + +import java.time.LocalDate; +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.locks.ReentrantReadWriteLock; + +/** + * 交易日志服务(主服务) + * 对应C++的CTranLogServer类 + * 负责接收和处理实时交易数据 + * @author ben.yang + */ +@Slf4j +@Service +@RequiredArgsConstructor +public class TranLogService { + + private final TranDatabaseService tranDatabaseService; + + /** + * 交易日索引 + */ + private final List dayIndex = new ArrayList<>(); + + /** + * 下一个交易日 + */ + private LocalDate nextTradeDate; + + /** + * 交易日锁 + */ + private final ReentrantReadWriteLock tradeDayLock = new ReentrantReadWriteLock(); + + /** + * 交易时间配置 + */ + private Byte tradeHour = 9; + private Byte tradeMin = 30; + private Integer tradePeriod = 240; // 交易时长(分钟) + + /** + * 处理数据项 + * 对应C++的Process方法 + */ + public boolean process(COPItem item) { + if (item == null) { + return false; + } + + Integer msgType = item.getMsgType(); + if (msgType == null) { + return false; + } + + switch (msgType) { + case COPItem.MsgType.MGT_UPDATE: + case COPItem.MsgType.MGT_FORCEUPDATE: + return processNonIntraDay(item); + case COPItem.MsgType.MGT_DROP: + return processDrop(item); + case COPItem.MsgType.MGT_CLOSINGRUN: + return closingRun(item); + default: + log.warn("[TranLogService] Unknown message type: {}", msgType); + return false; + } + } + + /** + * 处理非日内数据 + * 对应C++的ProcessNonIntraDay方法 + */ + private boolean processNonIntraDay(COPItem item) { + Long itemNo = item.getItemNo(); + if (itemNo == null) { + return false; + } + + // 处理交易记录 + if (item.getField(COPItem.FID.FID_TRAN_LOG) != null) { + return processTran(item); + } + + // 处理RIC名称 + if (item.getField(COPItem.FID.FID_RICNAME) != null) { + return processRicName(item); + } + + // 处理参考价格 + if (item.getField(COPItem.FID.FID_ASK) != null || item.getField(COPItem.FID.FID_BID) != null) { + return processAskBid(item); + } + + return true; + } + + /** + * 处理交易记录 + * 对应C++的ProcessTran方法 + */ + private boolean processTran(COPItem item) { + try { + Object tranData = item.getField(COPItem.FID.FID_TRAN_LOG); + if (tranData == null) { + return false; + } + + // 解析交易记录(简化处理,实际应该根据COP格式解析) + TranRecord record = parseTranRecord(tranData); + if (record == null) { + return false; + } + + // 验证交易时间 + if (!isValidTradeTime(record.getHour(), record.getMin())) { + return false; + } + + // 更新交易记录 + return tranDatabaseService.updateTran(item.getItemNo(), record); + + } catch (Exception e) { + log.error("[TranLogService] Error processing tran record", e); + return false; + } + } + + /** + * 处理RIC名称 + */ + private boolean processRicName(COPItem item) { + try { + Object ricNameData = item.getField(COPItem.FID.FID_RICNAME); + if (ricNameData == null) { + return false; + } + + String ricName = ricNameData.toString(); + return tranDatabaseService.updateRicName(item.getItemNo(), ricName); + + } catch (Exception e) { + log.error("[TranLogService] Error processing ric name", e); + return false; + } + } + + /** + * 处理买卖价 + */ + private boolean processAskBid(COPItem item) { + try { + Object askData = item.getField(COPItem.FID.FID_ASK); + Object bidData = item.getField(COPItem.FID.FID_BID); + + Float ask = askData != null ? Float.parseFloat(askData.toString()) : null; + Float bid = bidData != null ? Float.parseFloat(bidData.toString()) : null; + + int flag = 0; + if (ask != null) { + flag |= 0x1; // FLAG_ASK + } + if (bid != null) { + flag |= 0x2; // FLAG_BID + } + + return tranDatabaseService.updatePrice(item.getItemNo(), ask, bid, flag); + + } catch (Exception e) { + log.error("[TranLogService] Error processing ask/bid", e); + return false; + } + } + + public static void main(String[] args) { + int flag = 0; + System.out.println(flag |= 0x1); + + } + + /** + * 处理删除 + */ + private boolean processDrop(COPItem item) { + return tranDatabaseService.dropItem(item.getItemNo()); + } + + /** + * 收盘运行 + */ + private boolean closingRun(COPItem item) { + log.info("[TranLogService] Closing run for itemNo={}", item.getItemNo()); + // 执行收盘相关操作 + return true; + } + + /** + * 获取交易日 + */ + public LocalDate getTradeDay(int dayCode) { + tradeDayLock.readLock().lock(); + try { + if (dayCode >= 0 && dayCode < dayIndex.size()) { + return dayIndex.get(dayCode); + } + return null; + } finally { + tradeDayLock.readLock().unlock(); + } + } + + /** + * 获取下一个交易日 + */ + public LocalDate getNextTradeDay() { + tradeDayLock.readLock().lock(); + try { + return nextTradeDate; + } finally { + tradeDayLock.readLock().unlock(); + } + } + + /** + * 验证是否为有效交易时间 + */ + public boolean isValidTradeTime(Byte hour, Byte min) { + if (hour == null || min == null) { + return false; + } + + int cur = hour * 60 + min; + int start = tradeHour * 60 + tradeMin; + return (cur >= start) && (cur <= start + tradePeriod); + } + + /** + * 解析交易记录(简化实现) + */ + private TranRecord parseTranRecord(Object data) { + // 实际应该根据COP格式解析 + // 这里简化处理 + try { + // 假设data是Map或其他格式 + // 实际应该根据C++代码中的解析逻辑来实现 + return TranRecord.builder() + .initialized(true) + .build(); + } catch (Exception e) { + log.error("[TranLogService] Error parsing tran record", e); + return null; + } + } +} diff --git a/src/main/java/com/afe/dc/tranlog/service/TranTableService.java b/src/main/java/com/afe/dc/tranlog/service/TranTableService.java new file mode 100644 index 0000000..8b471e0 --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/service/TranTableService.java @@ -0,0 +1,237 @@ +package com.afe.dc.tranlog.service; + +import com.afe.dc.tranlog.domain.dto.AccumulatedData; +import com.afe.dc.tranlog.domain.dto.BuySellStrength; +import com.afe.dc.tranlog.domain.model.entity.MtcRecord; +import com.afe.dc.tranlog.domain.model.entity.TranRecord; +import com.afe.dc.tranlog.domain.model.entity.TranTable; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; + +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.locks.ReentrantReadWriteLock; + +/** + * 交易表服务 + * 对应C++的CTranTable类 + * 管理单个交易品种的数据和业务逻辑 + */ +@Slf4j +@Service +public class TranTableService { + + // 注意:MTCLogicService、AccumulatedLogicService等应该是每个TranTable实例独立的 + // 这里简化处理,使用单例模式 + // 实际应该为每个TranTable创建独立的Logic实例 + + /** + * 交易表缓存 + */ + private final ConcurrentHashMap tableCache = new ConcurrentHashMap<>(); + + /** + * Logic服务映射(每个itemNo对应一个Logic实例) + * 实际应该为每个TranTable创建独立的Logic实例 + */ + private final ConcurrentHashMap mtcLogicMap = new ConcurrentHashMap<>(); + private final ConcurrentHashMap accumulatedLogicMap = new ConcurrentHashMap<>(); + private final ConcurrentHashMap buySellStrengthMap = new ConcurrentHashMap<>(); + + /** + * 读写锁 + */ + private final ReentrantReadWriteLock lock = new ReentrantReadWriteLock(); + + /** + * 初始化交易表 + */ + public boolean initialize(Long itemNo) { + lock.writeLock().lock(); + try { + TranTable table = TranTable.builder() + .itemNo(itemNo) + .ask(-1.0f) + .bid(-1.0f) + .lastAsk(-1.0f) + .lastBid(-1.0f) + .sendValueAdded(true) + .build(); + + tableCache.put(itemNo, table); + + // 为每个TranTable创建独立的Logic实例 + MTCLogicService mtcLogic = new MTCLogicService(); + mtcLogic.initialize(itemNo, (short) 930, (short) 1600, (short) 930, (short) 1600, 240); + mtcLogicMap.put(itemNo, mtcLogic); + + AccumulatedLogicService accLogic = new AccumulatedLogicService(); + accumulatedLogicMap.put(itemNo, accLogic); + + BuySellStrengthService bssLogic = new BuySellStrengthService(); + buySellStrengthMap.put(itemNo, bssLogic); + + log.info("[TranTableService] Initialized table for itemNo={}", itemNo); + return true; + } finally { + lock.writeLock().unlock(); + } + } + + /** + * 更新交易记录 + */ + public boolean updateTran(Long itemNo, TranRecord record) { + TranTable table = getTranTable(itemNo); + if (table == null) { + log.warn("[TranTableService] Table not found for itemNo={}", itemNo); + return false; + } + + // 添加交易记录 + table.addRecord(record); + + // 触发业务逻辑更新 + postBusinessLogic(table, record); + + return true; + } + + /** + * 更新价格 + */ + public boolean updatePrice(Long itemNo, Float ask, Float bid, int flag) { + TranTable table = getTranTable(itemNo); + if (table == null) { + return false; + } + return table.updatePrice(ask, bid, flag); + } + + /** + * 更新参考价格(开盘价) + */ + public boolean updateRefPrice(Long itemNo, Float open) { + TranTable table = getTranTable(itemNo); + if (table == null) { + return false; + } + MTCLogicService mtcLogic = mtcLogicMap.get(itemNo); + if (mtcLogic != null) { + mtcLogic.updateReferencePrice(open, false); + } + return true; + } + + /** + * 更新RIC名称 + */ + public boolean updateRicName(Long itemNo, String ricName) { + TranTable table = getTranTable(itemNo); + if (table == null) { + return false; + } + table.setRicName(ricName); + return true; + } + + /** + * 获取交易表 + */ + public TranTable getTranTable(Long itemNo) { + lock.readLock().lock(); + try { + TranTable table = tableCache.get(itemNo); + if (table == null) { + // 如果不存在,自动创建 + lock.readLock().unlock(); + lock.writeLock().lock(); + try { + table = tableCache.get(itemNo); + if (table == null) { + initialize(itemNo); + table = tableCache.get(itemNo); + } + } finally { + lock.writeLock().unlock(); + lock.readLock().lock(); + } + } + return table; + } finally { + lock.readLock().unlock(); + } + } + + /** + * 获取MTC记录 + */ + public MtcRecord getMtcRecord(Long itemNo, int index, int dateIndex) { + MTCLogicService mtcLogic = mtcLogicMap.get(itemNo); + if (mtcLogic != null) { + return mtcLogic.getMtcRecord(index, dateIndex); + } + return null; + } + + /** + * 获取累计数据 + */ + public AccumulatedData getAccumulatedData(Long itemNo) { + AccumulatedLogicService accLogic = accumulatedLogicMap.get(itemNo); + if (accLogic != null) { + return accLogic.getData(); + } + return AccumulatedData.builder().volume(0.0).turnover(0.0).build(); + } + + /** + * 获取买卖强度 + */ + public BuySellStrength getBuySellStrength(Long itemNo, boolean is15Tick) { + BuySellStrengthService bssLogic = buySellStrengthMap.get(itemNo); + if (bssLogic != null) { + return bssLogic.getData(); + } + return BuySellStrength.builder().high(0.0f).mid(0.0f).low(0.0f).build(); + } + + /** + * 触发业务逻辑更新 + * 对应C++的PostBusinessLogic方法 + */ + private void postBusinessLogic(TranTable table, TranRecord record) { + Long itemNo = table.getItemNo(); + + // 更新MTC逻辑 + MTCLogicService mtcLogic = mtcLogicMap.get(itemNo); + if (mtcLogic != null) { + mtcLogic.update(record, 0); + } + + // 更新累计数据逻辑 + AccumulatedLogicService accLogic = accumulatedLogicMap.get(itemNo); + if (accLogic != null) { + accLogic.update(record); + } + + // 更新买卖强度逻辑 + BuySellStrengthService bssLogic = buySellStrengthMap.get(itemNo); + if (bssLogic != null) { + bssLogic.update(record); + } + } + + /** + * 删除交易品种 + */ + public boolean dropItem(Long itemNo) { + lock.writeLock().lock(); + try { + tableCache.remove(itemNo); + log.info("[TranTableService] Dropped table for itemNo={}", itemNo); + return true; + } finally { + lock.writeLock().unlock(); + } + } +} diff --git a/src/main/java/com/afe/dc/tranlog/service/multicast/COPMessageBuilder.java b/src/main/java/com/afe/dc/tranlog/service/multicast/COPMessageBuilder.java new file mode 100644 index 0000000..84ce4fc --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/service/multicast/COPMessageBuilder.java @@ -0,0 +1,253 @@ +package com.afe.dc.tranlog.service.multicast; + +import com.afe.dc.tranlog.domain.dto.COPItem; +import lombok.extern.slf4j.Slf4j; + +import java.nio.ByteBuffer; +import java.nio.ByteOrder; +import java.nio.charset.StandardCharsets; +import java.util.Map; + +/** + * COP协议消息构建器 + * 将COPItem对象转换为字节数组(COP协议格式) + * 对应C++的BUFFER_ITEM::GetMessage() + * + * @author ben.yang + */ +@Slf4j +public class COPMessageBuilder { + + /** + * COP消息头大小(字节) + * 消息类型(1) + ItemNo(4) + 其他头部信息(11) = 16字节 + */ + private static final int COP_HEADER_SIZE = 16; + + /** + * FID字段头部大小 + * FID编号(2) + 数据类型(1) + 数据长度(2) = 5字节 + */ + private static final int FID_HEADER_SIZE = 5; + + /** + * 构建COP协议消息 + * + * @param item COP数据项 + * @return 字节数组,如果构建失败返回null + */ + public static byte[] buildMessage(COPItem item) { + if (item == null) { + log.warn("[COPMessageBuilder] Cannot build message from null COPItem"); + return null; + } + + try { + // 计算总大小 + int totalSize = COP_HEADER_SIZE; + Map fields = item.getFields(); + + if (fields != null && !fields.isEmpty()) { + for (Map.Entry entry : fields.entrySet()) { + Object value = entry.getValue(); + int fieldSize = calculateFieldSize(value); + totalSize += FID_HEADER_SIZE + fieldSize; + } + } + + // 创建字节缓冲区 + ByteBuffer buffer = ByteBuffer.allocate(totalSize); + buffer.order(ByteOrder.LITTLE_ENDIAN); // COP协议使用小端序 + + // 写入消息头 + writeHeader(buffer, item); + + // 写入FID字段 + if (fields != null && !fields.isEmpty()) { + for (Map.Entry entry : fields.entrySet()) { + writeField(buffer, entry.getKey(), entry.getValue()); + } + } + + return buffer.array(); + + } catch (Exception e) { + log.error("[COPMessageBuilder] Error building message for itemNo: {}", item.getItemNo(), e); + return null; + } + } + + /** + * 写入消息头 + * + * @param buffer 字节缓冲区 + * @param item COP数据项 + */ + private static void writeHeader(ByteBuffer buffer, COPItem item) { + // 消息类型(1字节) + int msgType = item.getMsgType() != null ? item.getMsgType() : COPItem.MsgType.MGT_UPDATE; + buffer.put((byte) msgType); + + // ItemNo(4字节,小端序) + long itemNo = item.getItemNo() != null ? item.getItemNo() : 0; + buffer.putInt((int) itemNo); + + // 填充头部剩余部分(11字节) + // 实际COP协议可能包含更多头部信息,这里用0填充 + for (int i = 0; i < 11; i++) { + buffer.put((byte) 0); + } + } + + /** + * 写入FID字段 + * + * @param buffer 字节缓冲区 + * @param fid FID编号 + * @param value 字段值 + */ + private static void writeField(ByteBuffer buffer, Integer fid, Object value) { + if (value == null) { + return; + } + + try { + // 写入FID编号(2字节) + buffer.putShort(fid.shortValue()); + + // 确定数据类型和值 + byte dataType; + byte[] fieldData = getFieldData(value); + + if (fieldData == null) { + log.warn("[COPMessageBuilder] Cannot serialize field value for FID: {}", fid); + return; + } + + dataType = getDataType(value); + + // 写入数据类型(1字节) + buffer.put(dataType); + + // 写入数据长度(2字节) + buffer.putShort((short) fieldData.length); + + // 写入数据内容 + buffer.put(fieldData); + + } catch (Exception e) { + log.error("[COPMessageBuilder] Error writing field FID: {}", fid, e); + } + } + + /** + * 获取字段数据字节数组 + * + * @param value 字段值 + * @return 字节数组 + */ + private static byte[] getFieldData(Object value) { + if (value == null) { + return null; + } + + try { + if (value instanceof Byte || value instanceof Character) { + return new byte[] { value instanceof Byte ? (Byte) value : (byte) ((Character) value).charValue() }; + } else if (value instanceof Short) { + ByteBuffer bb = ByteBuffer.allocate(2); + bb.order(ByteOrder.LITTLE_ENDIAN); + bb.putShort((Short) value); + return bb.array(); + } else if (value instanceof Integer) { + ByteBuffer bb = ByteBuffer.allocate(4); + bb.order(ByteOrder.LITTLE_ENDIAN); + bb.putInt((Integer) value); + return bb.array(); + } else if (value instanceof Long) { + ByteBuffer bb = ByteBuffer.allocate(8); + bb.order(ByteOrder.LITTLE_ENDIAN); + bb.putLong((Long) value); + return bb.array(); + } else if (value instanceof Float) { + ByteBuffer bb = ByteBuffer.allocate(4); + bb.order(ByteOrder.LITTLE_ENDIAN); + bb.putFloat((Float) value); + return bb.array(); + } else if (value instanceof Double) { + ByteBuffer bb = ByteBuffer.allocate(8); + bb.order(ByteOrder.LITTLE_ENDIAN); + bb.putDouble((Double) value); + return bb.array(); + } else if (value instanceof String) { + return ((String) value).getBytes(StandardCharsets.UTF_8); + } else if (value instanceof byte[]) { + return (byte[]) value; + } else { + // 尝试转换为字符串 + return value.toString().getBytes(StandardCharsets.UTF_8); + } + } catch (Exception e) { + log.error("[COPMessageBuilder] Error converting value to bytes: {}", value.getClass().getName(), e); + return null; + } + } + + /** + * 获取数据类型标识 + * + * @param value 字段值 + * @return 数据类型字节 + */ + private static byte getDataType(Object value) { + if (value instanceof Byte || value instanceof Character) { + return 0x01; // CHAR + } else if (value instanceof Short) { + return 0x02; // SHORT + } else if (value instanceof Integer) { + return 0x04; // INT + } else if (value instanceof Long) { + return 0x08; // LONG + } else if (value instanceof Float) { + return 0x10; // FLOAT + } else if (value instanceof Double) { + return 0x20; // DOUBLE + } else if (value instanceof String) { + return 0x40; // STRING + } +// else if (value instanceof byte[]) { +// return 0x80; // BYTE_ARRAY +// } + else { + return 0x40; // 默认作为字符串处理 + } + } + + /** + * 计算字段大小 + * + * @param value 字段值 + * @return 字段大小(字节) + */ + private static int calculateFieldSize(Object value) { + if (value == null) { + return 0; + } + + if (value instanceof Byte || value instanceof Character) { + return 1; + } else if (value instanceof Short) { + return 2; + } else if (value instanceof Integer || value instanceof Float) { + return 4; + } else if (value instanceof Long || value instanceof Double) { + return 8; + } else if (value instanceof String) { + return ((String) value).getBytes(StandardCharsets.UTF_8).length; + } else if (value instanceof byte[]) { + return ((byte[]) value).length; + } else { + return value.toString().getBytes(StandardCharsets.UTF_8).length; + } + } +} diff --git a/src/main/java/com/afe/dc/tranlog/service/multicast/MulticastSender.java b/src/main/java/com/afe/dc/tranlog/service/multicast/MulticastSender.java new file mode 100644 index 0000000..ffd96b7 --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/service/multicast/MulticastSender.java @@ -0,0 +1,331 @@ +package com.afe.dc.tranlog.service.multicast; + +import com.afe.dc.tranlog.domain.dto.COPItem; +import lombok.extern.slf4j.Slf4j; + +import java.io.IOException; +import java.net.*; +import java.nio.ByteBuffer; +import java.nio.ByteOrder; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.atomic.AtomicBoolean; + +/** + * UDP多播发送器 + * 独立功能,不依赖Spring框架 + * 对应C++的CDataCtrl::m_SendCtrl + * 负责通过UDP多播发送COP协议数据 + * + * @author ben.yang + */ +@Slf4j +public class MulticastSender { + + /** + * 发送配置 + */ + public static class Config { + /** 发送端口 */ + private int sendPort; + /** 多播组IP地址 */ + private String groupIp; + /** 本地网络接口IP地址 */ + private String localIp = "0.0.0.0"; + /** TTL值(Time To Live) */ + private int ttl = 1; + /** 是否启用循环发送(用于测试) */ + private boolean loopbackDisabled = true; + + public Config(int sendPort, String groupIp) { + this.sendPort = sendPort; + this.groupIp = groupIp; + } + + public Config setLocalIp(String localIp) { + this.localIp = localIp; + return this; + } + + public Config setTtl(int ttl) { + this.ttl = ttl; + return this; + } + + public Config setLoopbackDisabled(boolean loopbackDisabled) { + this.loopbackDisabled = loopbackDisabled; + return this; + } + + // Getters + public int getSendPort() { return sendPort; } + public String getGroupIp() { return groupIp; } + public String getLocalIp() { return localIp; } + public int getTtl() { return ttl; } + public boolean isLoopbackDisabled() { return loopbackDisabled; } + } + + private final Config config; + private MulticastSocket socket; + private InetAddress groupAddress; + private final AtomicBoolean initialized = new AtomicBoolean(false); + private final List listeners = new CopyOnWriteArrayList<>(); + + /** + * 发送监听器接口 + */ + public interface SendListener { + /** + * 发送成功回调 + * + * @param item 发送的COPItem + * @param bytesSent 发送的字节数 + */ + void onSendSuccess(COPItem item, int bytesSent); + + /** + * 发送失败回调 + * + * @param item 发送的COPItem + * @param error 错误信息 + */ + void onSendError(COPItem item, Exception error); + } + + /** + * 构造函数 + * + * @param config 发送配置 + */ + public MulticastSender(Config config) { + this.config = config; + } + + /** + * 初始化并启动发送器 + * 对应C++的m_SendCtrl.Start()和AddChannel() + * + * @return 是否初始化成功 + */ + public boolean initialize() { + if (initialized.get()) { + log.warn("[MulticastSender] Already initialized"); + return true; + } + + try { + // 创建UDP多播Socket + socket = new MulticastSocket(); + + // 设置TTL(Time To Live) + socket.setTimeToLive(config.getTtl()); + + // 设置是否允许回环(使用StandardSocketOptions避免废弃警告) + socket.setOption(StandardSocketOptions.IP_MULTICAST_LOOP, !config.isLoopbackDisabled()); + + // 绑定本地网络接口 + if (config.getLocalIp() != null && !"0.0.0.0".equals(config.getLocalIp())) { + try { + InetAddress localInterface = InetAddress.getByName(config.getLocalIp()); + NetworkInterface networkInterface = NetworkInterface.getByInetAddress(localInterface); + if (networkInterface != null) { + socket.setOption(StandardSocketOptions.IP_MULTICAST_IF, networkInterface); + log.info("[MulticastSender] Set local interface to: {}", config.getLocalIp()); + } + } catch (Exception e) { + log.warn("[MulticastSender] Failed to set local interface: {}", config.getLocalIp(), e); + } + } + + // 解析多播组地址 + groupAddress = InetAddress.getByName(config.getGroupIp()); + if (!groupAddress.isMulticastAddress()) { + log.error("[MulticastSender] Invalid multicast address: {}", config.getGroupIp()); + close(); + return false; + } + + initialized.set(true); + log.info("[MulticastSender] Initialized successfully - Group: {}:{}, LocalIP: {}", + config.getGroupIp(), config.getSendPort(), config.getLocalIp()); + + return true; + + } catch (Exception e) { + log.error("[MulticastSender] Failed to initialize", e); + close(); + return false; + } + } + + /** + * 发送COPItem数据 + * 对应C++的m_SendCtrl.Send(item) + * + * @param item COP数据项 + * @return 是否发送成功 + */ + public boolean send(COPItem item) { + if (!initialized.get()) { + log.error("[MulticastSender] Not initialized, call initialize() first"); + return false; + } + + if (item == null) { + log.warn("[MulticastSender] Cannot send null COPItem"); + return false; + } + + try { + // 将COPItem转换为字节数组 + byte[] data = COPMessageBuilder.buildMessage(item); + + if (data == null || data.length == 0) { + log.warn("[MulticastSender] Failed to build message for COPItem: itemNo={}", item.getItemNo()); + notifyError(item, new IOException("Failed to build message")); + return false; + } + + // 创建数据包 + DatagramPacket packet = new DatagramPacket( + data, + data.length, + groupAddress, + config.getSendPort() + ); + + // 发送数据包 + socket.send(packet); + + int bytesSent = data.length; + log.debug("[MulticastSender] Sent COPItem - itemNo: {}, msgType: {}, size: {} bytes", + item.getItemNo(), item.getMsgType(), bytesSent); + + notifySuccess(item, bytesSent); + return true; + + } catch (IOException e) { + log.error("[MulticastSender] Error sending COPItem - itemNo: {}", item.getItemNo(), e); + notifyError(item, e); + return false; + } + } + + /** + * 发送原始字节数据 + * + * @param data 字节数据 + * @return 是否发送成功 + */ + public boolean sendRaw(byte[] data) { + if (!initialized.get()) { + log.error("[MulticastSender] Not initialized, call initialize() first"); + return false; + } + + if (data == null || data.length == 0) { + log.warn("[MulticastSender] Cannot send empty data"); + return false; + } + + try { + DatagramPacket packet = new DatagramPacket( + data, + data.length, + groupAddress, + config.getSendPort() + ); + + socket.send(packet); + log.debug("[MulticastSender] Sent raw data - size: {} bytes", data.length); + return true; + + } catch (IOException e) { + log.error("[MulticastSender] Error sending raw data", e); + return false; + } + } + + /** + * 添加发送监听器 + * + * @param listener 监听器 + */ + public void addListener(SendListener listener) { + if (listener != null) { + listeners.add(listener); + } + } + + /** + * 移除发送监听器 + * + * @param listener 监听器 + */ + public void removeListener(SendListener listener) { + listeners.remove(listener); + } + + /** + * 关闭发送器 + * 对应C++的m_SendCtrl.Stop()和DelChannel() + */ + public void close() { + if (!initialized.get()) { + return; + } + + initialized.set(false); + + if (socket != null && !socket.isClosed()) { + socket.close(); + log.info("[MulticastSender] Closed"); + } + + listeners.clear(); + } + + /** + * 检查是否已初始化 + * + * @return 是否已初始化 + */ + public boolean isInitialized() { + return initialized.get() && socket != null && !socket.isClosed(); + } + + /** + * 获取配置信息 + * + * @return 配置对象 + */ + public Config getConfig() { + return config; + } + + /** + * 通知发送成功 + */ + private void notifySuccess(COPItem item, int bytesSent) { + for (SendListener listener : listeners) { + try { + listener.onSendSuccess(item, bytesSent); + } catch (Exception e) { + log.warn("[MulticastSender] Error in send listener", e); + } + } + } + + /** + * 通知发送失败 + */ + private void notifyError(COPItem item, Exception error) { + for (SendListener listener : listeners) { + try { + listener.onSendError(item, error); + } catch (Exception e) { + log.warn("[MulticastSender] Error in send listener", e); + } + } + } +} diff --git a/src/main/java/com/afe/dc/tranlog/service/multicast/MulticastSenderExample.java b/src/main/java/com/afe/dc/tranlog/service/multicast/MulticastSenderExample.java new file mode 100644 index 0000000..bc47278 --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/service/multicast/MulticastSenderExample.java @@ -0,0 +1,79 @@ +package com.afe.dc.tranlog.service.multicast; + +import com.afe.dc.tranlog.domain.dto.COPItem; + +/** + * MulticastSender使用示例 + * 演示如何使用UDP多播发送器发送COP协议数据 + * + * @author ben.yang + */ +public class MulticastSenderExample { + + public static void main(String[] args) { + // 创建发送配置 + // 发送端口 + MulticastSender.Config config = new MulticastSender.Config( + 5001, + // 多播组IP + "225.6.7.9" + ); + // 本地IP(可选) + config.setLocalIp("0.0.0.0") + // TTL值(可选) + .setTtl(1) + // 禁用回环(可选) + .setLoopbackDisabled(true); + + // 创建发送器 + MulticastSender sender = new MulticastSender(config); + + // 添加发送监听器(可选) + sender.addListener(new MulticastSender.SendListener() { + @Override + public void onSendSuccess(COPItem item, int bytesSent) { + System.out.println("发送成功 - itemNo: " + item.getItemNo() + ", 大小: " + bytesSent + " 字节"); + } + + @Override + public void onSendError(COPItem item, Exception error) { + System.err.println("发送失败 - itemNo: " + item.getItemNo() + ", 错误: " + error.getMessage()); + } + }); + + // 初始化发送器 + if (!sender.initialize()) { + System.err.println("初始化失败"); + return; + } + + try { + // 创建COPItem数据 + COPItem item = COPItem.builder() + .itemNo(12345L) + .msgType(COPItem.MsgType.MGT_UPDATE) + .build(); + + // 添加FID字段 + item.addField(COPItem.FID.FID_TRAN_LOG, "交易数据"); + item.addField(COPItem.FID.FID_ASK, 100.5f); + item.addField(COPItem.FID.FID_BID, 100.3f); + + // 发送数据 + boolean success = sender.send(item); + if (success) { + System.out.println("数据发送成功"); + } else { + System.err.println("数据发送失败"); + } + + // 也可以发送原始字节数据 + byte[] rawData = new byte[]{0x01, 0x02, 0x03, 0x04}; + sender.sendRaw(rawData); + + } finally { + // 关闭发送器 + sender.close(); + } + } +} diff --git a/src/main/java/com/afe/dc/tranlog/task/ScheduledTask.java b/src/main/java/com/afe/dc/tranlog/task/ScheduledTask.java new file mode 100644 index 0000000..bf6dd3a --- /dev/null +++ b/src/main/java/com/afe/dc/tranlog/task/ScheduledTask.java @@ -0,0 +1,96 @@ +package com.afe.dc.tranlog.task; + +import com.afe.dc.tranlog.service.TranDatabaseService; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +/** + * 定时任务类 + * 对应C++的定时器事件处理 + */ +@Slf4j +@Component +@RequiredArgsConstructor +public class ScheduledTask { + + private final TranDatabaseService tranDatabaseService; + + /** + * MTC发送任务 + * 对应C++的MtcTimerEvent + * 每50毫秒执行一次(可配置) + */ + @Scheduled(fixedRate = 50) + public void mtcSend() { + try { + // 发送MTC数据 + // 实际实现应该遍历所有交易表并发送MTC数据 + log.debug("[ScheduledTask] MTC send task executed"); + } catch (Exception e) { + log.error("[ScheduledTask] Error in MTC send task", e); + } + } + + /** + * 定期重建任务 + * 对应C++的TimerEvent + * 每30秒执行一次(可配置) + */ + @Scheduled(fixedRate = 30000) + public void regularContribution() { + try { + // 定期发送增值数据 + // 实际实现应该调用TranDatabaseService的RegularContribution方法 + log.debug("[ScheduledTask] Regular contribution task executed"); + } catch (Exception e) { + log.error("[ScheduledTask] Error in regular contribution task", e); + } + } + + /** + * 价格-成交量发送任务 + * 对应C++的VolPriceTimerEvent + * 每50毫秒执行一次(可配置) + */ + @Scheduled(fixedRate = 50) + public void volPriceSend() { + try { + // 发送价格-成交量数据 + log.debug("[ScheduledTask] Volume price send task executed"); + } catch (Exception e) { + log.error("[ScheduledTask] Error in volume price send task", e); + } + } + + /** + * 开盘维护任务 + * 对应C++的OpenHouseKeepEvent + * 每天开盘前执行 + */ + @Scheduled(cron = "0 0 9 * * ?") + public void openHouseKeep() { + try { + log.info("[ScheduledTask] Open house keep task executed"); + // 执行开盘维护操作 + } catch (Exception e) { + log.error("[ScheduledTask] Error in open house keep task", e); + } + } + + /** + * 收盘维护任务 + * 对应C++的DayCloseHouseKeepEvent + * 每天收盘后执行 + */ + @Scheduled(cron = "0 0 16 * * ?") + public void dayCloseHouseKeep() { + try { + log.info("[ScheduledTask] Day close house keep task executed"); + // 执行收盘维护操作 + } catch (Exception e) { + log.error("[ScheduledTask] Error in day close house keep task", e); + } + } +} diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml new file mode 100644 index 0000000..602df14 --- /dev/null +++ b/src/main/resources/application.yml @@ -0,0 +1,58 @@ +server: + port: 8086 + +spring: + messages: + # 默认 messages, 这里我们多了一层目录i18n + basename: i18n/messages + # 如果默认 false, 则会出现匹配不到就会跑异常(NoSuchMessageException)的情况 + use-code-as-default-message: true + # 是否总是应用MessageFormat规则,即使是没有参数的消息也要解析, 默认 false + always-use-message-format: false + jackson: + time-zone: GMT+8 + date-format: yyyy-MM-dd HH:mm:ss + default-property-inclusion: ALWAYS + + +# MyBatis配置 +mybatis-plus: + # 搜索指定包别名 + typeAliasesPackage: com.afe.dc.** + # 配置mapper的扫描,找到所有的mapper.xml映射文件 + mapperLocations: classpath*:mapper/**/*Mapper.xml + global-config: + banner: false + configuration: + map-underscore-to-camel-case: true + cache-enabled: false + # spring boot集成mybatis的方式打印sql + log-impl: org.apache.ibatis.logging.stdout.StdOutImpl + +# 内部服务配置(用于Feign调用时添加Header) +internal: + service: + # 内部服务标识名称(必须与dc-data-service的配置一致) + name: dc-internal + # 内部服务密钥(如果dc-data-service配置了密钥,这里也需要配置相同的值) + secret: ${INTERNAL_SERVICE_SECRET:} + +# UDP多播接收配置 +# 对应C++的TLogServer.cpp中的接收配置 +multicast: + receiver: + # 是否启用UDP多播接收 + enabled: ${MULTICAST_RECEIVER_ENABLED:true} + # 接收端口(对应C++的RecvPort) + recv-port: ${MULTICAST_RECEIVER_PORT:5000} + # 本地网络接口IP地址(对应C++的IPAddress) + ip-address: ${MULTICAST_RECEIVER_IP:0.0.0.0} + # 多播组IP地址列表(对应C++的GroupIP0, GroupIP1, ...) + group-ips: + - ${MULTICAST_GROUP_IP_0:225.6.7.8} + # 可以配置多个多播组 + # - ${MULTICAST_GROUP_IP_1:225.6.7.9} + # 接收缓冲区大小(字节) + buffer-size: ${MULTICAST_BUFFER_SIZE:65536} + # 接收超时时间(毫秒) + timeout: ${MULTICAST_TIMEOUT:1000} diff --git a/src/main/resources/banner.txt b/src/main/resources/banner.txt new file mode 100644 index 0000000..69673d9 --- /dev/null +++ b/src/main/resources/banner.txt @@ -0,0 +1,24 @@ +Spring Boot Version: ${spring-boot.version} +Spring Application Name: ${spring.application.name} +//////////////////////////////////////////////////////////////////// +// _ooOoo_ // +// o8888888o // +// 88" . "88 // +// (| ^_^ |) // +// O\ = /O // +// ____/`---'\____ // +// .' \\| |// `. // +// / \\||| : |||// \ // +// / _||||| -:- |||||- \ // +// | | \\\ - /// | | // +// | \_| ''\---/'' | | // +// \ .-\__ `-` ___/-. / // +// ___`. .' /--.--\ `. . ___ // +// ."" '< `.___\_<|>_/___.' >'"". // +// | | : `- \`.;`\ _ /`;.`/ - ` : | | // +// \ \ `-. \_ __\ /__ _/ .-` / / // +// ========`-.____`-.___\_____/___.-`____.-'======== // +// `=---=' // +// ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ // +// 佛祖保佑 永不宕机 永无BUG // +//////////////////////////////////////////////////////////////////// diff --git a/src/main/resources/bootstrap.yml b/src/main/resources/bootstrap.yml new file mode 100644 index 0000000..b2a5d07 --- /dev/null +++ b/src/main/resources/bootstrap.yml @@ -0,0 +1,102 @@ +# Nacos帮助文档: https://nacos.io/zh-cn/docs/concepts.html +# Nacos认证信息 +# 指定的配置 +spring: + application: + # 应用名称 + name: dc-tranlog + profiles: + active: local + +--- +#默认环境(预生产部署环境) +spring: + config: + activate: + on-profile: pre-production + cloud: + nacos: + discovery: + # 服务注册地址 + server-addr: ai-nacos:8848 + config: + # 配置中心地址 + server-addr: ai-nacos:8848 + # 配置文件格式 + file-extension: yml + # 共享配置 + shared-configs: + - application.${spring.cloud.nacos.config.file-extension} + +--- + +#默认环境(jenkins部署环境) +spring: + config: + activate: + on-profile: dev + cloud: + nacos: + discovery: + # 服务注册地址 + server-addr: ai-nacos:8848 + config: + # 配置中心地址 + server-addr: ai-nacos:8848 + # 配置文件格式 + file-extension: yml + # 共享配置 + shared-configs: + - application-${spring.profiles.active}.${spring.cloud.nacos.config.file-extension} + +--- + +#localhost部署环境 +spring: + config: + activate: + on-profile: local + cloud: + nacos: + discovery: + # 服务注册地址 + server-addr: 127.0.0.1:8848 + config: + # 配置中心地址 + server-addr: 127.0.0.1:8848 + # 配置文件格式 + file-extension: yml + # 共享配置 + shared-configs: + - application-${spring.profiles.active}.${spring.cloud.nacos.config.file-extension} + + +--- + +#使用的远程的nacos在本地跑 +spring: + config: + activate: + on-profile: remote_nacos + cloud: + nacos: + discovery: + # 服务注册地址 + server-addr: 192.168.2.210:8848 + # # 命名空间 + namespace: ef971e56-0d03-4f14-b25e-7854881e87ed +# # 分组名称 +# group: LOCAL_GROUP + config: + # 配置中心地址 + server-addr: 192.168.2.210:8848 +# # 命名空间 + namespace: 41acd84d-bcda-4e9f-acc2-536bd8b9683c +# # 分组名称 +# group: LOCAL_GROUP + # 配置文件格式 + file-extension: yml + # 共享配置 + shared-configs: + #- application-${spring.profiles.active}.${spring.cloud.nacos.config.file-extension} + - application-common.${spring.cloud.nacos.config.file-extension} diff --git a/src/main/resources/logback.xml b/src/main/resources/logback.xml new file mode 100644 index 0000000..3d25d95 --- /dev/null +++ b/src/main/resources/logback.xml @@ -0,0 +1,74 @@ + + + + + + + + + + + ${log.pattern} + + + + + + ${log.path}/info.log + + + + ${log.path}/info.%d{yyyy-MM-dd}.log + + 60 + + + ${log.pattern} + + + + INFO + + ACCEPT + + DENY + + + + + ${log.path}/error.log + + + + ${log.path}/error.%d{yyyy-MM-dd}.log + + 60 + + + ${log.pattern} + + + + ERROR + + ACCEPT + + DENY + + + + + + + + + + + + + + + + + +