Prefect
用 Prefect 与 TwexAPI 只读任务在 Python 中构建定时 Twitter 搜索、资料、时间线与趋势工作流。
Prefect 是 Python 代码的开源工作流编排器。用 prefect-x-api-scraper 实现可重复的推文搜索、资料查询、时间线刷新与趋势检查。
该集合为只读,提供六个异步 Prefect 任务。每个任务将规范 TwexAPI JSON 响应作为 Python 字典返回。
搜索推文
搜索关键词、话题标签、账号、日期与 X 查询运算符。
获取推文详情
从数字 ID 获取一条公开推文。
搜索用户画像
按名称、用户名或主题查找公开 X 账号。
查询用户资料
按 screen name 获取一条公开资料。
刷新用户时间线
获取用户时间线页的最新推文。
追踪热门趋势
按国家、主题与内容标签获取全球热门推文。
用于研究、enrichment、仪表盘、告警与索引。写操作、粉丝分页与导出请用直接 REST、Python SDK 或 MCP。
安装
使用 Python 3.10 或更新版本。
python -m pip install "prefect>=3.0.0" "prefect-x-api-scraper"
安装 Prefect 后注册凭证 block:
prefect block register -m prefect_x_api_scraper
存储 TwexAPI API 密钥
在 TwexAPI dashboard 创建 API 密钥,存入 TwexApiCredentials block。
from prefect_x_api_scraper import TwexApiCredentials
credentials = TwexApiCredentials(
api_key="YOUR_API_KEY",
base_url="https://api.twexapi.io",
timeout_seconds=30,
)
credentials.save("twexapi", overwrite=True)
Prefect block 在 flow 与 deployment 间存储类型化配置。切勿将 API 密钥放入 deployment YAML、flow 参数、日志或仓库。
选择正确的 Prefect 任务
search_tweets
调用 POST /twitter/advanced_search/page。接受 query、cursor、sort_by 与分页字段。
get_tweet
调用 POST /v2/tweet/detail。传入一个数字 tweet ID。
search_users
调用 GET /twitter/search-user/{keyword}/{target_count}。接受资料查询。
get_user
调用 GET /twitter/{screen_name}/about。接受带或不带 @ 的 screen name。
get_user_tweets
调用 POST /twitter/{screen_name}/timeline/page。支持游标与页大小。
get_trends
调用 GET /twitter/global-trending/tweets。接受 country、topic、content、count。
在 Python 中构建 Twitter 自动化 flow
此 flow 搜索近期帖子并规范化推文行,保留 ID、作者、时间戳、指标、URL 与分页状态。
from __future__ import annotations
from typing import Any
from prefect import flow, task
from prefect_x_api_scraper import TwexApiCredentials, search_tweets
@task
async def normalize_tweet_page(page: dict[str, Any]) -> list[dict[str, Any]]:
rows: list[dict[str, Any]] = []
for tweet in page.get("data", {}).get("tweets", page.get("tweets", [])):
if not isinstance(tweet, dict):
continue
author = tweet.get("author") or tweet.get("user")
rows.append(
{
"tweet_id": tweet.get("id") or tweet.get("tweet_id"),
"text": tweet.get("text") or tweet.get("full_text"),
"created_at": tweet.get("created_at") or tweet.get("createdAt"),
"author": author if isinstance(author, dict) else {},
"like_count": tweet.get("like_count") or tweet.get("likeCount"),
"repost_count": tweet.get("retweet_count") or tweet.get("retweetCount"),
"reply_count": tweet.get("reply_count") or tweet.get("replyCount"),
}
)
return rows
@flow(name="TwexAPI Twitter Search")
async def social_signal_flow() -> dict[str, Any]:
credentials = TwexApiCredentials.load("twexapi")
page = await search_tweets(
credentials,
'"workflow orchestration" lang:en -filter:retweets',
sort_by="Latest",
)
rows = await normalize_tweet_page(page)
return {
"tweet_rows": rows,
"has_more": page.get("has_more") or page.get("has_next_page", False),
"next_cursor": page.get("next_cursor") or page.get("nextCursor"),
}
用 asyncio 运行:
import asyncio
result = asyncio.run(social_signal_flow())
编写聚焦推文搜索查询
精确短语匹配
精确短语用 "workflow orchestration"。
按账号筛选
单账号推文用 from:PrefectIO。
话题标签搜索
匹配话题标签用 #prefect #python。
时间范围窗口
查询内用 since: 与 until: 日期。
排除转推
只要原创时用 -filter:retweets。
时间顺序监控用 sort_by="Latest";按互动发现用 sort_by="Top"。排名可能变化,请持久化 tweet ID。
调度 Twitter 搜索管道
本地长期进程用 .serve()。
from prefect.schedules import Cron
if __name__ == "__main__":
social_signal_flow.serve(
name="twexapi-social-signals",
schedule=Cron("0 * * * *", timezone="UTC"),
)
Docker、Kubernetes 或无 serverless worker 用 work-pool deployment。
prefect deploy social_signal_flow.py:social_signal_flow \
--name twexapi-social-signals \
--pool production
分页推文与资料结果
游标分页在不猜测页码的情况下延续大搜索。保持原始请求不变,仅传入返回的游标。
from typing import Any, Optional
from prefect import flow
from prefect_x_api_scraper import TwexApiCredentials, search_tweets
@flow
async def collect_tweet_pages(query: str) -> list[dict[str, Any]]:
credentials = TwexApiCredentials.load("twexapi")
rows_by_id: dict[str, dict[str, Any]] = {}
cursor: Optional[str] = None
while True:
page = await search_tweets(
credentials,
query=query,
sort_by="Latest",
cursor=cursor,
)
tweets = page.get("data", {}).get("tweets", page.get("tweets", []))
for tweet in tweets:
if isinstance(tweet, dict):
tweet_id = tweet.get("id") or tweet.get("tweet_id")
if tweet_id:
rows_by_id[str(tweet_id)] = tweet
cursor_value = page.get("next_cursor") or page.get("nextCursor")
has_more = bool(page.get("has_more") or page.get("has_next_page", False))
if not has_more or not cursor_value:
break
cursor = str(cursor_value)
return list(rows_by_id.values())
勿解码游标,视为 opaque 字符串。持久化每页后存储对应游标。
使定时运行幂等
推文标识信息
按 tweet_id upsert 推文行。勿用文本作键。
账号身份信息
按数字用户 ID upsert 资料行。
游标检查点
每提交一页后存储 has_more 与 next_cursor。
仅重试 transient 失败
Prefect 支持重试延迟、抖动与条件重试。无效输入切勿原样重试。
from prefect_x_api_scraper import search_tweets
search_recent_tweets = search_tweets.with_options(
name="Search Recent Tweets",
retries=4,
retry_delay_seconds=[5, 15, 45, 120],
)
路由文档化错误
| 状态 | 操作 |
|---|---|
400 |
修正缺失查询或畸形输入。切勿原样重试。 |
401 |
从 credentials block 加载有效 API 密钥。 |
403 |
下次运行前解决账号访问或额度。 |
404 |
检查 tweet ID、用户名与 screen name。 |
429 |
放慢调度并以抖动延迟重试。 |
完整恢复指南见 Error Handling 与 Rate Limits。
| 5xx | 有上限地重试 transient 失败。 |
MCP 规划替代方案
规划阶段用 Twexapi MCP 发现端点,再将选定路由固化到 Prefect 任务或 REST 调用。
建议规划 prompt:
Use Twexapi MCP explore to choose the best endpoint for daily AI trend collection.
Return the exact method, relative path, query parameters, pagination fields, expected response fields, and retry guidance.
Do not execute write endpoints.
Prefect 集合或直接 TwexAPI API
六个支持的读取选 prefect-x-api-scraper。它提供 block、异步调用、校验与任务元数据。
以下请用直接 REST、Python SDK 或 MCP:
- 推文、回复、引用、点赞、关注或 DM 操作
- 粉丝与关注导出
- 列表与社区
recurring 检查请调度 Prefect flow 或 REST 分页,而非原生监控端点。TwexAPI 不将异步 CSV/JSON 提取作为独立产品暴露。
写操作放在明确审批之后。重试任何写操作前先存储已确认 ID。