Skip to content
Draft
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
added disconnect functionality when stream is empty for 20 secs and s…
…ome exception handling
  • Loading branch information
kisaga committed Jun 17, 2022
commit 39e326193ac3ff21e822f922478ad41816b09169
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@ public static void main(String[] args) {

class Responder implements TweetsStreamListener {
@Override
public synchronized void actionOnTweetsStream(StreamingTweet streamingTweet) {
public synchronized void onTweetArrival(StreamingTweet streamingTweet) {
if(streamingTweet == null) {
System.err.println("Error: actionOnTweetsStream - streamingTweet is null ");
return;
Expand Down
14 changes: 1 addition & 13 deletions src/main/java/com/twitter/clientlib/api/TweetsApi.java
Original file line number Diff line number Diff line change
Expand Up @@ -23,25 +23,16 @@
package com.twitter.clientlib.api;

import com.twitter.clientlib.ApiCallback;
import com.twitter.clientlib.ApiClient;
import com.twitter.clientlib.auth.*;
import com.twitter.clientlib.ApiException;
import com.twitter.clientlib.ApiResponse;
import com.twitter.clientlib.Configuration;
import com.twitter.clientlib.Pair;
import com.twitter.clientlib.ProgressRequestBody;
import com.twitter.clientlib.ProgressResponseBody;

import com.github.scribejava.core.model.OAuth2AccessToken;
import com.google.gson.reflect.TypeToken;

import java.io.IOException;


import com.twitter.clientlib.model.AddOrDeleteRulesRequest;
import com.twitter.clientlib.model.AddOrDeleteRulesResponse;
import com.twitter.clientlib.model.CreateTweetRequest;
import com.twitter.clientlib.model.Error;
import com.twitter.clientlib.model.FilteredStreamingTweet;
import com.twitter.clientlib.model.GenericTweetsTimelineResponse;
import com.twitter.clientlib.model.GetRulesResponse;
Expand All @@ -52,7 +43,7 @@
import com.twitter.clientlib.model.MultiTweetLookupResponse;
import com.twitter.clientlib.model.MultiUserLookupResponse;
import java.time.OffsetDateTime;
import com.twitter.clientlib.model.Problem;

import com.twitter.clientlib.model.QuoteTweetLookupResponse;
import com.twitter.clientlib.model.SingleTweetLookupResponse;
import com.twitter.clientlib.model.StreamingTweet;
Expand All @@ -75,11 +66,8 @@
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.io.InputStream;
import javax.ws.rs.core.GenericType;

import okio.BufferedSource;
import org.apache.commons.lang3.StringUtils;

public class TweetsApi extends ApiCommon {

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
package com.twitter.clientlib.exceptions;

public class AuthenticationException extends RuntimeException {

public AuthenticationException(String message) {
super(message);
}

public AuthenticationException(String message, Throwable cause) {
super(message, cause);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
package com.twitter.clientlib.exceptions;

public class EmptyStreamTimeoutException extends RuntimeException {
public EmptyStreamTimeoutException(String message) {
super(message);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
package com.twitter.clientlib.stream;


import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.BlockingQueue;
Expand All @@ -36,10 +37,14 @@
import com.fasterxml.jackson.databind.JsonMappingException;
import com.fasterxml.jackson.databind.ObjectMapper;

import com.twitter.clientlib.exceptions.EmptyStreamTimeoutException;
import com.twitter.clientlib.model.StreamingTweet;
import okio.BufferedSource;

public class TweetsStreamExecutor {

private static final long EMPTY_STREAM_TIMEOUT = 20000;
private static final int POLL_WAIT = 5;
private volatile BlockingQueue<String> rawTweets;
private volatile BlockingQueue<StreamingTweet> tweets;
private volatile boolean isRunning = true;
Expand Down Expand Up @@ -90,9 +95,12 @@ public synchronized void shutdown() {
shutDownServices();
try {
terminateServices();
stream.close();
} catch (InterruptedException ie) {
shutDownServices();
Thread.currentThread().interrupt();
} catch (IOException e) {

}
System.out.println("TweetsStreamListenersExecutor is shutting down.");
}
Expand All @@ -109,9 +117,9 @@ private void terminateServices() throws InterruptedException {
terminateService(listenersService);
}
private void terminateService(ExecutorService executorService) throws InterruptedException {
if (!executorService.awaitTermination(3, TimeUnit.SECONDS)) {
if (!executorService.awaitTermination(1500, TimeUnit.MILLISECONDS)) {
executorService.shutdownNow();
if (!executorService.awaitTermination(3, TimeUnit.SECONDS))
if (!executorService.awaitTermination(1500, TimeUnit.MILLISECONDS))
System.err.println("Pool did not terminate");
}
}
Expand All @@ -126,19 +134,32 @@ public void run() {
public void queueTweets() {
String line = null;
try {
boolean emptyResponse = false;
long firstEmptyResponseMillis = 0;
long lastEmptyReponseMillis;
while (isRunning) {
line = stream.readUtf8Line();
if(line == null || line.isEmpty()) {
if(!emptyResponse) {
firstEmptyResponseMillis = System.currentTimeMillis();
emptyResponse = true;
} else {
lastEmptyReponseMillis = System.currentTimeMillis();
if(lastEmptyReponseMillis - firstEmptyResponseMillis > EMPTY_STREAM_TIMEOUT) {
throw new EmptyStreamTimeoutException(String.format("Stream was empty for %d seconds consecutively", EMPTY_STREAM_TIMEOUT));
}
}
continue;
}
emptyResponse = false;
try {
rawTweets.put(line);
} catch (Exception interExcep) {
interExcep.printStackTrace();
} catch (Exception ignore) {

}
}
} catch (Exception e) {
e.printStackTrace();
System.out.println("Something went wrong. Closing stream... " + e.getMessage());
shutdown();
}
}
Expand All @@ -156,15 +177,14 @@ private DeserializeTweetsTask() {
public void run() {
while (isRunning) {
try {
String rawTweet = rawTweets.take();
String rawTweet = rawTweets.poll(POLL_WAIT, TimeUnit.MILLISECONDS);
if (rawTweet == null) continue;
StreamingTweet tweet = objectMapper.readValue(rawTweet, StreamingTweet.class);
tweets.put(tweet);
} catch (InterruptedException e) {
System.out.println("Fail 1");
} catch (JsonMappingException e) {
System.out.println("Fail 2");

} catch (JsonProcessingException e) {
System.out.println("Fail 3");
System.out.println("debug log here");
}
}
}
Expand All @@ -181,9 +201,10 @@ private void processTweets() {

while (isRunning) {
try {
streamingTweet = tweets.take();
streamingTweet = tweets.poll(POLL_WAIT, TimeUnit.MILLISECONDS);
if(streamingTweet == null) continue;
for (TweetsStreamListener listener : listeners) {
listener.actionOnTweetsStream(streamingTweet);
listener.onTweetArrival(streamingTweet);
}
tweetsCount++;
if(tweetsCount == tweetsLimit) {
Expand All @@ -194,7 +215,7 @@ private void processTweets() {
shutdown();
}
} catch (InterruptedException e) {
System.out.println("processTweets: Fail 1");

}

}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,5 +25,5 @@
import com.twitter.clientlib.model.StreamingTweet;

public interface TweetsStreamListener {
void actionOnTweetsStream(StreamingTweet streamingTweet);
void onTweetArrival(StreamingTweet streamingTweet);
}
11 changes: 5 additions & 6 deletions src/main/java/com/twitter/clientlib/stream/TwitterStream.java
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,10 @@
import com.twitter.clientlib.ApiException;
import com.twitter.clientlib.TwitterCredentialsBearer;
import com.twitter.clientlib.api.TweetsApi;
import com.twitter.clientlib.exceptions.AuthenticationException;
import com.twitter.clientlib.query.StreamQueryParameters;
import okio.BufferedSource;

import java.io.InputStream;
import java.util.LinkedList;
import java.util.List;

Expand Down Expand Up @@ -38,15 +38,14 @@ public void sampleStream(StreamQueryParameters streamParameters) {
listeners.forEach(executor::addListener);
executor.start();
} catch (ApiException e) {
System.err.println("Status code: " + e.getCode());
System.err.println("Reason: " + e.getResponseBody());
System.err.println("Response headers: " + e.getResponseHeaders());
e.printStackTrace();
if(e.getCode() == 401) {
throw new AuthenticationException("Not authenticated. Please check the credentials", e);
}
}
}

private void initBasePath() {
String basePath = "http://localhost:8080";
String basePath = System.getenv("TWITTER_API_BASE_PATH");
apiClient.setBasePath(basePath != null ? basePath : "https://api.twitter.com");
}

Expand Down