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

# 实时监测

发现特定账号的新推文，及时追踪关键词提及，并将实时 X 数据输入应用。本指南介绍如何基于 Sorsa API，使用主动拉取的轮询模式构建近实时监测管道。

检测延迟取决于轮询间隔、API 响应时间，以及帖子何时进入所选信息流或搜索索引。设计时应围绕检查点、分页和去重，不要假设每条新帖子都会出现在下一次响应中。

> **免费构建原型：** 初始赠送的 100 次请求可访问所有 Sorsa 端点，一次性赠送，无需信用卡，永不过期。足以搭建下方任意监测器、确认能够发现实时推文，并在选购套餐前验证 Slack 或 Discord 分流。轮询消耗较多请求，请根据下方用量表和轮询间隔选择付费套餐。

> **注意：** 更多架构模式和端到端示例请参阅博客上的[使用 REST API 实时监测 Twitter](https://api.sorsa.io/blog/real-time-twitter-monitoring)。

***

## 轮询监测的工作方式

从社交平台获取数据有推送式（流式传输、webhook）和拉取式（轮询）两种方式。Sorsa API 使用轮询，分为四步：

1. **定期轮询端点**，间隔为 1–30 秒。
2. **比较结果**与之前见过的推文 ID，识别新内容。
3. **处理新推文：** 发送提醒、保存，或分发至 Slack、Discord 等。
4. **重复。**

推文 ID 编码了创建时间，可通过 Python 整数或 JavaScript BigInt 比较。最大已见 ID 可作为检查点，但帖子可能延迟出现或乱序返回。生产环境应重复检查一段重叠时间窗口并按 ID 去重。崩溃后恢复需要显式保存和加载检查点。

### 选择合适的端点

| 监测目标           | 端点               | 方法   | 选择原因                               |
| :------------- | :--------------- | :--- | :--------------------------------- |
| 单个账号           | `/user-tweets`   | POST | 返回该用户时间线中的最新推文                     |
| 同时最多 5,000 个账号 | `/list-tweets`   | GET  | 一次请求覆盖 X 列表中的全部成员                  |
| 关键词或话题标签       | `/search-tweets` | POST | 完整支持搜索运算符，使用 `order: latest` 按时间排序 |
| 账号的 @提及        | `/mentions`      | POST | 专为提及追踪设计，支持互动筛选                    |

***

## 第一级：监测单个账号

这是最简单的情况。循环轮询 `/user-tweets` 的第一页，输出 ID 大于上次已见 ID 的推文。首次成功轮询时只建立基线，不输出现有推文。这是开发示例：如果两次轮询之间或进程停止期间新增内容超过一页，可能漏掉帖子。

### Python

```python theme={null}
import requests
import time

API_KEY = "YOUR_API_KEY"
USERNAME = "elonmusk"
POLL_INTERVAL = 5  # seconds

URL = "https://api.sorsa.io/v3/user-tweets"
HEADERS = {"ApiKey": API_KEY, "Content-Type": "application/json"}

last_seen_id = None

print(f"Monitoring @{USERNAME}...")

while True:
    try:
        resp = requests.post(URL, headers=HEADERS, json={"username": USERNAME}, timeout=30)
        resp.raise_for_status()
        tweets = resp.json().get("tweets", [])

        if tweets:
            # Snowflake IDs arrive as strings. Use the highest (newest) ID in the
            # batch; this stays correct even if a pinned tweet appears first.
            top_id = max(int(t["id"]) for t in tweets)

            if last_seen_id is None:
                last_seen_id = top_id
                print(f"Baseline set: {last_seen_id}")
            else:
                new_tweets = [t for t in tweets if int(t["id"]) > last_seen_id]
                for tweet in reversed(new_tweets):  # oldest first
                    print(f"[NEW] @{USERNAME}: {tweet['full_text'][:140]}")
                if new_tweets:
                    last_seen_id = top_id

    except requests.exceptions.RequestException as e:
        print(f"Error: {e}")
        time.sleep(POLL_INTERVAL * 2)
        continue

    time.sleep(POLL_INTERVAL)
```

### JavaScript

```javascript theme={null}
const API_KEY = "YOUR_API_KEY";
const USERNAME = "elonmusk";
const POLL_INTERVAL = 5000;

let lastSeenId = null;
console.log(`Monitoring @${USERNAME}...`);

while (true) {
  try {
    const resp = await fetch("https://api.sorsa.io/v3/user-tweets", {
      method: "POST",
      headers: { "ApiKey": API_KEY, "Content-Type": "application/json" },
      body: JSON.stringify({ username: USERNAME }),
    });
    if (!resp.ok) throw new Error(`HTTP ${resp.status}`);

    const tweets = (await resp.json()).tweets || [];

    if (tweets.length > 0) {
      // BigInt avoids precision loss on 64-bit Snowflake IDs.
      // Use the highest ID in the batch (robust if a pinned tweet appears first).
      let topId = 0n;
      for (const t of tweets) {
        const id = BigInt(t.id);
        if (id > topId) topId = id;
      }

      if (lastSeenId === null) {
        lastSeenId = topId;
        console.log(`Baseline set: ${lastSeenId}`);
      } else {
        const newTweets = tweets.filter((t) => BigInt(t.id) > lastSeenId);
        for (const t of [...newTweets].reverse()) {
          console.log(`[NEW] @${USERNAME}: ${t.full_text.slice(0, 140)}`);
        }
        if (newTweets.length) lastSeenId = topId;
      }
    }
  } catch (err) {
    console.error(`Error: ${err.message}`);
    await new Promise((r) => setTimeout(r, POLL_INTERVAL * 2));
    continue;
  }
  await new Promise((r) => setTimeout(r, POLL_INTERVAL));
}
```

这种方式不适合大量账号。监测 50 个账号需要 50 个独立循环和 50 倍请求。此时应使用 X 列表。

***

## 第二级：通过一次请求监测多个账号

X 列表最多可包含 5,000 个账号。`/list-tweets` 一次 API 调用即可返回全部成员合并后的最新推文。这是生产环境多账号监测的默认方案，详情见[列表与社群](https://docs.sorsa.io/zh-Hans/lists-and-communities)。

### 第 1 步：创建公开 X 列表

1. 前往 [X Lists](https://x.com/i/lists) 创建列表。
2. 添加要监测的账号，最多 5,000 个。
3. 将列表设为 **Public**。API 无法访问私有列表。
4. 从 URL 复制**列表 ID**。例如 `https://x.com/i/lists/1234567890` 的 ID 为 `1234567890`。

### 第 2 步：轮询列表

```python theme={null}
import requests
import time

API_KEY = "YOUR_API_KEY"
LIST_ID = "YOUR_LIST_ID"
POLL_INTERVAL = 5

URL = f"https://api.sorsa.io/v3/list-tweets?list_id={LIST_ID}"
HEADERS = {"ApiKey": API_KEY, "Accept": "application/json"}


def monitor_list(callback, interval=POLL_INTERVAL):
    """Poll an X List and call `callback` for each new tweet detected."""
    last_seen_id = None
    print(f"Monitoring List {LIST_ID} (interval: {interval}s)")

    while True:
        try:
            resp = requests.get(URL, headers=HEADERS, timeout=10)
            resp.raise_for_status()
            tweets = resp.json().get("tweets", [])

            if not tweets:
                time.sleep(interval)
                continue

            top_id = max(int(t["id"]) for t in tweets)

            if last_seen_id is None:
                last_seen_id = top_id
                print(f"Baseline set: {last_seen_id}")
            else:
                new_tweets = [t for t in tweets if int(t["id"]) > last_seen_id]
                if new_tweets:
                    for tweet in reversed(new_tweets):
                        callback(tweet)
                    last_seen_id = top_id

        except requests.exceptions.RequestException as e:
            print(f"Request error: {e}. Retrying in {interval * 2}s")
            time.sleep(interval * 2)
            continue

        time.sleep(interval)


def on_new_tweet(tweet):
    user = tweet["user"]
    print(f"[NEW] @{user['username']}: {tweet['full_text'][:120]}")
    print(
        f"       Likes: {tweet.get('likes_count', 0)} | "
        f"RTs: {tweet.get('retweet_count', 0)} | "
        f"Views: {tweet.get('view_count', 'N/A')}\n"
    )


if __name__ == "__main__":
    monitor_list(on_new_tweet)
```

**效率提升。** 以 10 秒间隔分别轮询 50 个账号，每天消耗 50 × 8,640 = 432,000 次请求。将同样 50 个账号放入一个列表，以 10 秒间隔轮询，每天仅需 8,640 次，减少 50 倍。更多方式见[优化 API 使用](https://docs.sorsa.io/zh-Hans/optimizing-api-usage)。

> `/list-tweets` 每页最多返回 20 条。如果成员在一次轮询间隔内发布更多推文，可将间隔缩短至 2–3 秒，或使用 `next_cursor` 分页，直到遇到已见过的 ID。

***

## 第三级：监测关键词或话题标签

不必追踪具体账号，也可以使用 `order: "latest"` 轮询 `/search-tweets`，按时间获取匹配查询的结果。

```python theme={null}
import requests
import time

API_KEY = "YOUR_API_KEY"
QUERY = '("your brand" OR @yourbrand) lang:en'
POLL_INTERVAL = 10

URL = "https://api.sorsa.io/v3/search-tweets"
HEADERS = {"ApiKey": API_KEY, "Content-Type": "application/json"}


def monitor_keyword(query, callback, interval=10):
    last_seen_id = None
    print(f"Monitoring: {query} (interval: {interval}s)")

    while True:
        try:
            resp = requests.post(
                URL,
                headers=HEADERS,
                json={"query": query, "order": "latest"},
                timeout=10,
            )
            resp.raise_for_status()
            tweets = resp.json().get("tweets", [])

            if tweets:
                top_id = max(int(t["id"]) for t in tweets)
                if last_seen_id is None:
                    last_seen_id = top_id
                    print(f"Baseline set: {last_seen_id}")
                else:
                    new_tweets = [t for t in tweets if int(t["id"]) > last_seen_id]
                    for tweet in reversed(new_tweets):
                        callback(tweet)
                    if new_tweets:
                        last_seen_id = top_id

        except requests.exceptions.RequestException as e:
            print(f"Error: {e}")
            time.sleep(interval * 2)
            continue

        time.sleep(interval)


monitor_keyword(QUERY, on_new_tweet, interval=10)
```

查询字符串可使用任意[搜索运算符](https://docs.sorsa.io/zh-Hans/search-operators)。例如监测品牌的高互动英文提及，并排除转推：

```python theme={null}
monitor_keyword('"your brand" min_faves:10 lang:en -filter:retweets', on_new_tweet)
```

***

## 将新推文发送到 Slack、Discord 或任意 HTTP 端点

轮询循环负责生成数据，回调函数决定如何处理每条新推文。回调只是一个函数，因此同一个监测器可连接任何支持 HTTP 的目标。

### 通过 Incoming Webhook 发送到 Slack

```python theme={null}
import requests

SLACK_WEBHOOK_URL = "https://hooks.slack.com/services/YOUR/SLACK/WEBHOOK"


def send_to_slack(tweet):
    user = tweet["user"]
    text = (
        f"*New tweet from @{user['username']}*\n"
        f"{tweet['full_text']}\n"
        f"Likes: {tweet.get('likes_count', 0)} | "
        f"RTs: {tweet.get('retweet_count', 0)} | "
        f"Views: {tweet.get('view_count', 'N/A')}\n"
        f"https://x.com/{user['username']}/status/{tweet['id']}"
    )
    requests.post(SLACK_WEBHOOK_URL, json={"text": text})


# Plug into any monitor:
monitor_list(send_to_slack)
# or: monitor_keyword("bitcoin lang:en min_faves:50", send_to_slack)
```

### Discord

```python theme={null}
DISCORD_WEBHOOK_URL = "https://discord.com/api/webhooks/YOUR/WEBHOOK"


def send_to_discord(tweet):
    user = tweet["user"]
    content = (
        f"**@{user['username']}** just tweeted:\n"
        f"{tweet['full_text']}\n"
        f"https://x.com/{user['username']}/status/{tweet['id']}"
    )
    requests.post(DISCORD_WEBHOOK_URL, json={"content": content})
```

### Telegram

```python theme={null}
TELEGRAM_BOT_TOKEN = "YOUR_BOT_TOKEN"
TELEGRAM_CHAT_ID = "YOUR_CHAT_ID"


def send_to_telegram(tweet):
    user = tweet["user"]
    text = (
        f"New tweet from @{user['username']}\n\n"
        f"{tweet['full_text']}\n\n"
        f"https://x.com/{user['username']}/status/{tweet['id']}"
    )
    requests.post(
        f"https://api.telegram.org/bot{TELEGRAM_BOT_TOKEN}/sendMessage",
        json={"chat_id": TELEGRAM_CHAT_ID, "text": text},
    )
```

### 自定义 HTTP 端点

```python theme={null}
def send_to_internal_api(tweet):
    requests.post(
        "https://internal.example.com/events/twitter",
        json={
            "tweet_id": tweet["id"],
            "username": tweet["user"]["username"],
            "text": tweet["full_text"],
            "metrics": {
                "likes": tweet.get("likes_count", 0),
                "retweets": tweet.get("retweet_count", 0),
                "views": tweet.get("view_count", 0),
            },
            "url": f"https://x.com/{tweet['user']['username']}/status/{tweet['id']}",
        },
        headers={"Authorization": "Bearer YOUR_INTERNAL_TOKEN"},
        timeout=5,
    )
```

***

## API 用量估算

下表假设每轮只有一页、固定调度且无需重试。请乘以监测器数量，并加上额外页面和重试次数。示例循环在收到响应后休眠，因此实际周期还包含网络和处理时间。

| 间隔   | 每小时请求 | 每天请求   | 每月请求（30 天） |
| :--- | :---- | :----- | :--------- |
| 1 秒  | 3,600 | 86,400 | 2,592,000  |
| 5 秒  | 720   | 17,280 | 518,400    |
| 10 秒 | 360   | 8,640  | 259,200    |
| 30 秒 | 120   | 2,880  | 86,400     |
| 1 分钟 | 60    | 1,440  | 43,200     |

根据应用可接受的延迟和信息流活跃程度选择间隔。缩短间隔会增加请求数，但不能保证帖子立即出现在搜索结果中。

免费 100 次请求足够构建并端到端验证原型。持续监测时，按表中月度用量选套餐：单个循环以 30–60 秒间隔运行适合 Pro（每月 100,000 次），10 秒间隔适合 Enterprise（每月 500,000 次）。这些数字按单个监测器计算，多个并行运行会相应增加总量，请按合计用量选择。完整套餐见[价格](https://api.sorsa.io/pricing)。

> 超过标准套餐的速率或用量需求，请[联系销售](https://api.sorsa.io/talk-to-sales)定制配额，或在 [Discord](https://discord.com/invite/uwAefKCj7X) 咨询。

***

## 生产环境加固

上方示例适用于开发。生产环境需要处理以下五点。

### 1. 在重启之间持久化 `last_seen_id`

脚本崩溃后如果不知道上次检查点，可能重复处理旧推文、发送重复提醒，或静默跳过缺口。将最后已见 ID 保存在文件、数据库或 Redis 中。

```python theme={null}
import json
import os

STATE_FILE = "monitor_state.json"


def load_state():
    if os.path.exists(STATE_FILE):
        with open(STATE_FILE) as f:
            return json.load(f).get("last_seen_id")
    return None


def save_state(last_seen_id):
    with open(STATE_FILE, "w") as f:
        json.dump({"last_seen_id": last_seen_id}, f)
```

将监测器初始化的 `last_seen_id = None` 替换为 `last_seen_id = load_state()`。在一轮所有页面处理完成或进入持久化队列后，保存新检查点。页面获取或投递失败时不要推进检查点。重启后应先分页补齐缺口，再接受更新的检查点。

### 2. 对错误使用指数退避

网络问题、速率限制（HTTP 429）和临时 API 错误不可避免。应逐步增加等待时间并设置上限，而不是立即不断重试。完整说明见[错误码](https://docs.sorsa.io/zh-Hans/error-codes)。

```python theme={null}
retry_delay = POLL_INTERVAL
MAX_DELAY = 60

while True:
    try:
        resp = requests.get(URL, headers=HEADERS, timeout=10)
        if resp.status_code == 429:
            print(f"Rate limited. Backing off {retry_delay}s")
            time.sleep(retry_delay)
            retry_delay = min(retry_delay * 2, MAX_DELAY)
            continue
        resp.raise_for_status()
        retry_delay = POLL_INTERVAL  # reset on success
        # process tweets
    except requests.exceptions.RequestException as e:
        print(f"Error: {e}")
        time.sleep(retry_delay)
        retry_delay = min(retry_delay * 2, MAX_DELAY)
        continue

    time.sleep(POLL_INTERVAL)
```

### 3. 将轮询与处理分离

不要在轮询循环内同步执行 NLP、数据库写入、外部 API 调用等耗时操作。下游变慢会导致轮询落后。将新推文放入队列，由独立工作进程处理。

```python theme={null}
from collections import deque
import threading

tweet_queue = deque()


def polling_loop():
    """Fast loop: poll and enqueue. No heavy work here."""
    # Standard polling code, but instead of calling callback(tweet):
    # tweet_queue.append(tweet)
    pass


def processing_worker():
    """Separate thread: dequeue and dispatch."""
    while True:
        if tweet_queue:
            tweet = tweet_queue.popleft()
            send_to_slack(tweet)
            save_to_database(tweet)
        else:
            time.sleep(0.1)


threading.Thread(target=processing_worker, daemon=True).start()
polling_loop()
```

更大工作负载可将内存 deque 替换为 Redis、RabbitMQ、SQS 或现有技术栈使用的消息系统。

### 4. 监测监测器本身

记录每轮的时间戳、新推文数、响应时间和错误。如果过去 N 分钟没有成功完成轮询，应发出提醒。静默故障会造成难以察觉的数据缺口。API 运行情况可查看 [Sorsa 状态页](https://uptime.sorsa.io/status/v3)。

### 5. 处理边界情况

**页面溢出和延迟到达：** 持续使用 `next_cursor`，直到覆盖上次检查点以来的时间段。保留少量重叠并按存储的 ID 去重，避免仅因延迟结果的 ID 较旧而丢弃。前面的首页示例没有实现这种回填。

**回调投递：** 检查 webhook 响应状态，使用有限重试或持久化队列。Sorsa 请求成功不代表 Slack、Discord 或数据库已接收事件。

* **已删除推文：** 如果推文在获取后、回调前被删除，URL 会返回 404，应视为正常情况。
* **受保护账号：** 追踪用户设为私有后，`/user-tweets` 返回空列表。记录并继续。
* **置顶推文：** `/user-tweets` 的第一条往往是置顶内容，而非最新内容。不要直接用 `tweets[0]` 作为最新 ID；应使用 `max(int(t["id"]) for t in tweets)`（如上方示例），或按 `created_at` 排序。
* **转推：** 转推的 `tweet["retweeted_status"]` 有值，可据此决定是否纳入。
* **回复限制：** `is_replies_limited` 表示作者限制了回复，对某些监测场景有价值。

***

## 后续步骤

* [搜索运算符](https://docs.sorsa.io/zh-Hans/search-operators)：通过高级筛选减少关键词监测噪声
* [追踪提及](https://docs.sorsa.io/zh-Hans/search-mentions)：支持互动筛选的专用 @提及端点
* [速率限制](https://docs.sorsa.io/zh-Hans/rate-limits)：处理 429 错误和请求节奏
* [分页](https://docs.sorsa.io/zh-Hans/pagination)：结合实时监测回填历史数据
* [API 参考](https://docs.sorsa.io/zh-Hans/api-reference-guide)：`/list-tweets`、`/user-tweets`、`/search-tweets` 及全部端点的完整规范
