Skip to main content

概述

Twitter Source Connector 通过 Twitter4J 连接旧版 Twitter Streaming API,根据关键词和可选的用户 ID 持续接收实时状态,并将每条状态转换为结构化 Kafka Connect 记录写入指定 Kafka Topic。它用于把账号获准访问的实时状态流接入 Kafka,供流处理、检索或归档系统消费。 正常状态写入 kafka.status.topic。启用删除通知处理后,Connector 也会把删除通知作为同一 Topic 中的 null value 记录写出。该 Connector 只处理连接建立后的实时回调,不提供历史搜索、存量回补或可恢复的上游读取位置。

前置条件

  • Kafka Connect 3.9.1 不能直接使用官方 0.3.34 ZIP。该发布包缺少运行时所需的 Guava;即使补入其依赖管理声明的 Guava 30.1.1-jre,包内 connect-utils 0.7.171 仍会在访问 Kafka 3.9.1 的 SystemTime 时触发 IllegalAccessError,导致 Task 无法启动。只有取得由 Connector 供应方提供并确认兼容 Kafka Connect 3.9.1 的正式制品后,才能继续部署。
  • Connector 使用 Twitter4J 4.0.6 对接旧版 statuses/filter.json 流式接口;部署前必须向当前 X/Twitter API 服务确认该接口仍对目标账号、产品层级和地区开放,不能仅凭已有 OAuth 凭证推断可用。
  • 准备一组已获准访问上述流式接口的 OAuth 1.0a consumer key、consumer secret、access token 和 access token secret,并确认其授权范围允许读取所需实时数据。
  • 明确要跟踪的关键词;filter.keywords 是必填项,仅配置用户 ID 不会创建采集 Task。

授权许可

使用 Apache License 2.0。

快速开始

以下配置仅供参考,不能证明官方 0.3.34 ZIP 可在 Kafka Connect 3.9.1 中运行。必须先取得供应方确认兼容 Kafka Connect 3.9.1 的正式制品;不要自行补入依赖或修改二进制文件来绕过兼容性错误。之后准备 Connect Cluster、Kafka、接收状态记录的 Topic,以及已确认可访问旧版流式接口的 X/Twitter 账号和 OAuth 凭证,并确认网络连通与访问权限。具体准备和管理操作请参阅管理 Connector
替换关键词、目标 Topic 和四项 OAuth 占位符。通过受控的配置管理方式提供凭证,不要将其提交到代码仓库或写入日志。示例继承 Worker 的 key/value Converter;所选 Converter 必须能序列化嵌套 Struct、数组、Map 和 Timestamp。提交配置前,仍需向 X/Twitter API 服务确认目标账号当前具有旧版流式接口权限;配置校验通过不代表外部实时流可用。

配置

身份认证

twitter.oauth.consumerKey

Twitter OAuth consumer key。
  • 类型password
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用与其他三项 OAuth 配置属于同一应用和授权关系的值;凭证有效性由外部 API 校验。
  • 必填:是

twitter.oauth.consumerSecret

Twitter OAuth consumer secret。
  • 类型password
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:作为敏感配置管理,不写入日志、工单或共享示例。
  • 必填:是

twitter.oauth.accessToken

Twitter OAuth access token。
  • 类型password
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:必须对应有权访问旧版流式接口的账号和应用;Connector 的配置校验不会确认实际权限。
  • 必填:是

twitter.oauth.accessTokenSecret

Twitter OAuth access token secret。
  • 类型password
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:作为敏感配置管理,并与对应的 access token 配套使用。
  • 必填:是

采集范围

filter.keywords

提交给流式过滤接口的跟踪关键词列表。
  • 类型list
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用逗号分隔关键词。至少提供一个有效关键词;Connector 按关键词数量和 tasks.max 划分 Task。
  • 必填:是

filter.userIds

提交给流式过滤接口的跟踪用户 ID 列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:高
  • 有效值 / 注意事项:使用逗号分隔十进制数字用户 ID,而不是用户名;每项必须能解析为 Java long。重复 ID 会被合并。该配置不能替代必填的 filter.keywords

Kafka 目标与删除通知

kafka.status.topic

接收正常状态和已启用删除通知的 Kafka Topic。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:填写一个明确的 Topic 名称。该版本不支持单独的 kafka.delete.topic
  • 必填:是

