跳到主要内容
版本:3.0.0

GooglePubSub

Google Pub/Sub Sink 连接器

描述​

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

支持这些引擎​

Spark
Flink
SeaTunnel Zeta

主要特性​

参数​

参数名类型是否必填默认值
project_idstring是-
topicstring是-
credentials_pathstring否-
emulator_hoststring否-
formatenum否json
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] [GooglePubSub] Add Google Pub/Sub source connector (#11989)https://github.com/apache/seatunnel/commit/c5b58851b3.0.0
[Feature][Connector-V2] [GooglePubSub] Add Google Pub/Sub sink connector (#11877)https://github.com/apache/seatunnel/commit/e97555a623.0.0
[Feature][Connector-V2] Add Google Pub/Sub source connector-Next
[Feature][Connector-V2] Add Google Pub/Sub sink connector-Next