PostgreSQL Sink 集成
本文介绍如何将 FlowMQ 数据写入 PostgreSQL。
前提条件
配置步骤
- 在“数据管道”中创建管道,完成基础信息配置。
- 在 Source 步骤选择
Kafka(Topic:flowmq.mqtt.kafka),Processors 按需配置。 - 在 Sink 步骤选择
PostgreSQL,配置host、database、table、columns等参数。 - 点击“测试”验证连接。
- 继续完成确认步骤。
配置示例
以下示例将 Kafka Topic flowmq.mqtt.kafka 中的样例消息写入目标系统。样例 JSON 与 Source 约定见 Sink · 示例场景。
可参考以下参数示例:
| 参数 | 建议值 |
|---|---|
| host | 127.0.0.1 |
| port | 5432 |
| username | postgres |
| database | flowmq |
| table | mqtt_kafka_events |
目标表示例:
sql
CREATE TABLE IF NOT EXISTS mqtt_kafka_events (
id BIGSERIAL PRIMARY KEY,
source VARCHAR(255),
value_time BIGINT,
value REAL,
unit_code VARCHAR(64),
create_time TIMESTAMP NULL
);columns 字段映射示例(对应样例消息):
| name | data_type | is_primary_key | value |
|---|---|---|---|
| source | VARCHAR | false | ${! json("source") } |
| value_time | BIGINT | false | ${! json("data.timestamp") } |
| value | FLOAT | false | ${! json("data.value") } |
| unit_code | VARCHAR | false | ${! json("type") } |
| create_time | TIMESTAMP | false | ${! json("time").ts_parse("2006-01-02T15:04:05.999Z07:00") } |
data.timestamp为 Unix 毫秒,建议以BIGINT落库。写入TIMESTAMP列时请用ts_parse(或${! now() }),勿直接填入 ISO 字符串。
必要表单参数
参数名后带 * 表示控制台必填。「支持表达式」为「是」时,可使用 ${! ... } 按消息动态取值,详见 Processors。
| 参数 | 支持表达式 | 说明 |
|---|---|---|
host* | 否 | 数据库地址 |
port* | 否 | 数据库端口 |
username* | 否 | 数据库用户名 |
password* | 否 | 数据库密码 |
database* | 否 | 目标数据库 |
table* | 否 | 目标表 |
columns* | 是 | 列映射(name / data_type / is_primary_key / value) |
参数建议
- 合理设置事务与批量写入参数
columns映射中注意时间字段与时区处理
注意事项
- 合理设置事务与批量写入参数
- 注意时间字段与时区处理