Sink
Sink 是数据管道的出口:负责将经 Processors 处理后的数据写入 FlowMQ 或外部系统。
示例场景
后文各 Sink 子文档共用同一示例:从 Kafka Topic flowmq.mqtt.kafka 读取消息,再写入目标系统。
Topic 名虽含
mqtt,仅为命名;本示例的 Source 类型是 Kafka。
Source(Kafka)建议配置:
| 参数 | 建议值 |
|---|---|
| seed_brokers | 127.0.0.1:9092 |
| topics | flowmq.mqtt.kafka |
| consumer_group | 按目标 Sink 区分,例如 flowmq-sink-mysql |
样例消息(JSON):
json
{
"id": "e6eb1d31-5ab7-4ace-8bf7-e7c3e7a30a73",
"source": "1SDA526VD_POP_P1",
"time": "2023-05-30T15:23:23.123+08:00",
"type": "1",
"data": {
"value": "280.44",
"quality": "0",
"timestamp": 1710400220000
}
}字段映射时常用:
| 字段路径 | 含义 |
|---|---|
source | 设备/测点标识 |
time | 事件时间 |
type | 类型或单位代码 |
data.value | 测点数值 |
data.timestamp | 时序时间戳(Unix 毫秒) |
适用场景
- 将业务流数据实时写入时序库、分析库或数据湖
- 对外分发事件到搜索、告警或业务服务系统
- 将处理后的结果写回 FlowMQ 数据流
配置要点
在创建管道向导的 Sink 步骤中:
- 选择目标连接器类型
- 填写连接参数
- 除
Drop sink外,需先测试通过再继续
完整创建流程见 数据集成。
部分 Sink 字段(如 subject、topic、table)支持 ${! ... } 按消息动态取值,详见各目标文档与 Processors。
Sink 选项(FlowMQ 控制台)
可选目标包括:
- RabbitMQ、DolphinDB、HTTP、IoTDB、MQTT、NATS
- Redis Hash、Redis List、Redis PubSub、Redis Streams
- MySQL、PostgreSQL、TDengine
- Drop sink(丢弃输出消息,不写入目标)
Sink 连接器
- RabbitMQ Sink
- DolphinDB Sink
- HTTP Sink
- IoTDB Sink
- MQTT Sink
- NATS Sink
- Redis Hash Sink
- Redis List Sink
- Redis PubSub Sink
- Redis Streams Sink
- MySQL Sink
- PostgreSQL Sink
- TDengine Sink
- Drop Sink
注意事项
- 建议先通过 Sink 连通性测试,再创建管道
- Source 与 Sink 均可连接 FlowMQ 或外部系统,按实际链路组合即可
- 目标端若支持幂等写入,建议启用
- 注意字段类型映射,避免写入时类型不兼容
- Processors 按单条消息处理;复杂关联或较重业务逻辑宜放在前置系统
- MQTT / NATS 等:发布目标勿选择会再次进入本管道 Source 的主题,以免形成消息环路
- Redis 使用
simple模式时,master仍需填写非空值(如mymaster);地址中的库号字段名为database,勿使用db - TDengine / PostgreSQL / MySQL:ISO 时间写入
TIMESTAMP请用ts_parse;Unix 毫秒写入BIGINT