bug fix: ETL, Mongodb

v2.0
Steve Nyemba 8 months ago
parent 715e40407a
commit e1763b1b19

@ -62,8 +62,14 @@ def wait(jobs):
time.sleep(1)
@app.command(name="apply")
def move (path,index=None):
def apply (path,index=None):
"""
This function applies data transport from one source to one or several others
:path path of the configuration file
:index index of the _item of interest (otherwise everything will be processed)
"""
_proxy = lambda _object: _object.write(_object.read())
if os.path.exists(path):
file = open(path)

@ -83,7 +83,12 @@ class Transporter(Process):
_reader = transport.factory.instance(**self._source)
#
# If arguments are provided then a query is to be executed (not just a table dump)
return _reader.read() if 'args' not in self._source else _reader.read(**self._source['args'])
if 'cmd' in self._source or 'query' in self._source :
_query = self._source['cmd'] if 'cmd' in self._source else self._source['query']
return _reader.read(**_query)
else:
return _reader.read()
# return _reader.read() if 'query' not in self._source else _reader.read(**self._source['query'])
def _delegate_write(self,_data,**_args):
"""

@ -218,7 +218,7 @@ class Writer(Mongo):
if type(info) == pd.DataFrame :
info = info.to_dict(orient='records')
# info if type(info) == list else info.to_dict(orient='records')
info = json.loads(json.dumps(info))
info = json.loads(json.dumps(info,cls=IEncoder))
self.db[_uid].insert_many(info)
else:
#

Loading…
Cancel
Save