Building a Resilient Twitter Filtered Stream Client in Python
A practical guide to Twitter API v2 Filtered Stream: rule management, resilient HTTP/2 connection handling, and a production-ready Python client with exponential back-off.
28 Jun 2026, 00:41 UTC

The Problem: Real-Time Tweet Ingestion Without the Firehose
You need to capture tweets matching specific criteria—say, English-language posts about AI—in real time. The full firehose is overkill and access-restricted. Twitter's API v2 Filtered Stream solves this by letting you define rules (keywords, hashtags, language filters) and opening a persistent HTTP/2 connection that pushes only matching tweets. The challenge isn't starting the stream; it's keeping it alive through rate limits, backpressure, and inevitable network hiccups.
How Filtered Stream Works
The endpoint GET /2/tweets/search/stream holds a long-lived connection. You manage rules separately via POST/GET/DELETE /2/tweets/search/stream/rules. Each project supports up to 512 rules; a rule like #AI lang:en matches English tweets containing that hashtag. When a tweet matches any active rule, Twitter streams a JSON line with the tweet object (default fields: id, text; expandable via tweet.fields and expansions query parameters).
Critical behavior: the stream does not guarantee delivery during extreme spikes. If your client falls behind, tweets may be dropped. Twitter signals backpressure via a disconnect message with code StreamingConnectionException—your cue to reconnect with exponential back-off.
Rule Management: Batch Operations and Rate Limits
Rules are added, listed, or deleted in batches. A single POST to /rules with {"add":[{"value":"#AI lang:en"},{"value":"machine learning lang:en"}]} creates two rules. The response includes rule IDs you'll need for deletion. Rate limits on rule endpoints are strict (undocumented per-window limits); exceeding them returns HTTP 429. Treat rule changes as infrequent administrative actions, not runtime toggles.
# Example: add a rule via curl
curl -X POST "https://api.twitter.com/2/tweets/search/stream/rules" \
-H "Authorization: Bearer $BEARER_TOKEN" \
-H "Content-Type: application/json" \
-d '{"add":[{"value":"#AI lang:en"}]}'
Run this from any machine with the bearer token. Requires tweet.read scope. Expect a 201 with the created rule's id and value. If you see 429, wait at least 15 minutes before retrying.
Worked Example: A Resilient Python Client
The script below uses httpx (async, HTTP/2 support) to connect, parse JSON lines, and handle reconnection with exponential back-off (capped at 60 s). It logs each reconnect attempt so you can observe stability in production.
import os, json, asyncio, logging, httpx
BEARER = os.environ["TWITTER_BEARER_TOKEN"]
RULE_VALUE = "#AI lang:en"
STREAM_URL = "https://api.twitter.com/2/tweets/search/stream"
RULES_URL = "https://api.twitter.com/2/tweets/search/stream/rules"
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
log = logging.getLogger(__name__)
async def ensure_rule(client: httpx.AsyncClient) -> str:
"""Create the rule if missing; return its ID."""
resp = await client.get(RULES_URL, headers={"Authorization": f"Bearer {BEARER}"})
resp.raise_for_status()
existing = resp.json().get("data", [])
for r in existing:
if r["value"] == RULE_VALUE:
return r["id"]
resp = await client.post(RULES_URL, headers={"Authorization": f"Bearer {BEARER}"},
json={"add": [{"value": RULE_VALUE}]})
if resp.status_code == 429:
raise RuntimeError("Rule endpoint rate limited (429)")
resp.raise_for_status()
return resp.json()["data"][0]["id"]
async def stream_tweets(client: httpx.AsyncClient):
backoff = 1
while True:
try:
async with client.stream("GET", STREAM_URL,
headers={"Authorization": f"Bearer {BEARER}"},
timeout=None) as resp:
if resp.status_code == 429:
raise httpx.HTTPStatusError("Rate limited", request=resp.request, response=resp)
if resp.status_code >= 500:
raise httpx.HTTPStatusError(f"Server error {resp.status_code}", request=resp.request, response=resp)
resp.raise_for_status()
backoff = 1 # reset on successful connection
async for line in resp.aiter_lines():
if not line:
continue
if line.startswith("{\"disconnect\"):
log.warning("Received disconnect signal: %s", line)
break
tweet = json.loads(line)
# Process tweet—here we just log id and text
log.info("Tweet %s: %s", tweet["data"]["id"], tweet["data"]["text"][:80])
except (httpx.HTTPStatusError, httpx.TransportError, httpx.StreamError) as exc:
log.warning("Stream error: %s. Reconnecting in %ds", exc, backoff)
await asyncio.sleep(backoff)
backoff = min(backoff * 2, 60)
async def main():
async with httpx.AsyncClient(http2=True) as client:
rule_id = await ensure_rule(client)
log.info("Using rule %s", rule_id)
await stream_tweets(client)
if __name__ == "__main__":
asyncio.run(main())
Run with TWITTER_BEARER_TOKEN set. The script creates the rule on first run, then streams indefinitely. Kill with Ctrl-C. No state is mutated beyond rule creation, so rollback isn't needed—just delete the rule via the API if you want to clean up.
Trade-offs and Limitations
- No exactly-once guarantee. During reconnects you may miss tweets or see duplicates if you don't track
ids. - Rule limits. 512 rules per project; complex boolean logic requires multiple rules.
- Backpressure drops. If your processing lags, Twitter silently drops tweets. Monitor the
disconnectmessage and consider a buffer (e.g., write to a local queue or file) before downstream processing. - Rate limits on rule changes. Automate rule updates sparingly; cache rule IDs locally.
Verify It Works
- Create a Twitter Developer Project & App, generate a Bearer Token with
tweet.readscope. - Run the curl example above to confirm rule creation returns 201.
- Execute the Python script. You should see log lines like
Tweet 1234567890: Exciting developments in #AI.... - Temporarily disconnect your network; the script should log reconnection attempts with increasing back-off.
If you receive no tweets after several minutes, double-check the rule value (case-insensitive, but syntax matters) and that your token has the correct scope.
Next Steps
Add tweet.fields=created_at,author_id,public_metrics and expansions=author_id to the stream URL for richer data. Persist tweets to a time-series database or object store. For production, run multiple client instances with different rule subsets behind a load balancer to scale ingestion horizontally.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.