process.deletes

控制是否处理 Twitter 删除通知。
  • 类型boolean
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项false 忽略删除通知;true 将删除通知写入 kafka.status.topic,记录 key 使用仅含 StatusIdStatusDeletionNoticeKey,value 为 null。
  • 必填:是

队列与批次

queue.empty.ms

内部队列为空时等待记录的时间,单位为毫秒。
  • 类型int
  • 默认值100
  • 重要级别:低
  • 有效值 / 注意事项:最小值为 10。该值控制空队列轮询等待,不是 Twitter API 连接超时。

queue.batch.size

一次 poll 从内部队列返回的目标批量大小。
  • 类型int
  • 默认值100
  • 重要级别:低
  • 有效值 / 注意事项:最小值为 1。该值控制 Connector 内部队列批次,不是 Twitter API 拉取批量。

Twitter 客户端诊断

twitter.debug

启用 Twitter4J 调试日志。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:仅在受控排障期间临时启用;输出可能较多,且不会改变 Kafka Connect Worker 的日志级别配置。

Connector 身份与任务

connector.class

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

tasks.max

Connector 允许创建的最大 Task 数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:最小值为 1。实际 Task 数量为 tasks.max 与关键词数量中的较小值;提高该值会建立更多独立外部流连接。

tasks.max.enforce

控制 Kafka Connect 是否强制执行 tasks.max 上限。
  • 类型boolean
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:Kafka Connect 3.9.1 已将该配置标记为弃用并计划移除;保持 true,通过 tasks.max 管理任务数。
  • 弃用:是

记录转换

key.converter

覆盖 Worker 用于序列化 SourceRecord key 的 Converter。
  • 类型class
  • 默认值null,继承 Worker 配置
  • 重要级别:低
  • 有效值 / 注意事项:如显式设置,类必须实现 Kafka Connect Converter 并可被实例化。状态和删除通知使用结构化 key。

value.converter

覆盖 Worker 用于序列化 SourceRecord value 的 Converter。
  • 类型class
  • 默认值null,继承 Worker 配置
  • 重要级别:低
  • 有效值 / 注意事项:如显式设置,类必须实现 Kafka Connect Converter 并可被实例化。正常状态包含嵌套结构、数组、Map 和 Timestamp;删除通知的 value 为 null。

最佳实践

调整持续采集的数据范围

适用业务场景:Connector 已完成首次接入,需要随着业务主题变化增加或收窄关键词,并可同时关注一组明确的数字用户 ID,使后续实时数据进入原有目标 Topic。 在快速开始配置中替换或添加以下属性:
关键说明:filter.userIds 只能填写数字 ID,并且不能单独驱动任务创建,因此始终保留至少一个关键词。变更过滤条件会重建实时流;Connector 没有历史回补或断点续传能力,应把切换窗口可能出现的缺口或重复纳入下游处理设计。

为多个关键词增加采集并行度

适用业务场景:现有单 Task 已持续运行,关键词集合扩大后需要把关键词分配给多个独立 Task 和流连接,以降低单个连接承担的过滤范围。 在快速开始配置中替换或添加以下属性:
关键说明:Connector 创建的 Task 数量不会超过关键词数量,并按迭代顺序轮询分配关键词。增加 tasks.max 会建立更多外部流连接;如果同时配置 filter.userIds,完整用户 ID 列表会复制到每个 Task,可能产生重复记录。Connector 不保证跨 Task 全局顺序或去重,下游应使用稳定业务标识处理重复。

保留上游删除信号

适用业务场景:持续采集已经运行,下游归档、索引或合规处理需要感知上游发出的状态删除通知,而不是静默忽略它们。 在快速开始配置中替换以下属性:
关键说明:删除通知写入与正常状态相同的 kafka.status.topic,记录 key 使用仅含 StatusIdStatusDeletionNoticeKey,value 为 null;该版本没有独立删除 Topic。下游必须显式识别这种记录形态。若依赖 Kafka 日志压缩执行删除,还需确认该 key 的最终序列化结果、分区方式和目标 Topic 配置符合下游语义。

监控

监控内容

