add recipe object

This commit is contained in:
Jeremy Carbaugh 2009-06-26 17:54:57 -04:00
parent a0b206e711
commit 6f969ac3d4
2 changed files with 32 additions and 17 deletions

View File

@ -4,26 +4,39 @@
import filters, emitters, sources, utils import filters, emitters, sources, utils
def run_recipe(source, *filter_args): class Recipe(filters.Filter):
""" Process data, taking it from a source and applying any number of filters
""" def __init__(self, *filter_args):
self._filter_args = filter_args
self._rejected = []
def process_record(self, record):
self.run(source=iter((record,)))
def run(self, source):
# connect datapath # connect datapath
data = source data = source
for filter_ in filter_args: for filter_ in self._filter_args:
data = filter_(data) data = filter_(self, data)
# actually run the data through (causes iterators to actually be called) # actually run the data through (causes iterators to actually be called)
for record in data: for record in data:
pass pass
# try and call done() on all filters # try and call done() on all filters
for filter_ in filter_args: for filter_ in self._filter_args:
try: try:
filter_.done() filter_.done()
except AttributeError: except AttributeError:
pass # don't care if there isn't a done method pass # don't care if there isn't a done method
def run_recipe(source, *filter_args):
""" Process data, taking it from a source and applying any number of filters
"""
Recipe(*filter_args).run(source)
# experiment with threading - do not use # experiment with threading - do not use

View File

@ -33,9 +33,11 @@ class Filter(object):
raise NotImplementedError('process_record not defined in ' + raise NotImplementedError('process_record not defined in ' +
self.__class__.__name__) self.__class__.__name__)
def __call__(self, source): def __call__(self, recipe, source):
for record in source: for record in source:
yield self.process_record(record) result = self.process_record(record)
if not result is None:
yield result
class YieldFilter(Filter): class YieldFilter(Filter):