2016-11-09 2 views
1

これに続いて、great tutorialはtweepyを使ってPythonでライブのTwitterストリームを活用しています。これにより、RxJava、RxPy、RxScala、またはReactiveXというライブタイムでツイートが印刷されます。RxPy - ライブTwitterストリームを受信可能なRxに変換しますか?

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']) 

これはRxPyを経由して観察可能ReactiveXに変身するのに最適な候補です。しかし、どのようにして正確にこれをホットソースObservableに変えるのですか? Observable.create()を実行する方法のどこにでもドキュメンテーションを見つけることができないようです...

+0

私はSubjectでこれを達成できると思いましたが、私は成功しました。しかし、私はこの件名を無料でできるかどうかまだ疑問に思っています... – tmn

答えて

0

私はこれをしばらく前に考え出しました。渡された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)) 
関連する問題