概述
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 包含createdat、channel、sender和message字段。
最佳实践
扩展到多个频道并分配独立连接
适用业务场景:完成单频道接入后,需要把同一 IRC 服务器上的多个频道汇聚到一个 Kafka Topic,并由不同 Task 分担频道连接。这适用于频道数量增加后的常规扩展,但不会在单个频道内并行采集。 配置示例:在快速开始配置的基础上,将频道列表扩展为三个频道,并把最大 Task 数调整为3。省略 irc.bot.name,使每个 Task 独立生成昵称,避免多个连接显式复用同一个固定昵称。
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
PRIVMSG;JOIN、PART、NOTICE、TOPIC、KICK和QUIT等事件不会写入 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 侧提供独立历史来源,或在下游设计中明确接受并监控采集缺口。