Skip to main content

概述

IRC Source Connector 连接 IRC 服务器并加入指定频道,将收到的 PRIVMSG 写入一个 Kafka Topic,可用于把公开频道或受控内部频道中的聊天消息接入流处理、检索和归档链路。Kafka 消息的 key 是消息目标(通常为频道名),value 包含接收时间、消息目标、发送者信息和消息文本。配置多个频道时,可由一个或多个 Source Task 分担采集工作。

前置条件

  • IRC 网络必须提供 Connect Worker 可访问的明文 TCP 端点,并允许所用机器人昵称完成注册;该 Connector 不支持 TLS,因此只应在可信或隔离的网络中使用。
  • 目标频道必须允许所用昵称直接加入;该 Connector 不能为需要频道密钥的频道提供 JOIN 密钥。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka、目标 Kafka Topic 和可访问的 IRC 服务器及频道,并确认网络连通和访问权限。具体准备和管理操作请参阅 管理 Connector
<irc-server> 替换为 IRC 服务器主机名或地址,将 <irc-channel> 替换为包含频道前缀的频道名(例如 #events),将 <kafka-topic> 替换为接收消息的 Kafka Topic。示例使用明文 IRC 的默认端口 6667;更换端口不会启用 TLS。省略机器人名称时,Connector 会动态生成昵称。

配置

IRC 连接

irc.server

IRC 服务器的主机名或地址。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:必须是 Connect Worker 可访问的非空地址,并与 irc.server.port 一起指向明文 IRC 服务。
  • 必填:是

irc.server.port

连接 IRC 服务器所用的 TCP 端口。
  • 类型int
  • 默认值6667
  • 重要级别:低
  • 有效值 / 注意事项:使用目标 IRC 服务实际监听的有效端口。修改端口不会启用 TLS,也不能用于安全连接仅支持 TLS 的端口。

IRC 身份与认证

irc.bot.name

Source Task 向 IRC 服务器注册时使用的昵称、用户名和真实名称。
  • 类型string
  • 默认值:动态生成,格式为 KafkaConnectBot_ 加 6 位随机字母或数字
  • 重要级别:低
  • 有效值 / 注意事项:必须符合目标 IRC 网络的昵称规则。省略时,每个 Task 会独立生成昵称;多 Task 场景不要显式复用同一个固定昵称。

irc.password

连接注册阶段发送给 IRC 服务器的可选服务器密码。
  • 类型password
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:空字符串表示不发送 PASS 命令。非空密码会通过明文 IRC 连接发送,不应在不可信网络中使用,也不要写入日志或版本库。

频道订阅与 Kafka 目标

irc.channels

Source Task 要加入并采集消息的 IRC 频道列表。
  • 类型list
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:至少提供一个非空频道;多个频道用逗号分隔,并使用目标 IRC 网络接受的频道前缀(通常为 #)。空元素、重复频道和频道语法不会由 Connector 提前拒绝。
  • 必填:是

kafka.topic

接收所有 IRC 消息的单个 Kafka Topic。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:填写一个有效的 Kafka Topic 名称;该值不是列表,也不支持按频道生成不同 Topic。
  • 必填:是

Kafka Connect 运行

connector.class

要加载的 IRC Source Connector 实现类。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用 com.github.cjmatta.kafka.connect.irc.IrcSourceConnector
  • 必填:是

tasks.max

用于分配已配置 IRC 频道的最大 Source Task 数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:最小值为 1。实际 Task 数量为频道列表元素数量与 tasks.max 中的较小值;每个 Task 建立独立 IRC 连接。

key.converter

用于序列化字符串类型消息 key 的 Connector 级 Converter 覆盖项。
  • 类型class
  • 默认值null(继承 Worker 配置)
  • 重要级别:低
  • 有效值 / 注意事项:填写可实例化的 Kafka Connect Converter 实现类;消息 key 是 PRIVMSG 的目标,通常为频道名。

value.converter

用于序列化带 Schema 的 IRC 消息 Struct 的 Connector 级 Converter 覆盖项。
  • 类型class
  • 默认值null(继承 Worker 配置)
  • 重要级别:低
  • 有效值 / 注意事项:填写可实例化且能处理带 Schema Struct 的 Kafka Connect Converter 实现类。消息 value 包含 createdatchannelsendermessage 字段。

最佳实践

扩展到多个频道并分配独立连接

适用业务场景:完成单频道接入后,需要把同一 IRC 服务器上的多个频道汇聚到一个 Kafka Topic,并由不同 Task 分担频道连接。这适用于频道数量增加后的常规扩展,但不会在单个频道内并行采集。 配置示例:在快速开始配置的基础上,将频道列表扩展为三个频道,并把最大 Task 数调整为 3。省略 irc.bot.name,使每个 Task 独立生成昵称,避免多个连接显式复用同一个固定昵称。
关键说明:实际 Task 数不会超过频道数;本例中,每个频道最多分配给一个独立 Task 和 IRC 连接。所有频道的消息仍写入同一个 kafka.topic,消息 key 标识 PRIVMSG 的目标。调整 tasks.max 会重建频道分组和连接。由于变更期间没有历史回放,也不会在旧、新 Task 之间接力消息,应在可接受短暂采集缺口的维护窗口内实施。

监控

监控内容

关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;仅在启用了相应错误处理时关注 DLQ 活动。还应结合消息吞吐和 IRC 服务端连接信息判断采集是否有效,因为运行期断线不一定会使 Task 状态变为失败。

导入 Grafana 大盘

确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。

限制条件

  • 仅采集 IRC PRIVMSGJOINPARTNOTICETOPICKICKQUIT 等事件不会写入 Kafka,直接发送给机器人昵称的私聊消息也可能被采集。
  • 只支持明文 IRC TCP 连接,不支持 TLS、STARTTLS 或证书配置;修改端口不能提供传输加密,非空服务器密码也会以明文发送。
  • 不支持为受频道密钥保护的频道提供 JOIN 密钥。
  • 每条记录不包含可恢复的 IRC Source Offset;重启后只能采集重新加入频道之后的新消息,停机期间或内存中尚未交付的消息无法从源端回放。
  • Connector 不提供 exactly-once、at-least-once、无重复或无丢失保障。
  • 运行期断线、昵称冲突和部分 IRC 服务端错误不会自动触发重连、重新加入频道或 Task 失败;Task 可能仍显示运行中,但不再产出消息。
  • 入站消息使用无容量上限的内存队列,且单次 poll 会排空当时的全部积压;持续高于 Kafka 写入能力的入站速率可能增加 Worker 内存压力。
  • 多 Task 显式使用同一个固定 irc.bot.name 时可能发生昵称冲突;Connector 不会为各 Task 派生唯一昵称或自动处理冲突。
  • Task 停止或重配时等待 IRC 读取线程退出且没有有界超时,可能延迟停止过程,并清空尚未交付的内存消息。

常见问题

Task 显示运行中,但 Kafka Topic 不再收到新消息

运行期 IRC 断线、昵称冲突、密码错误或频道拒绝不一定会使 Kafka Connect Task 进入失败状态,并且 Connector 不会自动重连。检查 Worker 日志、IRC 服务端连接记录、昵称占用情况和频道访问规则;确认配置有效后,重启 Connector 或 Task,并重新检查消息吞吐。重启期间的 IRC 消息无法回放。

重启后为什么没有补采停机期间的消息

该 Connector 不记录可恢复的 IRC 消息位置,也不向 IRC 服务器请求历史消息。它在重启后重新连接并加入频道,只接收此后到达的 PRIVMSG。如业务需要完整历史,应由 IRC 侧提供独立历史来源,或在下游设计中明确接受并监控采集缺口。