Skip to content

RabbitMQ Source 集成

本文介绍如何在管道的 Source 步骤中使用 RabbitMQ 作为数据来源。

适用场景

  • 接入已有 AMQP 生产链路数据
  • 从交换机/队列读取事件并汇入 FlowMQ

前提条件

  • 已了解管道结构与创建流程,参考 数据集成;Source 连接器总览见 Source
  • RabbitMQ 服务可访问
  • 已具备目标 vhost 与队列访问权限

配置步骤

  1. 在“数据管道”中创建管道,完成基础信息配置。
  2. 在 Source 步骤选择 RabbitMQ
  3. 配置 urlsqueueconsumer_tag 等参数。
  4. 点击“测试”验证连接。
  5. 继续完成 Processors、Sink 与确认步骤。

配置示例

可参考以下参数示例:

参数建议值
host127.0.0.1
port5672
usernameguest
passwordguest
vhost/
queueflowmq.source.rabbitmq
consumer_tagflowmq-source-1
queue_declare.enabledtrue
queue_declare.durablefalse
auto_acktrue
tls.enabledfalse

必要表单参数

以表单中带 * 的字段为必填:

参数说明
urls*连接地址列表(每项包含 host/port,可含用户名、密码、vhost)
queue*要消费的 AMQP 队列名
consumer_tag*消费者标识。用于在 RabbitMQ 侧区分本管道的消费连接;控制台要求必填,建议填可读且唯一的名称(如 flowmq-source-1
auto_ack*是否在消费时自动确认。开启可提高吞吐,但会削弱投递保证
prefetch_count*未确认消息的最大条数
prefetch_size*未确认消息的最大字节数(0 表示不按字节限制)

可选高级项:

参数说明
tls自定义 TLS
queue_declare被动声明目标队列
bindings_declare被动声明队列绑定

参数建议

  • 优先使用独立消费队列避免互相干扰
  • consumer_tag 同一连接上不要与其他消费者重复
  • 按业务类型规划 routing key

注意事项

  • 明确 ack 策略,避免消息堆积
  • 关注重试与死信策略
  • 接入后的字段清洗与映射见 Processors
  • 控制台 urls 按 host/port 对象填写(可含 username、password、vhost);Bento 按该对象格式解析,不要改成 amqp://... 字符串
  • 使用 FlowMQ 内置 AMQP(--class=amqp)时,建议开启 queue_declare.enabled,并将 queue_declare.durable 设为 false(内置服务不支持 durable queue)
  • 队列名不要与已有 Kafka Topic / FlowMQ 流同名(例如不要用 flowmq.mqtt.kafka),否则 queue_declare 可能返回 406
  • 非 TLS 场景下请将 tls.enabled 保持为 false;若误开 TLS,连通性测试可能仍通过,但管道任务无法正常运行
  • Sink 为 Drop 的 Source 测试场景,建议开启 auto_ack,便于观察处理计数