This commit is contained in:
yanghanbin
2026-01-26 18:08:01 +08:00
commit 45b4f79c16
36 changed files with 4390 additions and 0 deletions
+44
View File
@@ -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
+196
View File
@@ -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
+178
View File
@@ -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. 添加性能监控和统计
+216
View File
@@ -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` 文件中的完整示例代码。
+87
View File
@@ -0,0 +1,87 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<!-- 继承父工程 -->
<parent>
<groupId>com.afe.dc</groupId>
<artifactId>dc-parent</artifactId>
<version>1.0.0-SNAPSHOT</version>
<relativePath>../dc-parent/pom.xml</relativePath>
</parent>
<artifactId>dc-tranlog-server</artifactId>
<name>DC TranLog Server</name>
<description>交易日志服务器 - 处理交易数据、MTC(分钟交易数据)、价格成交量等业务逻辑</description>
<properties>
<openfeign.version>4.1.0</openfeign.version>
</properties>
<dependencies>
<!-- 公共模块 Spring包-->
<dependency>
<groupId>com.afe.dc</groupId>
<artifactId>dc-common-settting</artifactId>
<version>${project.version}</version>
</dependency>
<!-- 公共模块 基础包 -->
<dependency>
<groupId>com.afe.dc</groupId>
<artifactId>dc-common-core</artifactId>
<version>${project.version}</version>
</dependency>
<!-- Data Dao 基础包 -->
<dependency>
<groupId>com.afe.dc</groupId>
<artifactId>dc-data-feign</artifactId>
<version>${project.version}</version>
</dependency>
</dependencies>
<build>
<finalName>${project.artifactId}</finalName>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<source>17</source>
<target>17</target>
<compilerArgs>
<arg>-parameters</arg>
</compilerArgs>
</configuration>
</plugin>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
<configuration>
<excludes>
<exclude>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
</exclude>
</excludes>
</configuration>
<executions>
<execution>
<goals>
<goal>repackage</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
@@ -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;
}
}
@@ -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;
}
}
@@ -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<String> groupIps = new ArrayList<>();
/**
* 是否启用UDP多播接收
*/
private Boolean enabled = true;
/**
* 接收缓冲区大小(字节)
*/
private Integer bufferSize = 65536;
/**
* 接收超时时间(毫秒)
*/
private Integer timeout = 1000;
}
@@ -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类中定义
}
@@ -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<List<MtcRecord>> getMtcRecords(
@PathVariable Long itemNo,
@RequestParam(required = false, defaultValue = "0") Integer dateIndex,
@RequestParam(required = false) Integer startIndex,
@RequestParam(required = false) Integer endIndex) {
try {
List<MtcRecord> 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<AccumulatedData> 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<BuySellStrength> 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<TranTable> 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();
}
}
}
@@ -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<Boolean> 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<Integer> processBatchData(@RequestBody java.util.List<COPItem> 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);
}
}
}
@@ -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<TradeDayResponse> 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<Boolean> 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);
}
}
}
@@ -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;
}
@@ -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;
}
@@ -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<Integer, Object> 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;
}
}
@@ -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;
}
@@ -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;
}
}
@@ -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;
}
}
@@ -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;
}
}
@@ -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<TranRecord> 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();
}
}
}
@@ -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');
}
}
@@ -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);
}
}
@@ -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<Integer, Object> 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<Integer, Object> parseFidFields(ByteBuffer buffer) {
Map<Integer, Object> 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;
}
}
}
@@ -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<Integer, List<MtcRecord>> 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<MtcRecord> 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<MtcRecord> 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<MtcRecord> 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<MtcRecord> 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<MtcRecord> 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);
}
}
@@ -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();
}
}
@@ -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<Long, TranTable> 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<Long> 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();
}
}
}
@@ -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<LocalDate> 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;
}
}
}
@@ -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<Long, TranTable> tableCache = new ConcurrentHashMap<>();
/**
* Logic服务映射(每个itemNo对应一个Logic实例)
* 实际应该为每个TranTable创建独立的Logic实例
*/
private final ConcurrentHashMap<Long, MTCLogicService> mtcLogicMap = new ConcurrentHashMap<>();
private final ConcurrentHashMap<Long, AccumulatedLogicService> accumulatedLogicMap = new ConcurrentHashMap<>();
private final ConcurrentHashMap<Long, BuySellStrengthService> 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();
}
}
}
@@ -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<Integer, Object> fields = item.getFields();
if (fields != null && !fields.isEmpty()) {
for (Map.Entry<Integer, Object> 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<Integer, Object> 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;
}
}
}
@@ -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<SendListener> 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);
}
}
}
}
@@ -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();
}
}
}
@@ -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);
}
}
}
+58
View File
@@ -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}
+24
View File
@@ -0,0 +1,24 @@
Spring Boot Version: ${spring-boot.version}
Spring Application Name: ${spring.application.name}
////////////////////////////////////////////////////////////////////
// _ooOoo_ //
// o8888888o //
// 88" . "88 //
// (| ^_^ |) //
// O\ = /O //
// ____/`---'\____ //
// .' \\| |// `. //
// / \\||| : |||// \ //
// / _||||| -:- |||||- \ //
// | | \\\ - /// | | //
// | \_| ''\---/'' | | //
// \ .-\__ `-` ___/-. / //
// ___`. .' /--.--\ `. . ___ //
// ."" '< `.___\_<|>_/___.' >'"". //
// | | : `- \`.;`\ _ /`;.`/ - ` : | | //
// \ \ `-. \_ __\ /__ _/ .-` / / //
// ========`-.____`-.___\_____/___.-`____.-'======== //
// `=---=' //
// ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ //
// 佛祖保佑 永不宕机 永无BUG //
////////////////////////////////////////////////////////////////////
+102
View File
@@ -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}
+74
View File
@@ -0,0 +1,74 @@
<?xml version="1.0" encoding="UTF-8"?>
<configuration scan="true" scanPeriod="60 seconds" debug="false">
<!-- 日志存放路径 -->
<property name="log.path" value="/data/logs/dc-tranlog-server" />
<!-- 日志输出格式 -->
<property name="log.pattern" value="%d{HH:mm:ss.SSS} [%thread] [%X{traceId}] %-5level %logger{20} - [%method,%line] - %msg%n" />
<!-- 控制台输出 -->
<appender name="console" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>${log.pattern}</pattern>
</encoder>
</appender>
<!-- 系统日志输出 -->
<appender name="file_info" class="ch.qos.logback.core.rolling.RollingFileAppender">
<file>${log.path}/info.log</file>
<!-- 循环政策:基于时间创建日志文件 -->
<rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
<!-- 日志文件名格式 -->
<fileNamePattern>${log.path}/info.%d{yyyy-MM-dd}.log</fileNamePattern>
<!-- 日志最大的历史 60天 -->
<maxHistory>60</maxHistory>
</rollingPolicy>
<encoder>
<pattern>${log.pattern}</pattern>
</encoder>
<filter class="ch.qos.logback.classic.filter.LevelFilter">
<!-- 过滤的级别 -->
<level>INFO</level>
<!-- 匹配时的操作:接收(记录) -->
<onMatch>ACCEPT</onMatch>
<!-- 不匹配时的操作:拒绝(不记录) -->
<onMismatch>DENY</onMismatch>
</filter>
</appender>
<appender name="file_error" class="ch.qos.logback.core.rolling.RollingFileAppender">
<file>${log.path}/error.log</file>
<!-- 循环政策:基于时间创建日志文件 -->
<rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
<!-- 日志文件名格式 -->
<fileNamePattern>${log.path}/error.%d{yyyy-MM-dd}.log</fileNamePattern>
<!-- 日志最大的历史 60天 -->
<maxHistory>60</maxHistory>
</rollingPolicy>
<encoder>
<pattern>${log.pattern}</pattern>
</encoder>
<filter class="ch.qos.logback.classic.filter.LevelFilter">
<!-- 过滤的级别 -->
<level>ERROR</level>
<!-- 匹配时的操作:接收(记录) -->
<onMatch>ACCEPT</onMatch>
<!-- 不匹配时的操作:拒绝(不记录) -->
<onMismatch>DENY</onMismatch>
</filter>
</appender>
<!-- 系统模块日志级别控制 -->
<logger name="com.afe" level="info" />
<!-- Spring日志级别控制 -->
<logger name="org.springframework" level="warn" />
<root level="info">
<appender-ref ref="console" />
</root>
<!--系统操作日志-->
<root level="info">
<appender-ref ref="file_info" />
<appender-ref ref="file_error" />
</root>
</configuration>