用于行情/交易/事件流的实时流数据接入原子 skill,适用于通用行业数据接入场景。
Scanned 6/6/2026
Install via CLI
openskills install aifinlab/FinClaw---
name: realtime-stream-ingestion
description: 用于行情/交易/事件流的实时流数据接入原子 skill,适用于通用行业数据接入场景。
---
# 实时流数据接入 Skill
## 数据来源
本 Skill 支持多种实时流数据输入格式,核心数据来源包括:
### 1. 实时流类型
- **行情流**:股票行情、期货行情、外汇行情
- **交易流**:交易数据流、订单数据流
- **事件流**:实时事件流、监控事件流
- **其他实时流**:日志流、指标流
### 2. 数据源类型
- **消息队列**:Kafka、RabbitMQ、RocketMQ等
- **WebSocket**:WebSocket实时数据流
- **API接口**:实时API、流式API
- **数据库**:实时数据库、时序数据库
### 3. 数据格式要求
- **连接信息**:消息队列连接信息、WebSocket地址
- **认证信息**:用户名密码、API密钥、Token
- **订阅配置**:数据订阅配置、过滤规则
- **数据格式**:实时数据格式和编码方式
### 4. 实时流特征
- **数据频率**:高频数据(毫秒级)、中频数据(秒级)
- **数据规模**:小规模数据、大规模数据
- **数据实时性**:实时数据、近实时数据
- **数据连续性**:连续数据流、离散数据流
> 说明:本 Skill 不包含数据源配置功能,需要用户提供实时流连接信息。建议实时流稳定可靠,以便进行准确的实时数据接入。
---
## 功能
本 Skill 提供全面的实时流数据接入能力,涵盖多种接入功能:
### 1. 实时流连接
- **连接管理**:管理实时流连接
- **连接验证**:验证实时流连接
- **连接监控**:监控实时流状态
- **连接优化**:优化连接性能
### 2. 数据订阅
- **主题订阅**:订阅数据主题
- **过滤订阅**:基于条件过滤订阅
- **多主题订阅**:订阅多个数据主题
- **动态订阅**:动态调整订阅配置
### 3. 实时接收
- **实时接收**:实时接收数据流
- **流式处理**:流式处理数据
- **缓冲处理**:缓冲处理数据
- **背压处理**:处理数据背压
### 4. 数据解析
- **数据解析**:解析实时数据格式
- **数据分类**:对数据进行分类
- **数据标签**:为数据添加标签
- **数据提取**:提取数据关键信息
### 5. 数据处理
- **数据过滤**:过滤不需要的数据
- **数据转换**:转换数据格式
- **数据聚合**:聚合相关数据
- **数据采样**:采样数据流
### 6. 高级处理功能
- **数据存储**:存储实时数据
- **数据索引**:建立数据索引
- **数据查询**:支持数据查询功能
- **接入报告**:生成数据接入报告
---
## 使用示例
### 输出示例
```json
{
"source_info": {
"source_type": "websocket",
"source_name": "market_data_stream",
"endpoint": "ws://market.example.com/stream",
"data_type": "stock_quotes"
},
"ingestion_config": {
"subscription_mode": "subscribe",
"symbols": ["000001", "000002", "600519"],
"update_frequency": "real_time",
"buffer_size": 1000
},
"ingestion_results": {
"total_messages": 100000,
"processed_messages": 99500,
"failed_messages": 500,
"ingestion_start_time": "2024-03-15T10:00:00",
"ingestion_end_time": "2024-03-15T11:00:00",
"duration": "3600s"
},
"data_samples": [
{
"message_id": "MSG001",
"timestamp": "2024-03-15T10:00:00.123",
"symbol": "000001",
"data": {
"price": 10.50,
"volume": 1000000,
"bid": 10.49,
"ask": 10.51
},
"processed": true
}
],
"statistics": {
"messages_received": 100000,
"messages_processed": 99500,
"messages_stored": 99500,
"processing_rate": "27.6 messages/s",
"success_rate": 0.995,
"average_latency": "50ms"
}
}
```
---
## 注意事项与限制
### 1. 实时流要求
- 实时流需要稳定可靠
- 连接信息需要准确
- 网络延迟需要控制
### 2. 实时处理性能
- 高频数据可能影响性能
- 需要优化处理流程
- 需要合理设置缓冲
### 3. 数据质量
- 实时数据质量影响处理结果
- 需要验证数据完整性
- 异常数据需要处理
### 4. 数据延迟
- 实时处理可能存在延迟
- 需要监控处理延迟
- 需要优化处理速度
### 5. 使用限制
- 本 Skill 不包含数据源管理功能
- 接入结果需要人工复核
- 复杂实时流可能需要特殊处理
---
## 参考资料
- 见 references/ 目录中的相关文档,包括:
- 实时流数据接入方法手册
- WebSocket配置指南
- 实时处理策略说明
- 性能优化指南

No comments yet. Be the first to comment!