> ## Documentation Index
> Fetch the complete documentation index at: https://docs.automq.com/llms.txt
> Use this file to discover all available pages before exploring further.

# Twitter Source Connector

> 介绍如何在 AutoMQ Connect 中配置 Twitter Source Connector，包括兼容性限制、OAuth、过滤条件、监控和故障排查。

## 概述

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](../manage-connectors)。

```properties theme={null}
connector.class=com.github.jcustenborder.kafka.connect.twitter.TwitterSourceConnector
twitter.oauth.consumerKey=<twitter-consumer-key>
twitter.oauth.consumerSecret=<twitter-consumer-secret>
twitter.oauth.accessToken=<twitter-access-token>
twitter.oauth.accessTokenSecret=<twitter-access-token-secret>
filter.keywords=<keyword>
kafka.status.topic=<status-topic>
process.deletes=false
```

替换关键词、目标 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 使用仅含 `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。

在快速开始配置中替换或添加以下属性：

```properties theme={null}
filter.keywords=<keyword-a>,<keyword-b>
filter.userIds=<numeric-user-id-a>,<numeric-user-id-b>
```

关键说明：`filter.userIds` 只能填写数字 ID，并且不能单独驱动任务创建，因此始终保留至少一个关键词。变更过滤条件会重建实时流；Connector 没有历史回补或断点续传能力，应把切换窗口可能出现的缺口或重复纳入下游处理设计。

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

适用业务场景：现有单 Task 已持续运行，关键词集合扩大后需要把关键词分配给多个独立 Task 和流连接，以降低单个连接承担的过滤范围。

在快速开始配置中替换或添加以下属性：

```properties theme={null}
tasks.max=2
filter.keywords=<keyword-a>,<keyword-b>,<keyword-c>,<keyword-d>
```

关键说明：Connector 创建的 Task 数量不会超过关键词数量，并按迭代顺序轮询分配关键词。增加 `tasks.max` 会建立更多外部流连接；如果同时配置 `filter.userIds`，完整用户 ID 列表会复制到每个 Task，可能产生重复记录。Connector 不保证跨 Task 全局顺序或去重，下游应使用稳定业务标识处理重复。

### 保留上游删除信号

适用业务场景：持续采集已经运行，下游归档、索引或合规处理需要感知上游发出的状态删除通知，而不是静默忽略它们。

在快速开始配置中替换以下属性：

```properties theme={null}
process.deletes=true
```

关键说明：删除通知写入与正常状态相同的 `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](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 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 的压缩策略是否符合预期。
