-
Notifications
You must be signed in to change notification settings - Fork 0
/
history.py
51 lines (42 loc) · 1.51 KB
/
history.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
import sys
from disco.core import Disco, result_iterator
from disco.settings import DiscoSettings
from disco import func
def map(line, params):
"""
hackreduce:search:history format:
None, timestamp, id, search, frequency?
"""
from datetime import datetime, timedelta
from disco.util import msg
try:
unknown, timestamp, uid, query, frequency = line.split("','")
except ValueError:
print msg(line)
# bad hack :-(
time = timestamp.replace("'", "")
date_obj = datetime.fromtimestamp(float(time[:-3])) # timestamp has milliseconds, shave em off
nearest_minute = date_obj - timedelta(minutes=date_obj.minute % 1, seconds=date_obj.second, microseconds=date_obj.microsecond)
yield (nearest_minute, {'unique_id': uid, 'query': query, 'frequency': frequency})
def reduce(iter, params):
# This doesn't work at all, its from an old example.
for unique_id, counts in kvgroup(sorted(iter)):
yield unique_id, sum(counts)
disco = Disco(DiscoSettings()['DISCO_MASTER'])
print "Starting Disco job.."
print "Go to %s to see status of the job." % disco.master
"""
:clicks (ad id,people who clicked the ads)
"""
results = disco.new_job(name="bartekc",
input=["tag://hackreduce:search:history"],
map_input_stream=(
func.map_input_stream,
func.chain_reader,
),
map=map,
reduce=reduce,
save=True).wait()
print "Job done. Results:"
for word, count in result_iterator(results):
print word, count