You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何同时运行Twitter StreamingClient两条流规则?多进程可行吗?

同时运行两条Twitter流规则的实现方案

当然可以同时运行两条流规则,有两种常用实现方式,根据你的需求选择即可:

方案一:单流多规则(最简便)

Twitter API的流本身支持在同一个连接中添加多条规则,无需额外的多进程处理。你只需要取消注释原代码中的rule2,并一次性添加所有规则即可:

import sys
import time
from tweepy import StreamingClient, StreamRule

user = "目标用户名"  # 替换为实际目标用户
bearer_token = "你的Bearer Token"  # 替换为实际凭证

value1 = ('from:{user} -is:retweet').format(user=user)
value2 = ('to:{user} is:reply').format(user=user)

class CustomStream(StreamingClient):
    def on_connect(self):
        print("Connected!")
    def on_tweet(self, tweet):
        try:
            print(f"{tweet.id} {tweet.created_at} ({tweet.author_id}): {tweet.text}")
            print("-" * 50)
            time.sleep(0.2)
            return True
        except BaseException as e:
            sys.stderr.write('Error on_data: {}\n'.format(e))
            time.sleep(5)
        return True

stream = CustomStream(bearer_token=bearer_token, wait_on_rate_limit=True)
previousRules = stream.get_rules().data
if previousRules:
    stream.delete_rules(previousRules)

# 同时添加两条规则
rules = [StreamRule(value=value1), StreamRule(value=value2)]
stream.add_rules(rules)
stream.filter(tweet_fields=["created_at", "author_id"], expansions=["author_id"])

此方案下,同一个流会推送匹配任意一条规则的推文,所有推文共享同一个on_tweet处理逻辑,资源占用更低。

方案二:多进程独立流(适用于分场景处理)

如果两条规则的推文需要不同的处理逻辑,或者你需要完全独立的两个流连接,可以用Python的multiprocessing模块实现:

import multiprocessing
import sys
import time
from tweepy import StreamingClient, StreamRule

def run_stream(bearer_token, rule_value, stream_label):
    class CustomStream(StreamingClient):
        def on_connect(self):
            print(f"{stream_label} 已连接!")
        def on_tweet(self, tweet):
            try:
                print(f"[{stream_label}] {tweet.id} {tweet.created_at} ({tweet.author_id}): {tweet.text}")
                print("-" * 50)
                time.sleep(0.2)
                return True
            except BaseException as e:
                sys.stderr.write(f'[{stream_label}] 处理错误: {e}\n')
                time.sleep(5)
            return True

    stream = CustomStream(bearer_token=bearer_token, wait_on_rate_limit=True)
    previousRules = stream.get_rules().data
    if previousRules:
        stream.delete_rules(previousRules)
    stream.add_rules(StreamRule(value=rule_value))
    stream.filter(tweet_fields=["created_at", "author_id"], expansions=["author_id"])

if __name__ == "__main__":
    user = "目标用户名"
    bearer_token = "你的Bearer Token"
    value1 = ('from:{user} -is:retweet').format(user=user)
    value2 = ('to:{user} is:reply').format(user=user)

    # 创建两个独立进程分别运行不同流
    stream1_process = multiprocessing.Process(target=run_stream, args=(bearer_token, value1, "来自目标用户的推文"))
    stream2_process = multiprocessing.Process(target=run_stream, args=(bearer_token, value2, "@目标用户的回复"))

    stream1_process.start()
    stream2_process.start()

    stream1_process.join()
    stream2_process.join()

此方案中,两个流完全独立运行,各自拥有独立的连接和处理逻辑,通过stream_label可以清晰区分不同流的输出。

注意事项

  • 使用多进程时,需确保你的Bearer Token支持同时创建多个流连接(具体限制可参考Twitter API官方说明)。
  • 若无需分开处理推文逻辑,优先选择方案一,更节省系统资源。

内容的提问来源于stack exchange,提问作者Алеся

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 20:20:51