关注 Kafka Connect 通用健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;仅在部署启用了相应错误处理时关注 DLQ 活动。该 Connector 不保存可恢复的 Twitter 上游游标,因此 Task 处于 RUNNING 状态并不表示外部实时流完整;还应检查认证、连接、stall warning 和上游限流相关日志。

导入 Grafana 大盘

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

限制条件

  • 官方 0.3.34 ZIP 缺少 Guava,Kafka Connect 3.9.1 无法完成插件发现,因此该原始发布包不能直接用于该 Worker 版本。
  • 为 0.3.34 补齐 Guava 30.1.1-jre 后,connect-utils 0.7.171 仍会在 Task 启动时访问 Kafka 3.9.1 的 SystemTime 并触发 IllegalAccessError;失败发生在任何 X/Twitter 请求之前,Task 无法运行。
  • Connector 依赖 Twitter4J 4.0.6 使用的旧版 statuses/filter.json 流式接口;其当前可用性、产品层级和账号授权由 X/Twitter API 服务决定,必须由用户在部署前确认。
  • Connector 只接收连接建立后的实时状态和删除回调,不支持历史搜索、存量回补或指定时间范围读取。
  • Connector 不保存可用于恢复的上游游标;Task 重启、迁移或过滤条件变更后会重新建立实时连接,无法精确续传或回补中断窗口。
  • Connector 不保证跨 Task 全局顺序、幂等写入或去重;多 Task、过滤条件重叠和重连都可能带来重复或缺口。
  • 所有正常状态写入一个 kafka.status.topic,不按关键词、用户或事件类型路由到不同 Topic,也不提供原始 JSON 透传模式。
  • 删除通知只能作为 kafka.status.topic 中的 null value 记录写出;该版本不支持 kafka.delete.topic
  • 内部队列在实际负载下近似无界,也没有上游暂停机制;Kafka 写入或序列化持续变慢时,积压可能增加 Worker 内存压力。
  • Connector 内部捕获的单条转换或入队异常不会交给 Kafka Connect 的 DLQ 处理,也不会自动重试该事件。

常见问题

为什么官方 0.3.34 ZIP 在 Kafka Connect 3.9.1 中无法启动 Task

原始发布包缺少 Guava,Worker 会在插件发现阶段报告 com.google.common.collect.Multimap 缺失。仅补入 Guava 30.1.1-jre 后,插件虽然可以被发现并接受配置,但 Task 会因 connect-utils 0.7.171 访问 Kafka 3.9.1 的 SystemTime 而抛出 IllegalAccessError。该错误发生在创建 X/Twitter 客户端和发出外部请求之前,与 OAuth 凭证或外部 API 响应无关。不要把自行补包或修改二进制文件视为受支持方案;在供应方提供并确认兼容 Kafka Connect 3.9.1 的正式制品前,本页示例只能作为配置参考。

Connector 无法建立 Twitter 流连接

常见原因包括旧版流式接口已不对目标账号开放、OAuth 1.0a 四项凭证不属于同一授权关系、Token 权限不足或外部服务拒绝当前产品层级。先向 X/Twitter API 服务确认 statuses/filter.json 的当前可用性和账号授权,再核对四项凭证及 Worker 到外部服务的访问路径。不要通过在日志中打印完整凭证排障。

Connector 已启动但没有 Task 或没有记录

如果 filter.keywords 为空,仅配置 filter.userIds 不会生成采集 Task。确认至少存在一个有效关键词,检查 Connector 与 Task 状态和认证日志,并确认关键词或数字用户 ID 在当前流式接口的服务端规则内。流式接口只返回连接建立后的匹配事件,不会回补此前数据。

扩容后出现重复记录或顺序变化

每个 Task 使用独立流连接,关键词会分片,而 filter.userIds 会完整复制到每个 Task;同一事件可能因过滤范围重叠或重连而重复,跨 Task 也没有全局顺序保证。检查 tasks.max、关键词分组和用户 ID 配置,并在下游依据稳定状态 ID 去重,不要把 Kafka Offset 当作 Twitter 上游游标。

启用删除处理后为什么没有单独的删除 Topic

该版本没有 kafka.delete.topic 配置。process.deletes=true 时,删除通知以带删除 key、null value 的记录写入 kafka.status.topic。检查下游 Converter 和消费者是否保留并识别 null value,同时确认目标 Topic 的压缩策略是否符合预期。