| 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. 高级处理功能
- 数据存储:存储实时数据
- 数据索引:建立数据索引
- 数据查询:支持数据查询功能
- 接入报告:生成数据接入报告
使用示例
输出示例
{
"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":
注意事项与限制
1. 实时流要求
- 实时流需要稳定可靠
- 连接信息需要准确
- 网络延迟需要控制
2. 实时处理性能
- 高频数据可能影响性能
- 需要优化处理流程
- 需要合理设置缓冲
3. 数据质量
- 实时数据质量影响处理结果
- 需要验证数据完整性
- 异常数据需要处理
4. 数据延迟
- 实时处理可能存在延迟
- 需要监控处理延迟
- 需要优化处理速度
5. 使用限制
- 本 Skill 不包含数据源管理功能
- 接入结果需要人工复核
- 复杂实时流可能需要特殊处理
参考资料
- 见 references/ 目录中的相关文档,包括:
- 实时流数据接入方法手册
- WebSocket配置指南
- 实时处理策略说明
- 性能优化指南