跳到主要内容
版本:Next

GooglePubSub

Google Pub/Sub Sink 连接器

描述

将每条 SeaTunnel 输入行作为一条消息发布到 Google Pub/Sub 主题。

支持这些引擎

Spark
Flink
SeaTunnel Zeta

主要特性

参数

参数名类型是否必填默认值
project_idstring-
topicstring-
credentials_pathstring-
emulator_hoststring-
formatenumjson
field_delimiterstring,
common-options-

project_id [string]

目标主题所属的 Google Cloud 项目 ID。

topic [string]

目标 Pub/Sub 主题 ID。启动作业前必须先创建该主题。

credentials_path [string]

Google Cloud 服务账号 JSON 密钥文件的路径。未配置时,连接器使用 Application Default Credentials

emulator_host [string]

Pub/Sub 模拟器的主机和端口,例如 pubsub-emulator:8085。配置后,连接器使用无凭证的明文连接。生产环境中不要使用该选项。

format [enum]

消息负载格式。支持以下值:

  • json:将行写为 JSON 对象。
  • text:使用 field_delimiter 拼接行字段。

field_delimiter [string]

format = text 时使用的字段分隔符。默认值为 ,

common options

Sink 插件通用参数,请参考 Sink 通用选项

交付语义

连接器使用 Google Pub/Sub Publisher 的批处理、流量控制和重试机制。在检查点和关闭期间,连接器会等待所有已接收的发布操作完成;异步发布失败会使任务失败。

从 SeaTunnel 作业角度看,Pub/Sub 发布语义为至少一次。任务重试可能再次发布消息,因此下游消费者应根据业务需要处理重复消息。

当前版本只发布序列化后的行负载,不支持消息属性、排序键和按行选择主题。

任务示例

Application Default Credentials

env {
parallelism = 1
job.mode = "BATCH"
}

source {
FakeSource {
row.num = 10
schema = {
fields {
event_id = string
event_type = string
}
}
}
}

sink {
GooglePubSub {
project_id = "my-gcp-project"
topic = "events"
format = json
}
}

服务账号密钥文件

sink {
GooglePubSub {
project_id = "my-gcp-project"
topic = "events"
credentials_path = "/secrets/service-account.json"
format = text
field_delimiter = "|"
}
}

Pub/Sub 模拟器

sink {
GooglePubSub {
project_id = "local-project"
topic = "events"
emulator_host = "pubsub-emulator:8085"
}
}

Changelog

Change Log
ChangeCommitVersion
[Feature][Connector-V2] Add Google Pub/Sub sink connector-Next