Fix up the demos a bit

This commit is contained in:
Eric Tschetter 2012-10-24 04:51:32 -04:00
parent 9d41599967
commit f09f595c7c
10 changed files with 78 additions and 71 deletions

View File

@ -1,17 +1,19 @@
{
"queryType": "topN",
"queryType": "groupBy",
"dataSource": "randSeq",
"granularity": "all",
"dimension": "target",
"threshold": 10,
"metric": "randomNumberSum",
"dimensions": [],
"aggregations":[
{ "type": "count", "name": "rows"},
{ "type": "doubleSum", "fieldName": "events", "name": "e"},
{ "type": "doubleSum", "fieldName": "outColumn", "name": "randomNumberSum"}
],
"postAggregations":[
{"type":"arithmetic","name":"avg_random","fn":"/","fields":[{"type":"fieldAccess","name":"randomNumberSum","fieldName":"randomNumberSum"},{"type":"fieldAccess","name":"rows","fieldName":"rows"}]}
{ "type":"arithmetic",
"name":"avg_random",
"fn":"/",
"fields":[ {"type":"fieldAccess","name":"randomNumberSum","fieldName":"randomNumberSum"},
{"type":"fieldAccess","name":"rows","fieldName":"rows"} ]}
],
"intervals":["2012-10-01T00:00/2020-01-01T00"]
}

View File

@ -18,7 +18,8 @@ import java.util.Random;
import static java.lang.Thread.sleep;
/** Random value sequence Firehost Factory named "rand".
/**
* Random value sequence Firehost Factory named "rand".
* Builds a Firehose that emits a stream of random numbers (outColumn, a positive double)
* with timestamps along with an associated token (target). This provides a timeseries
* that requires no network access for demonstration, characterization, and testing.
@ -63,25 +64,26 @@ import static java.lang.Thread.sleep;
}]
* </pre>
*
* Example query using POST to /druid/v2/?w (where w is an arbitrary parameter and the UTC date and time
* MUST be adjusted for the current hour):
* Example query using POST to /druid/v2/ (where UTC date and time MUST include the current hour):
* <pre>
{
"queryType": "topN",
"queryType": "groupBy",
"dataSource": "randSeq",
"granularity": "all",
"dimension": "target",
"threshold": 10,
"metric": "randomNumberSum",
"dimensions": [],
"aggregations":[
{ "type": "count", "name": "rows"},
{ "type": "doubleSum", "fieldName": "events", "name": "e"},
{ "type": "doubleSum", "fieldName": "outColumn", "name": "randomNumberSum"}
],
"postAggregations":[
{"type":"arithmetic","name":"avg_random","fn":"/","fields":[{"type":"fieldAccess","name":"randomNumberSum","fieldName":"randomNumberSum"},{"type":"fieldAccess","name":"rows","fieldName":"rows"}]}
{ "type":"arithmetic",
"name":"avg_random",
"fn":"/",
"fields":[ {"type":"fieldAccess","name":"randomNumberSum","fieldName":"randomNumberSum"},
{"type":"fieldAccess","name":"rows","fieldName":"rows"} ]}
],
"intervals":["2012-10-16T20:03/2012-10-16T21"]
"intervals":["2012-10-01T00:00/2020-01-01T00"]
}
* </pre>
*/

View File

