【问题标题】:Using multiple threads to get data from twitter using twitter4j使用 twitter4j 使用多个线程从 twitter 获取数据
【发布时间】:2014-02-18 10:31:07
【问题描述】:

我有一组关键字(超过 600 个),我想使用流 api 来跟踪它们的推文。 Twitter api 将允许您跟踪的关键字数量限制为 200 个。因此,我决定使用多个线程来执行此操作,为此使用多个 OAuth 令牌。我就是这样做的:

String[] dbKeywords = KeywordImpl.listKeywords();
    List<String[]> keywords = ditributeKeywords(dbKeywords);
    for (String[] subList : keywords) {
        StreamCrawler streamCrawler = new StreamCrawler();
        streamCrawler.setKeywords(subList);
        Thread crawlerThread = new Thread(streamCrawler);
        crawlerThread.start();
    }

这就是单词在线程中的分布方式。每个线程接收不超过 200 个单词。 这是 StreamCrawler 的实现:

public class StreamCrawler extends Crawler implements Runnable {

...

    private String[] keywords;
    public void setKeywords(String[] keywords) {
    this.keywords = keywords;
}

@Override
public void run() {
    TwitterStream twitterStream = getTwitterInstance();
    StatusListener listener = new StatusListener() {
        ArrayDeque<Tweet> tweetbuffer = new ArrayDeque<Tweet>();
        ArrayDeque<TwitterUser> userbuffer = new ArrayDeque<TwitterUser>();


        @Override
        public void onException(Exception arg0) {
            System.out.println(arg0);
        }

        @Override
        public void onDeletionNotice(StatusDeletionNotice arg0) {
            System.out.println(arg0);
        }

        @Override
        public void onScrubGeo(long arg0, long arg1) {
            System.out.println(arg1);
        }

        @Override
        public void onStatus(Status status) {
                 ...Doing something with message
        }

        @Override
        public void onTrackLimitationNotice(int arg0) {
            System.out.println(arg0);
            try {
                Thread.sleep(5 * 60 * 1000);
                System.out.println("Will sleep for 5 minutes!");
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }

        @Override
        public void onStallWarning(StallWarning arg0) {
            System.out.println(arg0);
        }

    };

    FilterQuery fq = new FilterQuery();
    String keywords[] = getKeywords();
    System.out.println(keywords.length);
    System.out.println("Listening for " + Arrays.toString(keywords));
    fq.track(keywords);
    twitterStream.addListener(listener);
    twitterStream.filter(fq);
}

private long getCurrentThreadId() {
    return Thread.currentThread().getId();
}

private TwitterStream getTwitterInstance() {
    TwitterConfiguration configuration = null;
    TwitterStream twitterStream = null;
    while (configuration == null) {
        configuration = TokenFactory.getAvailableToken();
        if (configuration != null) {
            System.out
                    .println("Token was obtained " + getCurrentThreadId());
            System.out.println(configuration.getTwitterAccount());
            setToken(configuration);
            ConfigurationBuilder cb = new ConfigurationBuilder();
            cb.setDebugEnabled(true);
            cb.setOAuthConsumerKey(configuration.getConsumerKey());
            cb.setOAuthConsumerSecret(configuration.getConsumerSecret());
            cb.setOAuthAccessToken(configuration.getAccessToken());
            cb.setOAuthAccessTokenSecret(configuration.getAccessSecret());
            twitterStream = new TwitterStreamFactory(cb.build())
                    .getInstance();
        } else {
            // If there is no available configuration, wait for 2 minutes
            // and try again
            try {
                System.out
                        .println("There were no available tokens, sleeping for 2 minutes.");
                Thread.sleep(2 * 60 * 1000);
            } catch (InterruptedException e) {
                // TODO Auto-generated catch block
                e.printStackTrace();
            }
        }
    }
    return twitterStream;
    }
}

所以我的问题是,当我启动例如 2 个线程时,我收到通知,它们都在打开流并获取它。但实际上只有第一个真正得到流并分别调用 OnStatus 方法。在第二个线程中使用的数组不为空; Twitter 配置也是有效且唯一的。所以我不明白这种行为可能是什么原因。为什么唯一的第一个线程返回推文?

【问题讨论】:

  • 这很奇怪,因为一切似乎都是独一无二的。为什么不运行多个应用程序而不是在同一个应用程序中运行带有线程的流?
  • 我没有想到这种方法,因为关键字的数量会增加,所以我必须运行5个以上的实例。此外,我会说这是不明智的,因为我看不出用不同的线程来做这件事有什么障碍。
  • 我同意它应该以某种方式工作。如果您的密钥集很小,我只是提供了一个想法。 :D

标签: java multithreading twitter twitter4j


【解决方案1】:

据我所知,您正在尝试从同一个 IP 建立两个与公共流媒体端点(又名通用流或 stream.twitter.com)的同时连接。
更具体地说,我认为您希望从同一个 IP 到 stream.twitter.com/1.1/statuses/filter.json 的两个活动连接。

虽然 Twitter 流式 API 文档没有明确说明只有一个与公共端点的常设连接,但 Twitter 员工在开发网站 https://dev.twitter.com/discussions/7542 上澄清了这一点

对于一般流,您应该只从同一个 IP 建立一个连接。

这意味着您使用两个不同的 Twitter 应用程序/帐户连接到公共信息流并不重要;只要您从同一个 IP 地址进行连接,您就只能与公共流建立一个常设连接。你说你把两个流都连接了,这个行为的答案是由 Twitter 员工给出的:https://dev.twitter.com/discussions/14935

您可能会发现有时 stream.twitter.com 可以让您在这里或那里获得更多开放连接,但不应指望这种行为。

如果您尝试例如在第二个线程中连接到用户流(twitter4j TwitterStream user() 方法),那么您将真正开始同时获得过滤器和用户流。

关于 200 个跟踪关键字的限制,可能 twitter4j.org javadoc 有点过时了。这是 twitter api 文档所说的

默认访问级别最多允许 400 个跟踪关键字、5,000 个关注用户 ID 和 25 个 0.1-360 度位置框。如果您需要提升对 Streaming API 的访问权限,您应该探索我们的 Twitter 数据合作伙伴提供商...

因此,如果您需要超过 400 个,您可能需要要求 Twitter 提高您的 Twitter 帐户应用程序的跟踪访问级别,或者与 Twitter 数据的认证合作伙伴提供商合作。

您不一定需要的另一件事是启动新线程来获取流,因为 twitter4j 过滤器(或用户)“方法在内部创建了一个线程来操作 TwitterStream 并连续调用适当的侦听器方法”(引自示例Yusuke Yamamoto 的代码)。

我希望这会有所帮助。 (我无法发布更多链接,因为我收到“您需要至少 10 个声望才能发布超过 2 个链接”)

【讨论】:

  • 谢谢!似乎是这种行为的原因。此外,运行多个流线程的工作还不错,因此通常不应该指望它,但有时可以:)
猜你喜欢
  • 2019-06-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-08-01
  • 1970-01-01
  • 1970-01-01
  • 2012-06-15
相关资源
最近更新 更多