跳到主要内容
版本:Next

ActiveMQ

ActiveMQ Sink 连接器

描述

用于把 SeaTunnel 数据写入 ActiveMQ 队列。每一行数据都会被序列化成一条 JSON 文本消息。这个连接器只支持 Sink,SeaTunnel 目前没有提供 ActiveMQ Source 连接器。

关键特性

选项

名称类型是否必传默认值描述
uristring-ActiveMQ Broker 地址,例如 tcp://localhost:61616
queue_namestring-要写入的队列名称。
usernamestring-创建 ActiveMQ 连接时使用的用户名。配置该项时必须同时配置 password
passwordstring-创建 ActiveMQ 连接时使用的密码。配置该项时必须同时配置 username
client_idstring-连接工厂使用的 JMS 客户端 ID。
check_for_duplicateboolean-是否让 ActiveMQ 客户端检查重复消息。
always_session_asyncboolean-是否始终为每个 Session 使用独立线程分发消息。
always_sync_sendboolean-是否始终使用同步方式发送消息。
close_timeoutint-关闭连接的超时时间,单位毫秒。
dispatch_asyncboolean-Broker 是否异步分发消息。
nested_map_and_list_enabledboolean-是否允许结构化消息属性和 MapMessage 条目中包含嵌套的 MapList 对象。
warn_about_unstarted_connection_timeoutint-连接没有正确启动时,ActiveMQ 客户端发出警告前等待的毫秒数。设置为小于 0 的值可以关闭这个警告。
consumer_expiry_check_enabledboolean-是否在每个 MessageConsumer 分发消息前检查消息是否已经过期。

注意事项

  • uri 是连接入口,Broker 的主机和端口都写在这里,例如 tcp://activemq-host:61616
  • usernamepassword 是可选项;如果 Broker 需要认证,这两个配置要一起写。
  • 连接器会把每一行 SeaTunnel 数据作为一条 JSON 文本消息写入 queue_name,当前没有单独的 format 配置。
  • Broker 地址请使用 uri 配置,hostport 不是 ActiveMQ Sink 的配置项。
  • 该 Sink 前面可以接任意 SeaTunnel Source。ActiveMQ 连接器只负责把最终的数据行发送到队列。

示例

把 FakeSource 数据写入 ActiveMQ 队列:

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

source {
FakeSource {
schema = {
fields {
id = int
name = string
}
}
rows = [
{ kind = INSERT, fields = [1, "Alice"] }
{ kind = INSERT, fields = [2, "Bob"] }
]
}
}

sink {
ActiveMQ {
uri = "tcp://localhost:61616"
username = "admin"
password = "admin"
queue_name = "testQueue"
}
}

在流式模式下,Sink 会保持与 Broker 的连接持续打开,每来一行数据就写入一条。用户名 / 密码也可以直接 写在 uri 中,例如 tcp://admin:admin@localhost:61616

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

source {
FakeSource {
schema = {
fields {
id = int
name = string
}
}
rows = [
{ kind = INSERT, fields = [1, "Alice"] }
]
}
}

sink {
ActiveMQ {
uri = "tcp://admin:admin@localhost:61616"
queue_name = "testQueue"
}
}

FAQ

ActiveMQ Sink 支持 topic 吗?

不支持。当前 Sink 只会写入由 queue_name 指定的 JMS 队列,连接器工厂并不支持 topic 目的地。如果需要发布/订阅语义,请改用通用的 JMS 连接器或者其它专门面向 ActiveMQ topic 的桥接组件;要直接驱动本连接器,请继续用 queue。

usernamepassword 是怎么校验的?

这两项都是可选的。设置的时候必须同时设置(只设其中一个会让任务启动失败)。它们作用于 JMS 连接工厂这一层,因此会覆盖已经嵌在 uri 里的凭证。如果 Broker 需要鉴权,建议显式配置 username/password 而不是把它们写在 URL 里——这样凭证会出现在任务配置日志里,而不是埋在连接串里。

每条数据最终变成什么格式的消息?

每行 SeaTunnel 数据会被序列化为一条 JSON 文本消息写到配置的 queue_name 里。连接器没有 format 配置项,JSON 结构由 Sink 序列化器固定,所以对端如果想要其它编码,需要自己先解码 JSON body。

支持精确一次写入吗?

不支持。Sink 是尽力而为模式,重连行为由底层 JMS 客户端决定。如果可以接受来自上游的至少一次重放,请在任务级别启用 checkpoint。

变更日志