概述
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-utils0.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。配置
身份认证
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 使用仅含StatusId的StatusDeletionNoticeKey,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 和流连接,以降低单个连接承担的过滤范围。 在快速开始配置中替换或添加以下属性:tasks.max 会建立更多外部流连接;如果同时配置 filter.userIds,完整用户 ID 列表会复制到每个 Task,可能产生重复记录。Connector 不保证跨 Task 全局顺序或去重,下游应使用稳定业务标识处理重复。
保留上游删除信号
适用业务场景:持续采集已经运行,下游归档、索引或合规处理需要感知上游发出的状态删除通知,而不是静默忽略它们。 在快速开始配置中替换以下属性:kafka.status.topic,记录 key 使用仅含 StatusId 的 StatusDeletionNoticeKey,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-utils0.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 的压缩策略是否符合预期。