概述
RSS Source Connector 定期读取一个或多个 RSS 或 Atom Feed,并将 Feed 中的条目写入同一个 Kafka Topic。每个可处理的条目对应一条 Kafka 记录,记录值包含条目标题、标识、链接、description 内容、作者、日期,以及 Feed 的 URL 和标题;记录键为空。
Connector 每次轮询当前 Feed 内容,并使用为各 URL 保存的条目指纹避免重复发送仍可识别的旧条目。指纹由条目的标题、链接、标识、description 内容和作者生成,不包含日期或 Feed 标题。它适合把新闻、博客、公告或其他 Feed 更新接入 Kafka,供下游检索、通知、分析或流处理应用消费。
授权许可
使用 MIT License。快速开始
提前准备 Connect Cluster、Kafka 和可访问的 RSS 或 Atom Feed,并确认网络连通和访问权限。具体准备和管理操作请参阅 管理 Connector。<rss-feed-url> 替换为 Worker 可访问的 RSS 或 Atom URL,将 <topic-name> 替换为接收记录的 Kafka Topic。URL 中的空格必须进行百分号编码。该示例使用带 Schema 的 JSON 输出,以保留 Connector 定义的字段结构;提交配置后,首次轮询会读取 Feed 当前包含的条目,后续轮询只输出尚未记录或指纹发生变化的条目。仅日期或 Feed 标题变化不会生成新指纹。
配置
Feed 输入
rss.urls
要轮询的 RSS 或 Atom Feed URL。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:至少提供一个有效 URL;多个 URL 使用单个空格分隔。URL 自身包含的空格必须进行百分号编码,不要使用重复空格、制表符或换行符作为分隔符。应保持 URL 唯一;同一 Task 收到重复 URL 时会因初始化冲突而失败。URL 不应包含凭据、访问令牌或签名查询参数。
- 必填:是
Kafka 输出
topic
接收所有 Feed 条目的 Kafka Topic。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:使用有效且可写的 Kafka Topic 名称。一个 Connector 实例中的所有 Feed 都写入该 Topic,不支持按 Feed 分别配置 Topic。
- 必填:是
轮询与任务
sleep.seconds
首次轮询之后,每个 Task 开始下一轮轮询前等待的秒数。
- 类型:
int - 默认值:
60 - 重要级别:中
- 有效值 / 注意事项:必须大于或等于
0。设置为0会取消轮询间等待,可能形成高频循环;实际轮询周期还包含该 Task 对所有 Feed 的网络读取和解析时间。
tasks.max
该 Connector 允许创建的最大 Task 数量。
- 类型:
int - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:必须大于或等于
1。实际 Task 数不会超过rss.urls中的 URL 数量;多个 URL 会按输入顺序轮转分配给各 Task,同一 Task 内的 URL 仍然串行轮询。
Connector 与序列化
connector.class
要加载的 RSS Source Connector 实现类。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:使用
org.kaliy.kafka.connect.rss.RssSourceConnector。 - 必填:是
value.converter
在 Connector 级别指定记录值使用的 Converter。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:省略或使用
null时继承 Worker 的值 Converter。显式配置时必须是可实例化的 Kafka ConnectConverter实现,并按所选 Converter 配置其 Schema 或序列化选项;Connector 输出固定的 Connect Schema 和 Struct 值。
最佳实践
按 Feed 数量扩展轮询任务
适用业务场景:一个 Connector 需要持续读取多个彼此独立的 Feed,单个 Task 串行抓取所有 URL 已使更新延迟增加。此时可按 URL 数量逐步增加 Task,让不同 Feed 分配到不同 Task 并行轮询。 配置示例: 保留快速开始中的 Connector、Topic 和 Converter 配置,将rss.urls 扩展为唯一 URL 列表,并增加或覆盖:
tasks.max 与 URL 数量中的较小值;示例中的 5 个 URL 会轮转分配给最多 3 个 Task。一个 URL 不会拆分给多个 Task,同一 Task 内的 Feed 仍按顺序抓取,因此响应缓慢的 Feed 会延迟该 Task 中后续 Feed。应从较小并行度开始,根据 Task 延迟、源站承载能力和 Worker 资源逐步调整。还应确保 URL 列表没有重复项,避免重复 URL 被分到同一 Task 后导致初始化失败。
按更新频率安排轮询并保留足够内容
适用业务场景:Feed 会持续滚动更新,需要在源站请求频率、消息新鲜度和停机恢复范围之间取得平衡。此时应根据业务可接受延迟设置轮询间隔,并让上游 Feed 保留足够长的条目窗口。 配置示例: 在快速开始配置中增加或覆盖以下属性:300 秒只是示例。应根据业务可接受的发现延迟设置该值,并计入每轮抓取和解析所需的时间。较短间隔可更快发现新条目,但会增加源站和 Worker 负载;较长间隔会扩大条目在两次轮询之间出现后又滚出 Feed 的风险。Connector 的 Offset 是当前 Feed 中条目指纹的窗口,不是源站历史游标,因此还应在上游保留足够多的条目,并在维护或长时间停机后检查是否存在已滚出 Feed 的内容。
监控
监控内容
关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;仅在启用了相应错误处理时关注 DLQ 活动。对于 RSS 抓取,还应结合 Task 日志区分“Feed 没有新条目”和“URL 读取或 XML 解析失败”。导入 Grafana 大盘
确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。限制条件
- Connector 只读取每次轮询时 Feed 中仍可见的当前条目,没有可回放的上游日志;轮询间隙或停机期间已经滚出 Feed 的条目无法补取。
- Kafka 记录发送成功与 Source Offset 持久化之间存在重复窗口,而滚动 Feed 又存在不可恢复窗口,因此不能保证恰好一次、无重复、无遗漏或无条件的至少一次交付。
- 条目必须包含可用的
title、id和link。缺少这些必填字段的条目无法构造固定 Schema,会被记录异常后过滤,且不会进入 Connector 自有 DLQ。 - 版本 0.1.1 不提供 Connector 级 HTTP 认证、自定义请求头、条件请求、代理、连接超时或读取超时配置。不要把凭据或令牌放入
rss.urls,因为 Task 启动日志可能包含完整配置。 - 单次轮询会返回分配给该 Task 的所有 Feed 中全部新条目,没有 Connector 级记录数、字节数或内存上限。大型 Feed 或积压内容可能形成较大的单次批次。
- Kafka 记录键和目标 Partition 均未由 Connector 指定。Feed URL 只用于 Source Offset 分区,不提供固定的 Kafka Partition 路由,也不保证跨 Feed、跨 Task 或跨 Kafka Partition 的全局顺序。
content只映射 Feed 条目的 description;category、enclosure、comments、原始 XML 及其他扩展字段不会写入记录值。
常见问题
为什么 Task 在启动时失败?
检查rss.urls 是否至少包含一个有效 URL、多个 URL 之间是否只有一个空格、URL 中的空格是否已百分号编码,以及列表中是否存在重复 URL;同一 Task 收到重复 URL 时会发生初始化冲突。随后确认 topic 是有效且可写的 Kafka Topic,并检查 Task 日志中的配置校验或初始化异常。修正配置后重启失败的 Task。
为什么 Connector 显示运行中,但 Topic 中没有新记录?
Feed 可能没有新条目,也可能在本轮发生 URL 读取或 XML 解析失败。检查 Worker 是否能访问 Feed、响应内容是否仍是有效的 RSS 或 Atom XML,并查看 Task 日志中的抓取警告。若 Feed 有内容,再确认条目包含title、id 和 link;缺少这些字段的条目会被过滤。修复源站内容或网络问题后,Connector 会在后续轮询中再次读取该 URL。
为什么重启后出现重复记录或缺少停机期间的条目?
Kafka 已确认记录但对应 Source Offset 尚未持久化时发生故障,重启后仍在 Feed 中的条目可能再次输出。另一方面,Connector 只能查看当前 Feed 快照;停机期间已经滚出 Feed 的条目无法恢复。下游应使用稳定业务字段进行幂等处理,并根据最长维护窗口调整上游 Feed 的条目保留范围。为什么增加 tasks.max 后 Task 数量没有继续增加?
一个 Feed URL 最多分配给一个 Task,因此实际 Task 数不会超过 URL 数量。确认 rss.urls 中包含足够多且唯一的 URL,并检查 Connector 状态中的实际 Task 数。若多个 URL 仍位于同一 Task,可继续增加 tasks.max,但超过 URL 数量后不会获得额外并行度。
为什么输出字段与 Feed 原始 XML 不完全一致?
Connector 使用固定字段模型而不是转发原始 XML。content 只取条目 description,date 优先取更新时间、其次取发布时间,并作为 ISO-8601 字符串写入记录值,而不是 Kafka 记录时间戳。需要 category、enclosure、原始 XML 或其他扩展字段时,应在上游转换 Feed,或在下游使用能够保留这些字段的其他采集方式。