|
|
@ -182,7 +182,7 @@ class Reader(Worker):
|
|
|
|
try:
|
|
|
|
try:
|
|
|
|
|
|
|
|
|
|
|
|
self.reader = transport.factory.instance(**self._info) ;
|
|
|
|
self.reader = transport.factory.instance(**self._info) ;
|
|
|
|
# print (self.pipeline)
|
|
|
|
print (self._info)
|
|
|
|
# self.rows = self.reader.read(mongo=self.pipeline)
|
|
|
|
# self.rows = self.reader.read(mongo=self.pipeline)
|
|
|
|
self.rows = self.reader.read(**self.pipeline)
|
|
|
|
self.rows = self.reader.read(**self.pipeline)
|
|
|
|
|
|
|
|
|
|
|
@ -194,6 +194,7 @@ class Reader(Worker):
|
|
|
|
N = len(self.rows) / self.MAX_ROWS if len(self.rows) > self.MAX_ROWS else 1
|
|
|
|
N = len(self.rows) / self.MAX_ROWS if len(self.rows) > self.MAX_ROWS else 1
|
|
|
|
N = int(N)
|
|
|
|
N = int(N)
|
|
|
|
# self.rows = rows
|
|
|
|
# self.rows = rows
|
|
|
|
|
|
|
|
|
|
|
|
_log = {"context":self.name(), "status":1,"info":{"rows":len(self.rows),"table":self.table,"segments":N}}
|
|
|
|
_log = {"context":self.name(), "status":1,"info":{"rows":len(self.rows),"table":self.table,"segments":N}}
|
|
|
|
self.rows = np.array_split(self.rows,N)
|
|
|
|
self.rows = np.array_split(self.rows,N)
|
|
|
|
|
|
|
|
|
|
|
|