Skip to content

Sink

Sink 是数据管道的出口:负责将经 Processors 处理后的数据写入 FlowMQ 或外部系统。

示例场景

后文各 Sink 子文档共用同一示例:从 Kafka Topic flowmq.mqtt.kafka 读取消息,再写入目标系统。

Topic 名虽含 mqtt,仅为命名;本示例的 Source 类型是 Kafka。

Source(Kafka)建议配置:

参数建议值
seed_brokers127.0.0.1:9092
topicsflowmq.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 步骤中:

  1. 选择目标连接器类型
  2. 填写连接参数
  3. Drop sink 外,需先测试通过再继续

完整创建流程见 数据集成

部分 Sink 字段(如 subjecttopictable)支持 ${! ... } 按消息动态取值,详见各目标文档与 Processors

Sink 选项(FlowMQ 控制台)

可选目标包括:

  • RabbitMQ、DolphinDB、HTTP、IoTDB、MQTT、NATS
  • Redis Hash、Redis List、Redis PubSub、Redis Streams
  • MySQL、PostgreSQL、TDengine
  • Drop sink(丢弃输出消息,不写入目标)

Sink 连接器

注意事项

  • 建议先通过 Sink 连通性测试,再创建管道
  • Source 与 Sink 均可连接 FlowMQ 或外部系统,按实际链路组合即可
  • 目标端若支持幂等写入,建议启用
  • 注意字段类型映射,避免写入时类型不兼容
  • Processors 按单条消息处理;复杂关联或较重业务逻辑宜放在前置系统
  • MQTT / NATS 等:发布目标勿选择会再次进入本管道 Source 的主题,以免形成消息环路
  • Redis 使用 simple 模式时,master 仍需填写非空值(如 mymaster);地址中的库号字段名为 database,勿使用 db
  • TDengine / PostgreSQL / MySQL:ISO 时间写入 TIMESTAMP 请用 ts_parse;Unix 毫秒写入 BIGINT