Files
PolyHermes/docs/zh/copy-trading-implementation.md
WrBug a222c9a52f docs: 重构文档结构,使用 zh/ 和 en/ 目录区分中英文文档
- 创建 docs/zh/ 和 docs/en/ 目录结构
- 将所有中文文档移动到 docs/zh/
- 创建主要文档的英文版本:
  - DEPLOYMENT.md (651行)
  - DEVELOPMENT.md (514行)
  - VERSION_MANAGEMENT.md (已有)
- 更新所有文档中的内部链接
- 更新 README.md 和 README_EN.md 中的文档链接
- 在文档中添加中英文版本互链
2025-12-07 18:01:11 +08:00

551 lines
16 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 跟单买入和卖出实现方案
## 1. 核心思路
### 1.1 基本原理
- **买入跟单**:当 Leader 执行 `BUY` 交易时,系统自动创建 `BUY` 订单
- **卖出跟单**:当 Leader 执行 `SELL` 交易时,系统自动创建 `SELL` 订单
- **方向复制**:直接复制 Leader 交易的 `side` 字段(BUY 或 SELL
### 1.2 数据来源
- **方式1(优先)**WebSocket 推送(RTDS API
- 实时接收 Leader 的交易推送
- WebSocket URL: `wss://ws-live-data.polymarket.com`
- 订阅用户交易频道,实时获取交易数据
- **方式2(备选)**:轮询 CLOB API
- 通过 CLOB API `/trades?user={leaderAddress}` 获取 Leader 的交易记录
- 定期轮询(默认每 5 秒)
**交易数据包含**
- `side`: "BUY" 或 "SELL"(直接复制)
- `market`: 市场地址(直接复制)
- `price`: 交易价格(可调整)
- `size`: 交易数量(按比例或固定金额计算)
## 2. 实现流程
### 2.1 监控 Leader 交易
#### 2.1.1 WebSocket 推送模式(优先)
```kotlin
/**
* WebSocket 推送监控服务
*/
@Service
class CopyTradingWebSocketService(
private val leaderRepository: LeaderRepository,
private val configRepository: CopyTradingConfigRepository
) {
private var webSocketClient: WebSocketClient? = null
private val subscribedLeaders = mutableSetOf<String>()
/**
* 初始化 WebSocket 连接
*/
@PostConstruct
fun initWebSocket() {
val config = configRepository.findFirstByOrderByIdAsc() ?: getDefaultConfig()
if (config.useWebSocket) {
connectWebSocket()
}
}
/**
* 连接 WebSocket
*/
private fun connectWebSocket() {
try {
val wsUrl = "wss://ws-live-data.polymarket.com"
webSocketClient = WebSocketClient(wsUrl)
webSocketClient?.onMessage { message ->
handleWebSocketMessage(message)
}
webSocketClient?.onError { error ->
logger.error("WebSocket 连接错误", error)
// 降级到轮询模式
fallbackToPolling()
}
webSocketClient?.onClose {
logger.warn("WebSocket 连接关闭,尝试重连")
reconnectWebSocket()
}
// 订阅所有启用的 Leader
subscribeAllLeaders()
} catch (e: Exception) {
logger.error("WebSocket 连接失败,降级到轮询模式", e)
fallbackToPolling()
}
}
/**
* 订阅所有启用的 Leader
*/
private fun subscribeAllLeaders() {
val enabledLeaders = leaderRepository.findByEnabledTrue()
enabledLeaders.forEach { leader ->
subscribeLeader(leader.leaderAddress)
}
}
/**
* 订阅单个 Leader
*/
fun subscribeLeader(leaderAddress: String) {
val subscribeMessage = jsonObjectOf(
"type" to "subscribe",
"channel" to "user",
"user" to leaderAddress
)
webSocketClient?.send(subscribeMessage.toString())
subscribedLeaders.add(leaderAddress)
}
/**
* 处理 WebSocket 消息
*/
private fun handleWebSocketMessage(message: String) {
try {
val json = JSONObject(message)
val channel = json.optString("channel")
val eventType = json.optString("event")
if (channel == "user" && eventType == "trade") {
val trade = parseTradeMessage(json)
// 触发跟单逻辑
processTradeFromWebSocket(trade)
}
} catch (e: Exception) {
logger.error("处理 WebSocket 消息失败", e)
}
}
/**
* 重连 WebSocket
*/
private fun reconnectWebSocket() {
val config = configRepository.findFirstByOrderByIdAsc() ?: getDefaultConfig()
var retryCount = 0
while (retryCount < config.websocketMaxRetries) {
try {
Thread.sleep(config.websocketReconnectInterval.toLong())
connectWebSocket()
return
} catch (e: Exception) {
retryCount++
logger.warn("WebSocket 重连失败 (${retryCount}/${config.websocketMaxRetries})", e)
}
}
// 重连失败,降级到轮询
logger.error("WebSocket 重连失败,降级到轮询模式")
fallbackToPolling()
}
}
```
#### 2.1.2 轮询模式(备选)
```kotlin
/**
* 轮询监控服务(WebSocket 不可用时的备选方案)
*/
@Service
class CopyTradingPollingService(
private val clobService: PolymarketClobService,
private val leaderRepository: LeaderRepository,
private val configRepository: CopyTradingConfigRepository
) {
/**
* 定期轮询所有启用的 Leader 的交易
*/
@Scheduled(fixedDelayString = "\${copy.trading.poll.interval:5000}")
suspend fun monitorLeaders() {
val config = configRepository.findFirstByOrderByIdAsc() ?: getDefaultConfig()
// 如果配置了使用 WebSocket 且 WebSocket 可用,则不轮询
if (config.useWebSocket && isWebSocketAvailable()) {
return
}
val enabledLeaders = leaderRepository.findByEnabledTrue()
enabledLeaders.forEach { leader ->
try {
processLeaderTrades(leader)
} catch (e: Exception) {
logger.error("处理 Leader ${leader.id} 交易失败", e)
}
}
}
/**
* 处理单个 Leader 的交易
*/
private suspend fun processLeaderTrades(leader: Leader) {
// 1. 获取 Leader 的最新交易记录
val tradesResult = clobService.getTrades(
market = null,
user = leader.leaderAddress,
limit = 50,
offset = 0
)
tradesResult.fold(
onSuccess = { trades ->
trades.forEach { trade ->
// 2. 检查是否已处理过(去重)
if (!isProcessed(leader.id, trade.id)) {
// 3. 处理交易(买入或卖出)
processTrade(leader, trade)
// 4. 标记为已处理
markAsProcessed(leader.id, trade.id)
}
}
},
onFailure = { e ->
logger.error("获取 Leader ${leader.id} 交易失败", e)
}
)
}
}
```
### 2.2 处理交易(买入/卖出)
```kotlin
/**
* 处理单笔交易,自动识别买入或卖出
*/
private suspend fun processTrade(leader: Leader, trade: TradeResponse) {
// 1. 验证分类筛选
if (leader.category != null) {
val marketCategory = getMarketCategory(trade.market)
if (marketCategory != leader.category) {
logger.debug("跳过交易:分类不匹配 ${trade.market}")
return
}
}
// 2. 验证风险控制
if (!checkRiskControl(leader)) {
logger.warn("风险控制限制,跳过跟单 Leader ${leader.id}")
return
}
// 3. 确定使用的账户
val account = getAccountForLeader(leader)
if (account == null) {
logger.error("无法获取账户,跳过跟单 Leader ${leader.id}")
return
}
// 4. 计算跟单订单参数
val orderParams = calculateOrderParams(leader, trade)
// 5. 创建跟单订单(买入或卖出)
createCopyOrder(leader, account, trade, orderParams)
}
/**
* 计算跟单订单参数
*/
private fun calculateOrderParams(leader: Leader, trade: TradeResponse): OrderParams {
// 获取配置
val globalConfig = configRepository.findFirstByOrderByIdAsc() ?: getDefaultConfig()
// 计算订单大小
val leaderSize = trade.size.toSafeBigDecimal()
val copyRatio = leader.copyRatio ?: globalConfig.copyRatio
var orderSize = leaderSize.multiply(copyRatio)
// 应用限制
val maxSize = leader.maxOrderSize ?: globalConfig.maxOrderSize
val minSize = leader.minOrderSize ?: globalConfig.minOrderSize
orderSize = orderSize.coerceIn(minSize, maxSize)
// 计算价格(默认使用 Leader 的价格)
val price = trade.price.toSafeBigDecimal()
// TODO: 可以根据价格容忍度调整价格
return OrderParams(
market = trade.market,
side = trade.side, // 直接复制 BUY 或 SELL
price = price,
size = orderSize
)
}
data class OrderParams(
val market: String,
val side: String, // "BUY" 或 "SELL"
val price: BigDecimal,
val size: BigDecimal
)
```
### 2.3 创建跟单订单
```kotlin
/**
* 创建跟单订单(买入或卖出)
*/
private suspend fun createCopyOrder(
leader: Leader,
account: Account,
trade: TradeResponse,
params: OrderParams
) {
try {
// 1. 创建订单请求
val orderRequest = CreateOrderRequest(
market = params.market,
side = params.side, // "BUY" 或 "SELL"
price = params.price.toPlainString(),
size = params.size.toPlainString(),
type = "LIMIT"
)
// 2. 使用账户的 API Key 创建订单
val apiKey = account.apiKey ?: throw IllegalStateException("账户 ${account.id} 未配置 API Key")
val clobApi = createClobApiWithApiKey(apiKey)
val orderResult = clobApi.createOrder(orderRequest)
orderResult.fold(
onSuccess = { orderResponse ->
// 3. 保存跟单记录
val copyOrder = CopyOrder(
accountId = account.id!!,
leaderId = leader.id!!,
leaderAddress = leader.leaderAddress,
leaderTradeId = trade.id,
marketId = params.market,
category = getMarketCategory(params.market),
side = params.side, // 保存买入或卖出方向
price = params.price,
size = params.size,
copyRatio = leader.copyRatio ?: BigDecimal.ONE,
orderId = orderResponse.id,
status = "created"
)
copyOrderRepository.save(copyOrder)
logger.info("成功创建跟单订单: Leader=${leader.id}, Side=${params.side}, Market=${params.market}")
},
onFailure = { e ->
logger.error("创建跟单订单失败: Leader=${leader.id}, Side=${params.side}", e)
// 记录失败的跟单订单
val copyOrder = CopyOrder(
accountId = account.id!!,
leaderId = leader.id!!,
leaderAddress = leader.leaderAddress,
leaderTradeId = trade.id,
marketId = params.market,
category = getMarketCategory(params.market),
side = params.side,
price = params.price,
size = params.size,
copyRatio = leader.copyRatio ?: BigDecimal.ONE,
status = "failed"
)
copyOrderRepository.save(copyOrder)
}
)
} catch (e: Exception) {
logger.error("创建跟单订单异常: Leader=${leader.id}, Side=${params.side}", e)
}
}
```
## 3. 关键实现细节
### 3.1 买入和卖出的区别
**买入跟单(BUY**
- Leader 执行 `BUY` 交易 → 系统创建 `BUY` 订单
- 订单参数:`side = "BUY"`
- 表示买入该市场的 YES 或 NO 代币
**卖出跟单(SELL**
- Leader 执行 `SELL` 交易 → 系统创建 `SELL` 订单
- 订单参数:`side = "SELL"`
- 表示卖出持有的代币
**实现上无区别**
- 买入和卖出的处理逻辑完全相同
- 只是 `side` 字段的值不同("BUY" 或 "SELL"
- 都通过 `createOrder` API 创建订单
### 3.2 去重机制
```kotlin
/**
* 检查交易是否已处理
*/
private suspend fun isProcessed(leaderId: Long, tradeId: String): Boolean {
return processedTradeRepository.existsByLeaderIdAndTradeId(leaderId, tradeId)
}
/**
* 标记交易为已处理
*/
private suspend fun markAsProcessed(leaderId: Long, tradeId: String) {
val processed = ProcessedTrade(
leaderId = leaderId,
tradeId = tradeId,
processedAt = System.currentTimeMillis()
)
processedTradeRepository.save(processed)
}
```
### 3.3 账户选择
```kotlin
/**
* 获取 Leader 使用的账户
*/
private suspend fun getAccountForLeader(leader: Leader): Account? {
return if (leader.accountId != null) {
// 使用 Leader 指定的账户
accountRepository.findById(leader.accountId)
} else {
// 使用默认账户
accountRepository.findByIsDefaultTrue()
}
}
```
### 3.4 风险控制检查
```kotlin
/**
* 检查风险控制限制
*/
private suspend fun checkRiskControl(leader: Leader): Boolean {
val config = configRepository.findFirstByOrderByIdAsc() ?: getDefaultConfig()
// 检查每日亏损限制
val todayLoss = getTodayLoss(leader.accountId ?: getDefaultAccountId())
if (todayLoss >= config.maxDailyLoss) {
logger.warn("达到每日亏损限制: $todayLoss >= ${config.maxDailyLoss}")
return false
}
// 检查每日订单数限制
val todayOrderCount = getTodayOrderCount(leader.accountId ?: getDefaultAccountId())
if (todayOrderCount >= config.maxDailyOrders) {
logger.warn("达到每日订单数限制: $todayOrderCount >= ${config.maxDailyOrders}")
return false
}
return true
}
```
## 4. 数据模型
### 4.1 ProcessedTrade(已处理交易)
```kotlin
@Entity
@Table(name = "copy_trading_processed_trades")
data class ProcessedTrade(
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
val id: Long? = null,
@Column(name = "leader_id", nullable = false)
val leaderId: Long,
@Column(name = "trade_id", nullable = false, length = 100)
val tradeId: String,
@Column(name = "processed_at", nullable = false)
val processedAt: Long = System.currentTimeMillis(),
@UniqueConstraint(columnNames = ["leader_id", "trade_id"])
)
```
## 5. 完整示例
### 5.1 买入跟单示例
```
1. Leader 执行买入交易:
- Trade: { id: "trade_123", market: "0x...", side: "BUY", price: "0.5", size: "100" }
2. 系统检测到交易:
- 验证分类、风险控制
- 计算跟单参数:size = 100 × 1.0 = 100
3. 创建跟单订单:
- CreateOrderRequest: { market: "0x...", side: "BUY", price: "0.5", size: "100" }
- 调用 CLOB API 创建订单
4. 保存记录:
- CopyOrder: { side: "BUY", ... }
```
### 5.2 卖出跟单示例
```
1. Leader 执行卖出交易:
- Trade: { id: "trade_456", market: "0x...", side: "SELL", price: "0.6", size: "50" }
2. 系统检测到交易:
- 验证分类、风险控制
- 计算跟单参数:size = 50 × 1.0 = 50
3. 创建跟单订单:
- CreateOrderRequest: { market: "0x...", side: "SELL", price: "0.6", size: "50" }
- 调用 CLOB API 创建订单
4. 保存记录:
- CopyOrder: { side: "SELL", ... }
```
## 6. 注意事项
### 6.1 交易 vs 订单
- **交易(Trade)**:已成交的记录,包含 `side` 字段
- **订单(Order)**:挂单,可能未成交
- 跟单系统基于**交易记录**触发,创建**订单**
### 6.2 价格和数量
- **价格**:默认使用 Leader 的交易价格,可配置价格容忍度
- **数量**:按跟单比例计算,应用最大/最小限制
### 6.3 错误处理
- API 调用失败时记录失败状态
- 网络异常时重试机制
- 记录详细日志便于排查
### 6.4 性能优化
- 批量查询多个 Leader 的交易
- 使用缓存减少重复查询
- 异步处理订单创建
## 7. 总结
**买入和卖出的实现完全相同**
- 都通过监控 Leader 的交易记录触发
- 都通过 `side` 字段区分("BUY" 或 "SELL"
- 都调用相同的 `createOrder` API
- 区别仅在于 `side` 参数的值
**核心流程**
1. 轮询 Leader 交易 → 2. 识别买入/卖出 → 3. 计算参数 → 4. 创建订单 → 5. 保存记录