diff options
author | Marko Kreen | 2007-03-13 11:52:09 +0000 |
---|---|---|
committer | Marko Kreen | 2007-03-13 11:52:09 +0000 |
commit | 50abdba44a031ad40b1886f941479f203ca92039 (patch) | |
tree | 873e72d78cd48917b2907c4c63abf185649ebb54 /scripts/queue_splitter.py |
final public releaseskytools_2_1
Diffstat (limited to 'scripts/queue_splitter.py')
-rwxr-xr-x | scripts/queue_splitter.py | 33 |
1 files changed, 33 insertions, 0 deletions
diff --git a/scripts/queue_splitter.py b/scripts/queue_splitter.py new file mode 100755 index 00000000..c6714ca0 --- /dev/null +++ b/scripts/queue_splitter.py @@ -0,0 +1,33 @@ +#! /usr/bin/env python + +# puts events into queue specified by field from 'queue_field' config parameter + +import sys, os, pgq, skytools + +class QueueSplitter(pgq.SerialConsumer): + def __init__(self, args): + pgq.SerialConsumer.__init__(self, "queue_splitter", "src_db", "dst_db", args) + + def process_remote_batch(self, db, batch_id, ev_list, dst_db): + cache = {} + queue_field = self.cf.get('queue_field', 'extra1') + for ev in ev_list: + row = [ev.type, ev.data, ev.extra1, ev.extra2, ev.extra3, ev.extra4, ev.time] + queue = ev.__getattr__(queue_field) + if queue not in cache: + cache[queue] = [] + cache[queue].append(row) + ev.tag_done() + + # should match the composed row + fields = ['type', 'data', 'extra1', 'extra2', 'extra3', 'extra4', 'time'] + + # now send them to right queues + curs = dst_db.cursor() + for queue, rows in cache.items(): + pgq.bulk_insert_events(curs, rows, fields, queue) + +if __name__ == '__main__': + script = QueueSplitter(sys.argv[1:]) + script.start() + |