@ -18,7 +18,8 @@ import org.codehaus.jackson.map.jsontype.NamedType;
import java.io.File;
import java.io.IOException;
/** Standalone Demo Realtime process.
/**
* Standalone Demo Realtime process.
* Created: 20121009T2050
*
* @author pbaclace

View File

@ -1,25 +1,22 @@
{
"queryType": "topN",
"queryType": "groupBy",
"dataSource": "twitterstream",
"granularity": "all",
"dimension": "lang",
"dimension": ["lang"],
"threshold": 10,
"metric": "totally",
"aggregations":[
{ "type": "count", "name": "rows"},
{ "type": "doubleSum", "fieldName": "events", "name": "e"},
{ "type": "doubleSum", "fieldName": "tweets", "name": "tweets"},
{ "type": "max", "fieldName": "maxStatusesCount", "name": "theMaxStatusesCount"},
{ "type": "max", "fieldName": "maxRetweetCount", "name": "theMaxRetweetCount"},
{ "type": "max", "fieldName": "max_statuses_count", "name": "theMaxStatusesCount"},
{ "type": "max", "fieldName": "max_retweet_count", "name": "theMaxRetweetCount"},
{ "type": "max", "fieldName": "maxFriendsCount", "name": "theMaxFriendsCount"},
{ "type": "max", "fieldName": "maxFollowerCount", "name": "theMaxFollowerCount"},
{ "type": "max", "fieldName": "max_friends_count", "name": "theMaxFriendsCount"},
{ "type": "max", "fieldName": "max_follower_count", "name": "theMaxFollowerCount"},
{ "type": "doubleSum", "fieldName": "totalStatusesCount", "name": "totally"}
{ "type": "doubleSum", "fieldName": "total_statuses_count", "name": "total_tweets_all_time"}
],
"postAggregations":[
{"type":"arithmetic","name":"avg_f","fn":"/","fields":[{"type":"fieldAccess","name":"dummy","fieldName":"totally"},{"type":"fieldAccess","name":"rows","fieldName":"rows"}]}
],
"intervals":["2012-10-01T00:00/2020-01-01T00"]
}

View File

@ -9,12 +9,12 @@ trap "exit 1" 1 2 3 15
[ ! -e query.body ] && echo "expecting file query.body to be in current directory" && exit 2
for delay in 1 30
for delay in 5 15 15 15 15 15 15 15 15 15 15
do
echo "sleep for $delay seconds..."
echo " "
sleep $delay
curl -X POST 'http://localhost:8080/druid/v2/?w' -H 'content-type: application/json' -d @query.body
curl -X POST 'http://localhost:8080/druid/v2/' -H 'content-type: application/json' -d @query.body
echo " "
echo " "
done

View File

@ -25,9 +25,6 @@ else
fi
trap "${PF_CLEANUP} ; exit 1" 1 2 3 15
# be sure to use UTC
export TZ=UTC
# props are set in src/main/resources/runtime.properties
[ -d /tmp/twitter_realtime ] && echo "cleaning up from previous run.." && /bin/rm -fr /tmp/twitter_realtime
@ -39,7 +36,7 @@ OPT_PROPS=""
# start RealtimeNode process
#
java -Xmx600m -classpath target/druid-examples-twitter-*-selfcontained.jar $OPT_PROPS -Dtwitter4j.http.useSSL=true -Ddruid.realtime.specFile=twitter_realtime.spec druid.examples.RealtimeStandaloneMain >RealtimeNode.out 2>&1 &
java -Duser.timezone=UTC -Dfile.encoding=UTF-8 -Xmx600m -classpath target/druid-examples-twitter-*-selfcontained.jar $OPT_PROPS -Dtwitter4j.http.useSSL=true -Ddruid.realtime.specFile=twitter_realtime.spec druid.examples.RealtimeStandaloneMain >RealtimeNode.out 2>&1 &
PID=$!
sleep 4
grep com.metamx.druid.realtime.TwitterSpritzerFirehoseFactory RealtimeNode.out | awk '{ print $7,$8,$9,$10,$11,$12,$13,$14,$15 }'

View File

@ -40,6 +40,17 @@ import static java.lang.Thread.*;
* is UTC):
* <pre>
* </pre>
*
*
* Notes on twitter.com HTTP (REST) API: v1.0 will be disabled around 2013-03 so v1.1 should be used;
* twitter4j 3.0 (not yet released) will support the v1.1 api.
* Specifically, we should be using https://stream.twitter.com/1.1/statuses/sample.json
* See: http://jira.twitter4j.org/browse/TFJ-186
*
* Notes on JSON parsing: as of twitter4j 2.2.x, the json parser has some bugs (ex: Status.toString()
* can have number format exceptions), so it might be necessary to extract raw json and process it
* separately. If so, set twitter4.jsonStoreEnabled=true and look at DataObjectFactory#getRawJSON();
* org.codehaus.jackson.map.ObjectMapper should be used to parse.
* @author pbaclace
*/
@JsonTypeName("twitzer")
@ -106,12 +117,12 @@ public class TwitterSpritzerFirehoseFactory implements FirehoseFactory {
final long startMsec = System.currentTimeMillis();
dimensions.add("htags");
dimensions.add("retweetCount");
dimensions.add("followerCount");
dimensions.add("friendsCount");
dimensions.add("retweet_count");
dimensions.add("follower_count");
dimensions.add("friends_count");
dimensions.add("lang");
dimensions.add("utcOffset");
dimensions.add("statusesCount");
dimensions.add("utc_offset");
dimensions.add("statuses_count");
//
// set up Twitter Spritzer
@ -185,7 +196,7 @@ public class TwitterSpritzerFirehoseFactory implements FirehoseFactory {
if (maxRunMinutes <= 0) {
return false;
} else {
return (System.currentTimeMillis() - startMsec) / 10000L >= maxRunMinutes;
return (System.currentTimeMillis() - startMsec) / 60000L >= maxRunMinutes;
}
}
@ -249,23 +260,23 @@ public class TwitterSpritzerFirehoseFactory implements FirehoseFactory {
}
long retweetCount = status.getRetweetCount();
theMap.put("retweetCount", retweetCount);
theMap.put("retweet_count", retweetCount);
User u = status.getUser();
if (u != null) {
theMap.put("followerCount", u.getFollowersCount());
theMap.put("friendsCount", u.getFriendsCount());
theMap.put("follower_count", u.getFollowersCount());
theMap.put("friends_count", u.getFriendsCount());
theMap.put("lang", u.getLang());
theMap.put("utcOffset", u.getUtcOffset()); // resolution in seconds, -1 if not available?
theMap.put("statusesCount", u.getStatusesCount());
theMap.put("utc_offset", u.getUtcOffset()); // resolution in seconds, -1 if not available?
theMap.put("statuses_count", u.getStatusesCount());
} else {
log.error("status.getUser() is null");
}
if (rowCount % 10 == 0) {
log.info("" + status.getCreatedAt() +
" followerCount=" + u.getFollowersCount() +
" friendsCount=" + u.getFriendsCount() +
" statusesCount=" + u.getStatusesCount() +
" retweetCount=" + retweetCount
" follower_count=" + u.getFollowersCount() +
" friends_count=" + u.getFriendsCount() +
" statuses_count=" + u.getStatusesCount() +
" retweet_count=" + retweetCount
);
}

View File

@ -107,8 +107,6 @@ druid.service=foo
druid.pusher.s3.bucket=
druid.pusher.s3.baseKey=
# TODO: should the next prop also work via runtime.properties ?
# next MUST be on command line, does not work here
druid.realtime.specFile=twitter_realtime.spec
#

View File

@ -1,24 +1,23 @@
[{
"schema" : { "dataSource":"twitterstream",
"aggregators":[
{"type":"count", "name":"events"},
{"type":"doubleSum","fieldName":"followerCount","name":"totalFollowerCount"},
{"type":"doubleSum","fieldName":"retweetCount","name":"totalRetweetCount"},
{"type":"doubleSum","fieldName":"friendsCount","name":"totalFriendsCount"},
{"type":"doubleSum","fieldName":"statusesCount","name":"totalStatusesCount"},
{"type":"count", "name":"tweets"},
{"type":"doubleSum","fieldName":"follower_count","name":"total_follower_count"},
{"type":"doubleSum","fieldName":"retweet_count","name":"tota_retweet_count"},
{"type":"doubleSum","fieldName":"friends_count","name":"total_friends_count"},
{"type":"doubleSum","fieldName":"statuses_count","name":"total_statuses_count"},
{"type":"min","fieldName":"followerCount","name":"minFollowerCount"},
{"type":"max","fieldName":"followerCount","name":"maxFollowerCount"},
{"type":"min","fieldName":"follower_count","name":"min_follower_count"},
{"type":"max","fieldName":"follower_count","name":"max_follower_count"},
{"type":"min","fieldName":"friendsCount","name":"minFriendsCount"},
{"type":"max","fieldName":"friendsCount","name":"maxFriendsCount"},
{"type":"min","fieldName":"friends_count","name":"min_friends_count"},
{"type":"max","fieldName":"friends_count","name":"max_friends_count"},
{"type":"min","fieldName":"statusesCount","name":"minStatusesCount"},
{"type":"max","fieldName":"statusesCount","name":"maxStatusesCount"},
{"type":"min","fieldName":"retweetCount","name":"minRetweetCount"},
{"type":"max","fieldName":"retweetCount","name":"maxRetweetCount"}
{"type":"min","fieldName":"statuses_count","name":"min_statuses_count"},
{"type":"max","fieldName":"statuses_count","name":"max_statuses_count"},
{"type":"min","fieldName":"retweet_count","name":"min_retweet_count"},
{"type":"max","fieldName":"retweet_count","name":"max_retweet_count"}
],
"indexGranularity":"minute",
"shardSpec" : { "type": "none" } },
@ -26,8 +25,8 @@
"intermediatePersistPeriod" : "PT2m" },
"firehose" : { "type" : "twitzer",
"maxEventCount": 10000,
"maxRunMinutes" : 5
"maxEventCount": 50000,
"maxRunMinutes" : 10
},
"plumber" : { "type" : "realtime",