RxPy - 将实时 Twitter 流转换为 Rx Observable?
RxPy - Turn Live Twitter Stream into Rx Observable?
我按照这个 great tutorial 使用 tweepy 在 Python 中利用实时 Twitter 流。这将实时打印提及 RxJava、RxPy、RxScala 或 ReactiveX 的推文。
from tweepy.streaming import StreamListener
from tweepy import OAuthHandler
from tweepy import Stream
from rx import Observable, Observer
#Variables that contains the user credentials to access Twitter API
access_token = "CONFIDENTIAL"
access_token_secret = "CONFIDENTIAL"
consumer_key = "CONFIDENTIAL"
consumer_secret = "CONFIDENTIAL"
#This is a basic listener that just prints received tweets to stdout.
class TweetObserver(StreamListener):
def on_data(self, data):
print(data)
return True
def on_error(self, status):
print(status)
if __name__ == '__main__':
#This handles Twitter authetification and the connection to Twitter Streaming API
l = TweetObserver()
auth = OAuthHandler(consumer_key, consumer_secret)
auth.set_access_token(access_token, access_token_secret)
stream = Stream(auth, l)
#This line filter Twitter Streams to capture data by the keywords: 'python', 'javascript', 'ruby'
stream.filter(track=['rxjava','rxpy','reactivex','rxscala'])
这是成为 ReactiveX Observable via RxPy 的完美人选。但我究竟如何将它变成一个热点Observable
?我似乎无法在任何地方找到有关如何执行 Observable.create()
...
的文档
我很久以前就想通了。您必须定义一个函数来处理传递的 Observer
参数。然后你将它传递给 Observable.create()
。
from tweepy.streaming import StreamListener
from tweepy import OAuthHandler
from tweepy import Stream
import json
from rx import Observable
# Variables that contains the user credentials to access Twitter API
access_token = "PUT YOURS HERE"
access_token_secret = "PUT YOURS HERE"
consumer_key = "PUT YOURS HERE"
consumer_secret = "PUT YOURS HERE"
def tweets_for(topics):
def observe_tweets(observer):
class TweetListener(StreamListener):
def on_data(self, data):
observer.on_next(data)
return True
def on_error(self, status):
observer.on_error(status)
# This handles Twitter authetification and the connection to Twitter Streaming API
l = TweetListener()
auth = OAuthHandler(consumer_key, consumer_secret)
auth.set_access_token(access_token, access_token_secret)
stream = Stream(auth, l)
stream.filter(track=topics)
return Observable.create(observe_tweets).share()
topics = ['Britain', 'France']
tweets_for(topics) \
.map(lambda d: json.loads(d)) \
.subscribe(on_next=lambda s: print(s), on_error=lambda e: print(e))
我按照这个 great tutorial 使用 tweepy 在 Python 中利用实时 Twitter 流。这将实时打印提及 RxJava、RxPy、RxScala 或 ReactiveX 的推文。
from tweepy.streaming import StreamListener
from tweepy import OAuthHandler
from tweepy import Stream
from rx import Observable, Observer
#Variables that contains the user credentials to access Twitter API
access_token = "CONFIDENTIAL"
access_token_secret = "CONFIDENTIAL"
consumer_key = "CONFIDENTIAL"
consumer_secret = "CONFIDENTIAL"
#This is a basic listener that just prints received tweets to stdout.
class TweetObserver(StreamListener):
def on_data(self, data):
print(data)
return True
def on_error(self, status):
print(status)
if __name__ == '__main__':
#This handles Twitter authetification and the connection to Twitter Streaming API
l = TweetObserver()
auth = OAuthHandler(consumer_key, consumer_secret)
auth.set_access_token(access_token, access_token_secret)
stream = Stream(auth, l)
#This line filter Twitter Streams to capture data by the keywords: 'python', 'javascript', 'ruby'
stream.filter(track=['rxjava','rxpy','reactivex','rxscala'])
这是成为 ReactiveX Observable via RxPy 的完美人选。但我究竟如何将它变成一个热点Observable
?我似乎无法在任何地方找到有关如何执行 Observable.create()
...
我很久以前就想通了。您必须定义一个函数来处理传递的 Observer
参数。然后你将它传递给 Observable.create()
。
from tweepy.streaming import StreamListener
from tweepy import OAuthHandler
from tweepy import Stream
import json
from rx import Observable
# Variables that contains the user credentials to access Twitter API
access_token = "PUT YOURS HERE"
access_token_secret = "PUT YOURS HERE"
consumer_key = "PUT YOURS HERE"
consumer_secret = "PUT YOURS HERE"
def tweets_for(topics):
def observe_tweets(observer):
class TweetListener(StreamListener):
def on_data(self, data):
observer.on_next(data)
return True
def on_error(self, status):
observer.on_error(status)
# This handles Twitter authetification and the connection to Twitter Streaming API
l = TweetListener()
auth = OAuthHandler(consumer_key, consumer_secret)
auth.set_access_token(access_token, access_token_secret)
stream = Stream(auth, l)
stream.filter(track=topics)
return Observable.create(observe_tweets).share()
topics = ['Britain', 'France']
tweets_for(topics) \
.map(lambda d: json.loads(d)) \
.subscribe(on_next=lambda s: print(s), on_error=lambda e: print(e))