From 56bdda17b70b23c0b0af4b413432bc65b43a3654 Mon Sep 17 00:00:00 2001 From: Steve Nyemba Date: Thu, 30 Nov 2023 12:40:57 -0600 Subject: [PATCH 1/4] minor bug fix --- transport/disk.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/transport/disk.py b/transport/disk.py index 956386d..42b5b33 100644 --- a/transport/disk.py +++ b/transport/disk.py @@ -219,9 +219,10 @@ class SQLiteWriter(SQLite,DiskWriter) : elif type(info) == pd.DataFrame : info = info.fillna('') info = info.to_dict(orient='records') - if not self.fields : - _rec = info[0] - self.init(list(_rec.keys())) + + if not self.fields : + _rec = info[0] + self.init(list(_rec.keys())) SQLiteWriter.LOCK.acquire() try: -- 2.34.1 From d74372f645630fb100a4cb7d2afa2c421b426df4 Mon Sep 17 00:00:00 2001 From: Steve Nyemba Date: Fri, 8 Dec 2023 18:19:46 -0600 Subject: [PATCH 2/4] bug fixes: mongodb, common, nextcloud --- transport/__init__.py | 24 ++++++++++++++---------- transport/common.py | 15 ++++++++++++++- transport/disk.py | 7 ++++++- transport/mongo.py | 6 +++--- transport/nextcloud.py | 4 ++-- 5 files changed, 39 insertions(+), 17 deletions(-) diff --git a/transport/__init__.py b/transport/__init__.py index e139aa5..234c418 100644 --- a/transport/__init__.py +++ b/transport/__init__.py @@ -27,6 +27,7 @@ import json import importlib import sys import sqlalchemy +from datetime import datetime if sys.version_info[0] > 2 : # from transport.common import Reader, Writer,Console #, factory from transport import disk @@ -83,16 +84,19 @@ import os # PGSQL = POSTGRESQL # import providers -class IEncoder (json.JSONEncoder): - def default (self,object): - if type(object) == np.integer : - return int(object) - elif type(object) == np.floating: - return float(object) - elif type(object) == np.ndarray : - return object.tolist() - else: - return super(IEncoder,self).default(object) +# class IEncoder (json.JSONEncoder): +def IEncoder (self,object): + if type(object) == np.integer : + return int(object) + elif type(object) == np.floating: + return float(object) + elif type(object) == np.ndarray : + return object.tolist() + elif type(object) == datetime : + return o.isoformat() + else: + return super(IEncoder,self).default(object) + class factory : TYPE = {"sql":{"providers":["postgresql","mysql","neteeza","bigquery","mariadb","redshift"]}} PROVIDERS = { diff --git a/transport/common.py b/transport/common.py index 59f57ea..8b9f718 100644 --- a/transport/common.py +++ b/transport/common.py @@ -25,7 +25,7 @@ from multiprocessing import RLock import queue # import couch # import mongo - +from datetime import datetime class IO: def init(self,**args): @@ -39,6 +39,19 @@ class IO: continue value = args[field] setattr(self,field,value) +class IEncoder (json.JSONEncoder): + def default (self,object): + if type(object) == np.integer : + return int(object) + elif type(object) == np.floating: + return float(object) + elif type(object) == np.ndarray : + return object.tolist() + elif type(object) == datetime : + return object.isoformat() + else: + return super(IEncoder,self).default(object) + class Reader (IO): """ This class is an abstraction of a read functionalities of a data store diff --git a/transport/disk.py b/transport/disk.py index 42b5b33..2c9f6c8 100644 --- a/transport/disk.py +++ b/transport/disk.py @@ -12,6 +12,8 @@ import json import sqlite3 import pandas as pd from multiprocessing import Lock +from transport.common import Reader, Writer, IEncoder + class DiskReader(Reader) : """ This class is designed to read data from disk (location on hard drive) @@ -221,6 +223,8 @@ class SQLiteWriter(SQLite,DiskWriter) : info = info.to_dict(orient='records') if not self.fields : + + _rec = info[0] self.init(list(_rec.keys())) @@ -231,7 +235,8 @@ class SQLiteWriter(SQLite,DiskWriter) : sql = " " .join(["INSERT INTO ",self.table,"(", ",".join(self.fields) ,")", "values(:values)"]) for row in info : stream =["".join(["",value,""]) if type(value) == str else value for value in row.values()] - stream = json.dumps(stream).replace("[","").replace("]","") + stream = json.dumps(stream,cls=IEncoder) + stream = stream.replace("[","").replace("]","") self.conn.execute(sql.replace(":values",stream) ) diff --git a/transport/mongo.py b/transport/mongo.py index c24b4b8..bac1780 100644 --- a/transport/mongo.py +++ b/transport/mongo.py @@ -15,7 +15,7 @@ import gridfs # from transport import Reader,Writer import sys if sys.version_info[0] > 2 : - from transport.common import Reader, Writer + from transport.common import Reader, Writer, IEncoder else: from common import Reader, Writer import json @@ -102,7 +102,7 @@ class MongoReader(Mongo,Reader): if 'pipeline' in args : cmd['pipeline']= args['pipeline'] if 'aggregate' not in cmd : - cmd['aggregate'] = self.collection + cmd['aggregate'] = self.uid if 'pipeline' not in args or 'aggregate' not in cmd : cmd = args['mongo'] if 'mongo' in args else args['cmd'] if "aggregate" in cmd : @@ -182,7 +182,7 @@ class MongoWriter(Mongo,Writer): for row in rows : if type(row['_id']) == ObjectId : row['_id'] = str(row['_id']) - stream = Binary(json.dumps(collection).encode()) + stream = Binary(json.dumps(collection,cls=IEncoder).encode()) collection.delete_many({}) now = "-".join([str(datetime.now().year()),str(datetime.now().month), str(datetime.now().day)]) name = ".".join([self.uid,'archive',now])+".json" diff --git a/transport/nextcloud.py b/transport/nextcloud.py index f096f70..2eefd51 100644 --- a/transport/nextcloud.py +++ b/transport/nextcloud.py @@ -3,7 +3,7 @@ We are implementing transport to and from nextcloud (just like s3) """ import os import sys -from transport.common import Reader,Writer +from transport.common import Reader,Writer, IEncoder import pandas as pd from io import StringIO import json @@ -73,7 +73,7 @@ class NextcloudWriter (Nextcloud,Writer): _data.to_csv(f,index=False) _content = f.getvalue() elif type(_data) == dict : - _content = json.dumps(_data) + _content = json.dumps(_data,cls=IEncoder) else: _content = str(_data) self._handler.put_file_contents(_uri,_content) -- 2.34.1 From e46ebadcc2c0a383d838c83e6ddba7a82dae8162 Mon Sep 17 00:00:00 2001 From: Steve Nyemba Date: Mon, 11 Dec 2023 22:10:53 -0600 Subject: [PATCH 3/4] bug fix --- transport/disk.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/transport/disk.py b/transport/disk.py index 2c9f6c8..424e95e 100644 --- a/transport/disk.py +++ b/transport/disk.py @@ -223,8 +223,6 @@ class SQLiteWriter(SQLite,DiskWriter) : info = info.to_dict(orient='records') if not self.fields : - - _rec = info[0] self.init(list(_rec.keys())) -- 2.34.1 From fbfaaebbdc4300c137a0209d1a83f9ebecc6ac21 Mon Sep 17 00:00:00 2001 From: Steve Nyemba Date: Wed, 20 Dec 2023 10:14:15 -0600 Subject: [PATCH 4/4] adding alias CALLBACK == QLISTENER --- transport/providers.py | 1 + 1 file changed, 1 insertion(+) diff --git a/transport/providers.py b/transport/providers.py index 23843e7..93b6f53 100644 --- a/transport/providers.py +++ b/transport/providers.py @@ -51,6 +51,7 @@ RABBIT = RABBITMQ QLISTENER = 'qlistener' QUEUE = QLISTENER +CALLBACK = QLISTENER DATABRICKS= 'databricks+connector' DRIVERS = {PG:pg,REDSHIFT:pg,MYSQL:my,MARIADB:my,NETEZZA:nz,SQLITE:sqlite3} CATEGORIES ={'sql':[NETEZZA,PG,MYSQL,REDSHIFT,SQLITE,MARIADB],'nosql':[MONGODB,COUCHDB],'cloud':[NEXTCLOUD,S3,BIGQUERY,DATABRICKS],'file':[FILE], -- 2.34.1